diff --git a/docs/architecture/media-requests.md b/docs/architecture/media-requests.md index e2fd20cf11..95f839268f 100644 --- a/docs/architecture/media-requests.md +++ b/docs/architecture/media-requests.md @@ -119,6 +119,10 @@ marked 4K (the plugin's `is_4k` switch, or its older `is_default_4k`), HD versions only to one that is not. Saving a route that breaks this is refused, and so, under Advanced, is changing a server's 4K switch while a route sends it the other version. A server of another plugin (Seerr) takes either version. +Changing a server's type or media types while a route sends it a media type it +would no longer take is refused too. A route save and a server save each check +the other again with the server row locked, so two admins saving at once +cannot leave a route pointing at a server that no longer fits. The migration that introduced routes carried the Sonarr/Radarr plugin's routing over unchanged: each media type's first usable default and default-4K servers @@ -305,15 +309,20 @@ row until the follow commits, so a follow cannot land just after the request was declined, cancelled or completed, and miss that transition's follow cleanup. -A follow belongs to the title and the profile (`media_request_follows`, keyed -by account and profile id, since profile ids repeat across accounts), not to -one request, so it survives the request failing and being retried or requested -again. Declining or cancelling the request clears the title's follows: the title -is no longer on its way, and the follower can request it themselves. The -requesting profile never needs a follow: the fulfilled notification always -reaches it. When a request's fulfilled notification goes out, it is also sent -to every follower of the title, marked `follower` so its wording does not say -"your request", and those follows are then cleared. A dispatch failure leaves +A follow belongs to a profile and the request that was open when it was made +(`media_request_follows`, keyed by account, profile id and request, since +profile ids repeat across accounts). A series can have completed requests still +waiting for the library beside a newer open request for other seasons; each +request's notification goes to its own follows, and a profile can follow each +of them. Unfollowing a title removes the profile's follows on all of them. A follow +survives its request failing: the title's next request takes over the follows +of a failed request, or one its requester replaced. Declining or cancelling a +request clears its follows: the title is no longer on its way, and the follower +can request it themselves. The requesting profile never needs a follow: the +fulfilled notification always reaches it. When a request's fulfilled +notification goes out, it is also sent to the request's followers, marked +`follower` so its wording does not say "your request", and those follows are +then cleared. A dispatch failure leaves the follows for the retry, and the server-channel announcement waits until an attempt has reached every recipient, so a retry does not repeat it. diff --git a/go.mod b/go.mod index a77f215cba..bf0511ae62 100644 --- a/go.mod +++ b/go.mod @@ -129,7 +129,7 @@ require ( ) require ( - github.com/Silo-Server/silo-plugin-sdk v0.18.1-0.20260928194839-1abd582d0304 + github.com/Silo-Server/silo-plugin-sdk v0.19.0 github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.14 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.30 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.30 // indirect diff --git a/go.sum b/go.sum index c6fd89f43d..0d726fca83 100644 --- a/go.sum +++ b/go.sum @@ -6,8 +6,8 @@ github.com/PuerkitoBio/goquery v1.8.0 h1:PJTF7AmFCFKk1N6V6jmKfrNH9tV5pNE6lZMkG0g github.com/PuerkitoBio/goquery v1.8.0/go.mod h1:ypIiRMtY7COPGk+I/YbZLbxsxn9g5ejnI2HSMtkjZvI= github.com/SherClockHolmes/webpush-go v1.4.0 h1:ocnzNKWN23T9nvHi6IfyrQjkIc0oJWv1B1pULsf9i3s= github.com/SherClockHolmes/webpush-go v1.4.0/go.mod h1:XSq8pKX11vNV8MJEMwjrlTkxhAj1zKfxmyhdV7Pd6UA= -github.com/Silo-Server/silo-plugin-sdk v0.18.1-0.20260928194839-1abd582d0304 h1:D8lB4rqvGxOQTw9DeJiE7Dwr6qTVk2xqjI5uuyJM+i0= -github.com/Silo-Server/silo-plugin-sdk v0.18.1-0.20260928194839-1abd582d0304/go.mod h1:abwsCEKuPAAgeAqpNGbwoaut2eQlC/Kj97u89Vvg9qM= +github.com/Silo-Server/silo-plugin-sdk v0.19.0 h1:LWYI9x6OxBr8d+IxsHjJpcABpjvx/lA4UsDD13ynlAg= +github.com/Silo-Server/silo-plugin-sdk v0.19.0/go.mod h1:abwsCEKuPAAgeAqpNGbwoaut2eQlC/Kj97u89Vvg9qM= github.com/TwiN/go-color v1.4.1 h1:mqG0P/KBgHKVqmtL5ye7K0/Gr4l6hTksPgTgMk3mUzc= github.com/TwiN/go-color v1.4.1/go.mod h1:WcPf/jtiW95WBIsEeY1Lc/b8aaWoiqQpu5cf8WFxu+s= github.com/abadojack/whatlanggo v1.0.1 h1:19N6YogDnf71CTHm3Mp2qhYfkRdyvbgwWdd2EPxJRG4= diff --git a/internal/requests/editor_concurrency.go b/internal/requests/editor_concurrency.go index 73726d3acd..71ed83e6e9 100644 --- a/internal/requests/editor_concurrency.go +++ b/internal/requests/editor_concurrency.go @@ -113,10 +113,8 @@ func (r *Repository) UpdateIntegrationConditional(ctx context.Context, in Integr if err = lockRevision(ctx, tx, `SELECT revision FROM request_integrations WHERE id=$1 FOR UPDATE`, []any{in.ID}, expected, false); err != nil { return nil, err } - if !standard { - if err = ensureTierKeptUnderAdvanced(ctx, tx, in); err != nil { - return nil, err - } + if err = ensureRoutesStillFit(ctx, tx, in, !standard); err != nil { + return nil, err } out, err := r.updateIntegration(ctx, tx, in) if err == nil && standard { diff --git a/internal/requests/follows.go b/internal/requests/follows.go index 8dc8210f9f..fec62ad96c 100644 --- a/internal/requests/follows.go +++ b/internal/requests/follows.go @@ -3,17 +3,23 @@ package requests import ( "context" "fmt" + "maps" + "slices" "strings" ) // Following a title: a profile that finds a title someone else has already // requested can ask to be notified when it becomes available, instead of -// requesting it again. A follow belongs to the title, not to one request, so it -// survives the request failing and being retried or requested again. It is -// cleared once the fulfilled notification has gone out, and when the request -// is declined or withdrawn: the title is then no longer on its way, and the -// follower can request it themselves. The requester is always notified and -// never needs a follow. +// requesting it again. A follow belongs to the request that was open when it +// was made, since a series can have completed requests still waiting for the +// library beside a newer open request for other seasons; each request's +// notification goes to its own follows, and a profile can follow each of them. +// A follow survives its request failing: the title's next request takes over +// the follows of a failed or replaced one. +// It is cleared once the fulfilled notification has gone out, and when its +// request is declined or withdrawn: the title is then no longer on its way, +// and the follower can request it themselves. The requester is always +// notified and never needs a follow. // Follower is a profile waiting to hear that a requested title is available. type Follower struct { @@ -97,7 +103,7 @@ func (s *Service) Unfollow(ctx context.Context, viewer Viewer, mediaType MediaTy // is waiting on, either as the requesting profile or as a follower. func (s *Service) followedTitles(ctx context.Context, viewer Viewer, mediaType MediaType, active map[int]*Request) (map[int]bool, error) { out := map[int]bool{} - var others []int + others := map[string]int{} for tmdbID, req := range active { if req == nil { continue @@ -106,18 +112,18 @@ func (s *Service) followedTitles(ctx context.Context, viewer Viewer, mediaType M out[tmdbID] = true continue } - others = append(others, tmdbID) + others[req.ID] = tmdbID } if len(others) == 0 || strings.TrimSpace(viewer.ProfileID) == "" { return out, nil } - followed, err := s.store.FollowedTitles(ctx, mediaType, others, viewer) + followed, err := s.store.FollowedRequests(ctx, slices.Collect(maps.Keys(others)), viewer) if err != nil { return nil, err } - for tmdbID, ok := range followed { + for id, ok := range followed { if ok { - out[tmdbID] = true + out[others[id]] = true } } return out, nil @@ -127,7 +133,8 @@ func (s *Service) followedTitles(ctx context.Context, viewer Viewer, mediaType M // the same statement, so a follow cannot land just after the request // completed and never be told. It answers ErrNotRequested when there is none. // -// FOR SHARE holds the open request until the follow commits. Every transition +// The follow records the open request it read. FOR SHARE holds that request +// until the follow commits. Every transition // that closes a request (decline, cancel, completion) updates its row, and // that row lock conflicts with FOR SHARE, so the close cannot commit, and its // follow cleanup cannot run, between the read and the insert. A follow that @@ -137,15 +144,15 @@ func (r *Repository) FollowTitle(ctx context.Context, mediaType MediaType, tmdbI var open bool if err := r.pool.QueryRow(ctx, ` WITH open_request AS ( - SELECT 1 FROM media_requests + SELECT id FROM media_requests WHERE media_type = $1 AND provider = 'tmdb' AND tmdb_id = $2 AND outcome = 'active' AND status <> 'completed' LIMIT 1 FOR SHARE ), inserted AS ( - INSERT INTO media_request_follows (media_type, tmdb_id, user_id, profile_id) - SELECT $1, $2, $3, $4 FROM open_request - ON CONFLICT (media_type, tmdb_id, user_id, profile_id) DO NOTHING + INSERT INTO media_request_follows (media_type, tmdb_id, user_id, profile_id, request_id) + SELECT $1, $2, $3, $4, id FROM open_request + ON CONFLICT (user_id, profile_id, request_id) DO NOTHING ) SELECT EXISTS (SELECT 1 FROM open_request) `, mediaType, tmdbID, viewer.UserID, viewer.ProfileID).Scan(&open); err != nil { @@ -157,23 +164,46 @@ func (r *Repository) FollowTitle(ctx context.Context, mediaType MediaType, tmdbI return nil } -// forgetTitleFollows removes the follows on the title of a request that was -// just declined or withdrawn, unless the title has another open request whose -// followers are still waiting for it. +// forgetTitleFollows removes the follows of a request that was just declined +// or withdrawn. func forgetTitleFollows(ctx context.Context, exec requestExecutor, closed *Request) error { - if _, err := exec.Exec(ctx, ` - DELETE FROM media_request_follows - WHERE media_type = $1 AND tmdb_id = $2 - AND NOT EXISTS ( - SELECT 1 FROM media_requests - WHERE media_type = $1 AND provider = 'tmdb' AND tmdb_id = $2 - AND outcome = 'active' AND status <> 'completed' AND id <> $3) - `, closed.MediaType, closed.TMDBID, closed.ID); err != nil { + if _, err := exec.Exec(ctx, `DELETE FROM media_request_follows WHERE request_id = $1`, closed.ID); err != nil { return fmt.Errorf("forget title follows: %w", err) } return nil } +// adoptTitleFollows gives a new request the follows of the title's failed +// requests, so a follow survives its request failing. The caller creates the +// request in the same transaction, before deleting any failed request it +// replaces. +// +// The follows move in place: an UnfollowTitle that waited on a moved row +// re-checks the moved row, still on the title and the profile, and deletes it, +// where a delete and re-insert would leave it a row it cannot see. +func adoptTitleFollows(ctx context.Context, exec requestExecutor, req *Request) error { + const failed = `SELECT id FROM media_requests + WHERE media_type = $1 AND provider = 'tmdb' AND tmdb_id = $2 AND outcome = 'failed'` + // A profile that followed two failed requests keeps its earliest follow. + if _, err := exec.Exec(ctx, ` + DELETE FROM media_request_follows f + USING media_request_follows keep + WHERE f.request_id IN (`+failed+`) AND keep.request_id IN (`+failed+`) + AND keep.user_id = f.user_id AND keep.profile_id = f.profile_id + AND (keep.created_at, keep.request_id) < (f.created_at, f.request_id) + `, req.MediaType, req.TMDBID); err != nil { + return fmt.Errorf("adopt title follows: %w", err) + } + if _, err := exec.Exec(ctx, ` + UPDATE media_request_follows SET request_id = $3 + WHERE request_id IN (`+failed+`) + `, req.MediaType, req.TMDBID, req.ID); err != nil { + return fmt.Errorf("adopt title follows: %w", err) + } + return nil +} + +// UnfollowTitle removes the profile's follows on every request of the title. func (r *Repository) UnfollowTitle(ctx context.Context, mediaType MediaType, tmdbID int, viewer Viewer) error { if _, err := r.pool.Exec(ctx, ` DELETE FROM media_request_follows @@ -184,37 +214,39 @@ func (r *Repository) UnfollowTitle(ctx context.Context, mediaType MediaType, tmd return nil } -func (r *Repository) FollowedTitles(ctx context.Context, mediaType MediaType, tmdbIDs []int, viewer Viewer) (map[int]bool, error) { - out := map[int]bool{} - if len(tmdbIDs) == 0 { +// FollowedRequests reports which of the requests the profile follows. +func (r *Repository) FollowedRequests(ctx context.Context, requestIDs []string, viewer Viewer) (map[string]bool, error) { + out := map[string]bool{} + if len(requestIDs) == 0 { return out, nil } rows, err := r.pool.Query(ctx, ` - SELECT tmdb_id FROM media_request_follows - WHERE media_type = $1 AND tmdb_id = ANY($2) AND user_id = $3 AND profile_id = $4 - `, mediaType, tmdbIDs, viewer.UserID, viewer.ProfileID) + SELECT request_id FROM media_request_follows + WHERE request_id = ANY($1) AND user_id = $2 AND profile_id = $3 + `, requestIDs, viewer.UserID, viewer.ProfileID) if err != nil { - return nil, fmt.Errorf("list followed titles: %w", err) + return nil, fmt.Errorf("list followed requests: %w", err) } defer rows.Close() for rows.Next() { - var tmdbID int - if err := rows.Scan(&tmdbID); err != nil { + var id string + if err := rows.Scan(&id); err != nil { return nil, err } - out[tmdbID] = true + out[id] = true } return out, rows.Err() } -func (r *Repository) ListTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int) ([]Follower, error) { +// ListRequestFollowers lists the follows a request's notification goes to. +func (r *Repository) ListRequestFollowers(ctx context.Context, req Request) ([]Follower, error) { rows, err := r.pool.Query(ctx, ` SELECT user_id, profile_id FROM media_request_follows - WHERE media_type = $1 AND tmdb_id = $2 + WHERE request_id = $1 ORDER BY created_at, user_id, profile_id - `, mediaType, tmdbID) + `, req.ID) if err != nil { - return nil, fmt.Errorf("list title followers: %w", err) + return nil, fmt.Errorf("list request followers: %w", err) } defer rows.Close() var out []Follower @@ -228,7 +260,10 @@ func (r *Repository) ListTitleFollowers(ctx context.Context, mediaType MediaType return out, rows.Err() } -func (r *Repository) ClearTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int, followers []Follower) error { +// ClearRequestFollowers removes the listed follows once the request's +// notification has gone out. A profile that unfollowed and followed again +// since, for a newer request, keeps its new follow. +func (r *Repository) ClearRequestFollowers(ctx context.Context, req Request, followers []Follower) error { if len(followers) == 0 { return nil } @@ -240,10 +275,10 @@ func (r *Repository) ClearTitleFollowers(ctx context.Context, mediaType MediaTyp } if _, err := r.pool.Exec(ctx, ` DELETE FROM media_request_follows - WHERE media_type = $1 AND tmdb_id = $2 - AND (user_id, profile_id) IN (SELECT * FROM unnest($3::int[], $4::text[])) - `, mediaType, tmdbID, userIDs, profileIDs); err != nil { - return fmt.Errorf("clear title followers: %w", err) + WHERE request_id = $1 + AND (user_id, profile_id) IN (SELECT * FROM unnest($2::int[], $3::text[])) + `, req.ID, userIDs, profileIDs); err != nil { + return fmt.Errorf("clear request followers: %w", err) } return nil } diff --git a/internal/requests/follows_test.go b/internal/requests/follows_test.go index ba9470d675..1ab8a0fb08 100644 --- a/internal/requests/follows_test.go +++ b/internal/requests/follows_test.go @@ -3,9 +3,14 @@ package requests import ( "context" "errors" + "fmt" + "slices" "testing" "time" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/Silo-Server/silo-server/internal/metadata/tmdb" ) @@ -34,8 +39,8 @@ func TestFollowTitleSomeoneElseRequested(t *testing.T) { if !state.Following || state.RequestedByViewer || state.Requestable || state.Reason != "already_requested" || state.RequestID != "" { t.Fatalf("state = %+v, want following, not requestable, request id hidden from another account", state) } - followed, _ := store.FollowedTitles(context.Background(), MediaTypeMovie, []int{949}, testViewer(1)) - if !followed[949] { + followed, _ := store.FollowedRequests(context.Background(), []string{"req-owner"}, testViewer(1)) + if !followed["req-owner"] { t.Fatal("follow was not stored") } if _, err := svc.Follow(context.Background(), testViewer(1), MediaTypeMovie, 949); err != nil { @@ -45,8 +50,8 @@ func TestFollowTitleSomeoneElseRequested(t *testing.T) { if err := svc.Unfollow(context.Background(), testViewer(1), MediaTypeMovie, 949); err != nil { t.Fatalf("Unfollow: %v", err) } - followed, _ = store.FollowedTitles(context.Background(), MediaTypeMovie, []int{949}, testViewer(1)) - if followed[949] { + followed, _ = store.FollowedRequests(context.Background(), []string{"req-owner"}, testViewer(1)) + if followed["req-owner"] { t.Fatal("follow survived Unfollow") } } @@ -86,7 +91,7 @@ func TestFollowSameProfileIDOnAnotherAccount(t *testing.T) { if !state.Following || state.RequestedByViewer { t.Fatalf("state = %+v, want following and not requested by the viewer", state) } - followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 949) + followers, _ := store.titleFollowers(MediaTypeMovie, 949) if len(followers) != 1 || followers[0] != (Follower{UserID: 1, ProfileID: "default"}) { t.Fatalf("followers = %+v, want the other account's default profile", followers) } @@ -159,7 +164,7 @@ func TestWithdrawingRequestForgetsFollows(t *testing.T) { if err := tc.withdraw(svc); err != nil { t.Fatalf("%s: %v", tc.name, err) } - if followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 949); len(followers) != 0 { + if followers, _ := store.titleFollowers(MediaTypeMovie, 949); len(followers) != 0 { t.Fatalf("followers after %s = %+v, want none", tc.name, followers) } }) @@ -213,11 +218,35 @@ func TestNotifyFulfilledTellsFollowersAndClearsThem(t *testing.T) { if len(notifier.followers) != 1 || len(notifier.followers[0]) != 1 || notifier.followers[0][0] != (Follower{UserID: 3, ProfileID: "follower-profile"}) { t.Fatalf("followers handed to the notifier = %+v, want the one follower", notifier.followers) } - if followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42); len(followers) != 0 { + if followers, _ := store.titleFollowers(MediaTypeMovie, 42); len(followers) != 0 { t.Fatalf("followers after notifying = %+v, want cleared", followers) } } +// A series can have a completed request waiting for the library beside a newer +// open request for other seasons. Its notification goes to its own follows; +// the open request's follows wait for it. +func TestNotifyFulfilledLeavesFollowsOfNewerRequest(t *testing.T) { + store := newFakeStore() + store.requests["req1"] = completedRequestFixture("req1", 42) + store.unnotified = []string{"req1"} + store.seedFollowFor(MediaTypeMovie, 42, Viewer{UserID: 3, ProfileID: "early"}, "req1") + store.seedFollowFor(MediaTypeMovie, 42, Viewer{UserID: 4, ProfileID: "late"}, "req2") + notifier := &fakeNotifier{} + svc := NewService(store, &fakeTMDBClient{}, presentMovie(42)) + svc.SetFulfillmentNotifier(notifier) + + svc.notifyFulfilledPending(context.Background()) + + if len(notifier.followers) != 1 || !slices.Equal(notifier.followers[0], []Follower{{UserID: 3, ProfileID: "early"}}) { + t.Fatalf("followers handed to the notifier = %+v, want only req1's follower", notifier.followers) + } + left, _ := store.ListRequestFollowers(context.Background(), Request{ID: "req2", MediaType: MediaTypeMovie, TMDBID: 42}) + if !slices.Equal(left, []Follower{{UserID: 4, ProfileID: "late"}}) { + t.Fatalf("req2's followers = %+v, want kept", left) + } +} + func TestNotifyFulfilledKeepsFollowersWhenDispatchFails(t *testing.T) { store := newFakeStore() store.requests["req1"] = completedRequestFixture("req1", 42) @@ -228,7 +257,7 @@ func TestNotifyFulfilledKeepsFollowersWhenDispatchFails(t *testing.T) { svc.notifyFulfilledPending(context.Background()) - if followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42); len(followers) != 1 { + if followers, _ := store.titleFollowers(MediaTypeMovie, 42); len(followers) != 1 { t.Fatalf("followers after a failed dispatch = %+v, want kept for the retry", followers) } } @@ -254,7 +283,7 @@ func TestNotifyFulfilledRetriesWhenClearingFollowersFails(t *testing.T) { if len(store.unnotified) != 0 { t.Fatalf("unnotified after the retry = %v, want none", store.unnotified) } - if followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42); len(followers) != 0 { + if followers, _ := store.titleFollowers(MediaTypeMovie, 42); len(followers) != 0 { t.Fatalf("followers after the retry = %+v, want cleared", followers) } } @@ -286,31 +315,31 @@ func TestFollowsDatabase(t *testing.T) { t.Fatal(err) } - followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949) + followers, err := repo.ListRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}) if err != nil { t.Fatal(err) } if len(followers) != 2 { t.Fatalf("movie followers = %+v, want two (the series follow is a different title)", followers) } - followed, err := repo.FollowedTitles(ctx, MediaTypeMovie, []int{949, 950}, a) + followed, err := repo.FollowedRequests(ctx, []string{"movie-949", "series-949", "other"}, b) if err != nil { t.Fatal(err) } - if !followed[949] || followed[950] { - t.Fatalf("followed = %v, want only 949", followed) + if !followed["movie-949"] || len(followed) != 1 { + t.Fatalf("followed = %v, want only movie-949", followed) } - if err := repo.ClearTitleFollowers(ctx, MediaTypeMovie, 949, []Follower{{UserID: a.UserID, ProfileID: a.ProfileID}}); err != nil { + if err := repo.ClearRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}, []Follower{{UserID: a.UserID, ProfileID: a.ProfileID}}); err != nil { t.Fatal(err) } if err := repo.UnfollowTitle(ctx, MediaTypeMovie, 949, b); err != nil { t.Fatal(err) } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949); len(followers) != 0 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}); len(followers) != 0 { t.Fatalf("movie followers after clear and unfollow = %+v, want none", followers) } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeSeries, 949); len(followers) != 1 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{ID: "series-949", MediaType: MediaTypeSeries, TMDBID: 949}); len(followers) != 1 { t.Fatalf("series followers = %+v, want the one untouched follow", followers) } @@ -322,36 +351,36 @@ func TestFollowsDatabase(t *testing.T) { t.Fatalf("follow as account %d: %v", v.UserID, err) } } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949); len(followers) != 2 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}); len(followers) != 2 { t.Fatalf("followers sharing a profile id = %+v, want one per account", followers) } - if followed, _ := repo.FollowedTitles(ctx, MediaTypeMovie, []int{949}, Viewer{UserID: 3, ProfileID: "default"}); followed[949] { + if followed, _ := repo.FollowedRequests(ctx, []string{"movie-949"}, Viewer{UserID: 3, ProfileID: "default"}); followed["movie-949"] { t.Fatal("a third account's default profile sees the others' follow") } if err := repo.UnfollowTitle(ctx, MediaTypeMovie, 949, mine); err != nil { t.Fatal(err) } - followers, err = repo.ListTitleFollowers(ctx, MediaTypeMovie, 949) + followers, err = repo.ListRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}) if err != nil { t.Fatal(err) } if len(followers) != 1 || followers[0] != (Follower{UserID: theirs.UserID, ProfileID: theirs.ProfileID}) { t.Fatalf("followers after one account unfollowed = %+v, want only the other account's", followers) } - if err := repo.ClearTitleFollowers(ctx, MediaTypeMovie, 949, []Follower{{UserID: mine.UserID, ProfileID: mine.ProfileID}}); err != nil { + if err := repo.ClearRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}, []Follower{{UserID: mine.UserID, ProfileID: mine.ProfileID}}); err != nil { t.Fatal(err) } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949); len(followers) != 1 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}); len(followers) != 1 { t.Fatalf("clearing one account's follow removed %+v, want the other account's kept", followers) } if _, err := repo.SetOutcome(ctx, "series-949", guardWithdrawable, OutcomeCancelled, Viewer{}, ""); err != nil { t.Fatal(err) } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeSeries, 949); len(followers) != 0 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{ID: "series-949", MediaType: MediaTypeSeries, TMDBID: 949}); len(followers) != 0 { t.Fatalf("series followers after the withdrawal = %+v, want none", followers) } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949); len(followers) != 1 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}); len(followers) != 1 { t.Fatalf("movie followers after the series withdrawal = %+v, want the one left", followers) } } @@ -378,7 +407,7 @@ func TestClosingRequestForgetsFollowsDatabase(t *testing.T) { if _, err := repo.SetOutcome(ctx, tc.id, guardWithdrawable, tc.outcome, Viewer{}, ""); err != nil { t.Fatal(err) } - if followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, tc.tmdbID); err != nil || len(followers) != 0 { + if followers, err := titleFollowers(ctx, pool, MediaTypeMovie, tc.tmdbID); err != nil || len(followers) != 0 { t.Fatalf("followers once %s committed = %+v, err = %v; want none", tc.id, followers, err) } } @@ -396,11 +425,229 @@ func TestClosingRequestForgetsFollowsDatabase(t *testing.T) { if _, err := repo.SetOutcome(ctx, "failed", StateGuard{Outcomes: []Outcome{OutcomeFailed}}, OutcomeCancelled, Viewer{}, ""); err != nil { t.Fatal(err) } - if followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 973); err != nil || len(followers) != 1 { + if followers, err := titleFollowers(ctx, pool, MediaTypeMovie, 973); err != nil || len(followers) != 1 { t.Fatalf("followers of the open request after closing the failed one = %+v, err = %v; want one", followers, err) } } +// titleFollowers lists every follow on a title, whichever request it waits for. +func titleFollowers(ctx context.Context, pool *pgxpool.Pool, mediaType MediaType, tmdbID int) ([]Follower, error) { + rows, err := pool.Query(ctx, `SELECT user_id, profile_id FROM media_request_follows + WHERE media_type = $1 AND tmdb_id = $2 ORDER BY user_id, profile_id`, mediaType, tmdbID) + if err != nil { + return nil, err + } + return pgx.CollectRows(rows, func(row pgx.CollectableRow) (Follower, error) { + var f Follower + err := row.Scan(&f.UserID, &f.ProfileID) + return f, err + }) +} + +func completeLifecycleRequest(t *testing.T, pool *pgxpool.Pool, id string) { + t.Helper() + if _, err := pool.Exec(t.Context(), `UPDATE media_requests SET status = 'completed', completed_at = now() WHERE id = $1`, id); err != nil { + t.Fatal(err) + } +} + +// A series can have completed requests still waiting for the library beside a +// newer open request for other seasons. Each request's notification goes to +// the follows made while it was open; clearing them spares a profile that +// followed again since, and declining the open request keeps the others. +func TestRequestFollowersDatabase(t *testing.T) { + repo, pool := lifecycleTestRepository(t) + ctx := t.Context() + a := Viewer{UserID: 1, ProfileID: "profile-a"} + b := Viewer{UserID: 2, ProfileID: "profile-b"} + follow := func(v Viewer) { + t.Helper() + if err := repo.FollowTitle(ctx, MediaTypeMovie, 975, v); err != nil { + t.Fatal(err) + } + } + list := func(id string) []Follower { + t.Helper() + followers, err := repo.ListRequestFollowers(ctx, Request{ID: id, MediaType: MediaTypeMovie, TMDBID: 975}) + if err != nil { + t.Fatal(err) + } + return followers + } + insertLifecycleRequest(t, repo, "older", 5, 975, StatusPending) + follow(a) + completeLifecycleRequest(t, pool, "older") + insertLifecycleRequest(t, repo, "newer", 6, 975, StatusPending) + follow(b) + completeLifecycleRequest(t, pool, "newer") + if got := list("older"); !slices.Equal(got, []Follower{{UserID: 1, ProfileID: "profile-a"}}) { + t.Fatalf("older's followers = %+v, want profile-a", got) + } + if got := list("newer"); !slices.Equal(got, []Follower{{UserID: 2, ProfileID: "profile-b"}}) { + t.Fatalf("newer's followers = %+v, want profile-b", got) + } + + // profile-b unfollows and follows a third request while newer's + // notification is going out; newer's clear leaves the new follow. + if err := repo.UnfollowTitle(ctx, MediaTypeMovie, 975, b); err != nil { + t.Fatal(err) + } + insertLifecycleRequest(t, repo, "third", 7, 975, StatusPending) + follow(b) + if err := repo.ClearRequestFollowers(ctx, Request{ID: "newer", MediaType: MediaTypeMovie, TMDBID: 975}, []Follower{{UserID: 2, ProfileID: "profile-b"}}); err != nil { + t.Fatal(err) + } + if got := list("third"); len(got) != 1 { + t.Fatalf("third's followers after newer's clear = %+v, want profile-b kept", got) + } + // profile-a, still waiting for older, follows third too. + follow(a) + if got := list("third"); len(got) != 2 { + t.Fatalf("third's followers = %+v, want profile-a and profile-b", got) + } + if followed, err := repo.FollowedRequests(ctx, []string{"third"}, a); err != nil || !followed["third"] { + t.Fatalf("profile-a follows third = %v, err = %v; want true", followed, err) + } + + if _, err := repo.SetOutcome(ctx, "third", guardWithdrawable, OutcomeDeclined, Viewer{}, ""); err != nil { + t.Fatal(err) + } + if got, err := titleFollowers(ctx, pool, MediaTypeMovie, 975); err != nil || !slices.Equal(got, []Follower{{UserID: 1, ProfileID: "profile-a"}}) { + t.Fatalf("follows after declining third = %+v, err = %v; want older's kept", got, err) + } +} + +// A follow that commits while a completion is under way belongs to the +// request it read, even though the completion's timestamp, taken when its +// transaction began, is earlier than the follow's. +func TestFollowDuringCompletionDatabase(t *testing.T) { + repo, pool := lifecycleTestRepository(t) + ctx := t.Context() + insertLifecycleRequest(t, repo, "req", 5, 976, StatusPending) + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer func() { _ = tx.Rollback(ctx) }() + if _, err := tx.Exec(ctx, `SELECT now()`); err != nil { + t.Fatal(err) + } + if err := repo.FollowTitle(ctx, MediaTypeMovie, 976, Viewer{UserID: 1, ProfileID: "profile-a"}); err != nil { + t.Fatal(err) + } + if _, err := tx.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() WHERE id = 'req'`); err != nil { + t.Fatal(err) + } + if err := tx.Commit(ctx); err != nil { + t.Fatal(err) + } + var inverted bool + if err := pool.QueryRow(ctx, `SELECT f.created_at > r.completed_at FROM media_request_follows f + JOIN media_requests r ON r.id = 'req' WHERE f.tmdb_id = 976`).Scan(&inverted); err != nil || !inverted { + t.Fatalf("follow made after the completion's timestamp = %v, err = %v; the race was not reproduced", inverted, err) + } + followers, err := repo.ListRequestFollowers(ctx, Request{ID: "req", MediaType: MediaTypeMovie, TMDBID: 976}) + if err != nil || len(followers) != 1 { + t.Fatalf("followers = %+v, err = %v; want the follow made during the completion", followers, err) + } +} + +// A follow survives its request failing: the title's next request takes it, +// including when the requester's new request replaces the failed one. +func TestNewRequestAdoptsFollowsOfFailedRequestDatabase(t *testing.T) { + repo, pool := lifecycleTestRepository(t) + ctx := t.Context() + for _, tc := range []struct { + tmdbID int + replace bool + }{{977, false}, {978, true}} { + insertLifecycleRequest(t, repo, fmt.Sprintf("failed-%d", tc.tmdbID), 5, tc.tmdbID, StatusApproved) + if err := repo.FollowTitle(ctx, MediaTypeMovie, tc.tmdbID, Viewer{UserID: 1, ProfileID: "profile-a"}); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `UPDATE media_requests SET outcome = 'failed' WHERE id = $1`, fmt.Sprintf("failed-%d", tc.tmdbID)); err != nil { + t.Fatal(err) + } + retry := fmt.Sprintf("retry-%d", tc.tmdbID) + if _, err := repo.CreateRequest(ctx, CreateRequestRecord{ + ID: retry, + Input: CreateRequestInput{MediaType: MediaTypeMovie, TMDBID: tc.tmdbID, Title: "Retry"}, + Status: StatusPending, + Outcome: OutcomeActive, + Requester: Viewer{UserID: 5, ProfileID: "profile"}, + ReplaceFailed: tc.replace, + }); err != nil { + t.Fatal(err) + } + followers, err := repo.ListRequestFollowers(ctx, Request{ID: retry, MediaType: MediaTypeMovie, TMDBID: tc.tmdbID}) + if err != nil || len(followers) != 1 { + t.Fatalf("replace=%v: the new request's followers = %+v, err = %v; want the failed request's follow", tc.replace, followers, err) + } + } +} + +// An unfollow that runs while a new request takes over a failed request's +// follows waits for it and removes the moved follow, rather than returning +// with the profile still following the new request. +func TestUnfollowDuringAdoptionDatabase(t *testing.T) { + repo, pool := lifecycleTestRepository(t) + ctx := t.Context() + follower := Viewer{UserID: 1, ProfileID: "profile-a"} + insertLifecycleRequest(t, repo, "failed", 5, 979, StatusApproved) + if err := repo.FollowTitle(ctx, MediaTypeMovie, 979, follower); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `UPDATE media_requests SET outcome = 'failed' WHERE id = 'failed'`); err != nil { + t.Fatal(err) + } + + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer func() { _ = tx.Rollback(context.Background()) }() + var adopter int + if err := tx.QueryRow(ctx, `SELECT pg_backend_pid()`).Scan(&adopter); err != nil { + t.Fatal(err) + } + retry, err := repo.insertRequest(ctx, tx, CreateRequestRecord{ + ID: "retry", + Input: CreateRequestInput{MediaType: MediaTypeMovie, TMDBID: 979, Title: "Retry"}, + Requester: Viewer{UserID: 6, ProfileID: "profile"}, + }, StatusPending, OutcomeActive, time.Now(), nil) + if err != nil { + t.Fatal(err) + } + if err := adoptTitleFollows(ctx, tx, retry); err != nil { + t.Fatal(err) + } + + unfollowed := make(chan error, 1) + go func() { unfollowed <- repo.UnfollowTitle(ctx, MediaTypeMovie, 979, follower) }() + for blocked := false; !blocked; { + select { + case err := <-unfollowed: + t.Fatalf("unfollow finished while the adoption was open: err = %v, want it to wait", err) + default: + } + if err := pool.QueryRow(ctx, `SELECT EXISTS (SELECT 1 FROM pg_stat_activity WHERE $1 = ANY(pg_blocking_pids(pid)))`, adopter).Scan(&blocked); err != nil { + t.Fatal(err) + } + if !blocked { + time.Sleep(5 * time.Millisecond) + } + } + if err := tx.Commit(ctx); err != nil { + t.Fatal(err) + } + if err := <-unfollowed; err != nil { + t.Fatal(err) + } + if followers, err := titleFollowers(ctx, pool, MediaTypeMovie, 979); err != nil || len(followers) != 0 { + t.Fatalf("follows after the unfollow = %+v, err = %v; want none", followers, err) + } +} + // A follow racing a withdrawal must not outlive it: the follow waits for the // withdrawal to commit, sees the request closed, and inserts nothing, so the // follow cleanup that ran with the withdrawal leaves no stray follower behind. @@ -451,7 +698,7 @@ func TestFollowWaitsForConcurrentWithdrawalDatabase(t *testing.T) { if err := <-followed; !errors.Is(err, ErrNotRequested) { t.Fatalf("follow after the withdrawal: err = %v, want ErrNotRequested", err) } - if followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 959); err != nil || len(followers) != 0 { + if followers, err := titleFollowers(ctx, pool, MediaTypeMovie, 959); err != nil || len(followers) != 0 { t.Fatalf("followers after the withdrawal = %+v, err = %v; want none", followers, err) } } diff --git a/internal/requests/notify.go b/internal/requests/notify.go index b045d99d3f..5f40eca612 100644 --- a/internal/requests/notify.go +++ b/internal/requests/notify.go @@ -129,7 +129,7 @@ func (s *Service) notifyFulfilledPending(ctx context.Context) { s.markNotifyChecked(ctx, req.ID) continue } - followers, err := s.store.ListTitleFollowers(ctx, req.MediaType, req.TMDBID) + followers, err := s.store.ListRequestFollowers(ctx, *req) if err != nil { slog.WarnContext(ctx, "request fulfill-notify: list followers failed", "component", "requests", "request_id", req.ID, "err", err) @@ -141,13 +141,12 @@ func (s *Service) notifyFulfilledPending(ctx context.Context) { "request_id", req.ID, "err", err) continue } - // The followers have been told. Only the listed rows are cleared: - // FollowTitle needs an open request, so none can have been added - // since this request completed. They are cleared before the request - // is stamped, so a failed clear leaves it unstamped and the next run - // retries it (the deliveries dedupe) instead of leaving follows - // behind for a later request of the title. - if err := s.store.ClearTitleFollowers(ctx, req.MediaType, req.TMDBID, followers); err != nil { + // The followers have been told. Only the listed rows are cleared, + // leaving the follows made for other requests of the title. They are + // cleared before the request is stamped, so a failed clear leaves it + // unstamped and the next run retries it (the deliveries dedupe) + // instead of leaving follows behind for a later request of the title. + if err := s.store.ClearRequestFollowers(ctx, *req, followers); err != nil { slog.WarnContext(ctx, "request fulfill-notify: clear followers failed", "component", "requests", "request_id", req.ID, "err", err) continue diff --git a/internal/requests/repository.go b/internal/requests/repository.go index b782eaf56c..259caa16a1 100644 --- a/internal/requests/repository.go +++ b/internal/requests/repository.go @@ -237,20 +237,6 @@ func (r *Repository) CreateRequest(ctx context.Context, input CreateRequestRecor return nil, fmt.Errorf("acquire request quota lock: %w", err) } } - if input.ReplaceFailed { - // Only the requester's own rows: other accounts' failed requests for - // the title are their history. - if _, err := tx.Exec(ctx, ` - DELETE FROM media_requests - WHERE requested_by_user_id = $1 - AND media_type = $2 - AND provider = 'tmdb' - AND tmdb_id = $3 - AND outcome = 'failed' - `, input.Requester.UserID, input.Input.MediaType, input.Input.TMDBID); err != nil { - return nil, fmt.Errorf("replace failed requests: %w", err) - } - } if input.Quota != nil { var count int if err := tx.QueryRow(ctx, ` @@ -292,6 +278,25 @@ func (r *Repository) CreateRequest(ctx context.Context, input CreateRequestRecor } return nil, err } + if err := adoptTitleFollows(ctx, tx, req); err != nil { + return nil, err + } + // Failed requests do not count against the quota, so the ones this + // replaces go after the insert, once their follows have moved to it. + if input.ReplaceFailed { + // Only the requester's own rows: other accounts' failed requests for + // the title are their history. + if _, err := tx.Exec(ctx, ` + DELETE FROM media_requests + WHERE requested_by_user_id = $1 + AND media_type = $2 + AND provider = 'tmdb' + AND tmdb_id = $3 + AND outcome = 'failed' + `, input.Requester.UserID, input.Input.MediaType, input.Input.TMDBID); err != nil { + return nil, fmt.Errorf("replace failed requests: %w", err) + } + } if err := r.recordEvent(ctx, tx, req.ID, "created", input.Requester, ""); err != nil { return nil, err } @@ -1057,11 +1062,7 @@ func (r *Repository) SaveIntegrationWithDefaults(ctx context.Context, in Integra if err == nil { err = defaultFirstServer(ctx, tx, out) } - } else if !standard { - if err = ensureTierKeptUnderAdvanced(ctx, tx, in); err == nil { - out, err = r.updateIntegration(ctx, tx, in) - } - } else { + } else if err = ensureRoutesStillFit(ctx, tx, in, !standard); err == nil { out, err = r.updateIntegration(ctx, tx, in) } if err == nil && standard { diff --git a/internal/requests/routes_admin.go b/internal/requests/routes_admin.go index 9dd84bbc33..c0db94ca2c 100644 --- a/internal/requests/routes_admin.go +++ b/internal/requests/routes_admin.go @@ -398,6 +398,30 @@ func (s *Service) ensureRoutesKeepServerKind(ctx context.Context, in Integration } checkTier = settings.Mode != RoutingStandard } + fields := routeKindConflicts(in, routes) + if checkTier { + current, err := s.store.GetIntegration(ctx, in.ID) + if err != nil && !errors.Is(err, ErrNotFound) { + return err + } + if current == nil || is4KServer(*current) != is4KServer(in) { + if msg := tierConflict(in, routes); msg != "" { + fields["plugin_config."+configIs4K] = msg + } + } + } + if len(fields) == 0 { + return nil + } + return &ValidationError{FieldErrors: fields} +} + +// routeKindConflicts explains, by field, why the routes sending to a server +// keep its type and media types where they are: they send it requests it +// would no longer take. It is empty when none do. The repository checks again +// under the server's row lock (ensureRoutesStillFit), since a route save can +// commit between this check and the server save. +func routeKindConflicts(in Integration, routes []Route) map[string]string { var wrongKind, unsupported []string for _, r := range routes { if r.HD.IntegrationID != in.ID && r.UHD.IntegrationID != in.ID { @@ -419,21 +443,7 @@ func (s *Service) ensureRoutesKeepServerKind(ctx context.Context, in Integration fields["supported_media_types"] = "Routing sends a media type this server would no longer take to it (" + strings.Join(unsupported, ", ") + "); change those routes first." } - if checkTier { - current, err := s.store.GetIntegration(ctx, in.ID) - if err != nil && !errors.Is(err, ErrNotFound) { - return err - } - if current == nil || is4KServer(*current) != is4KServer(in) { - if msg := tierConflict(in, routes); msg != "" { - fields["plugin_config."+configIs4K] = msg - } - } - } - if len(fields) == 0 { - return nil - } - return &ValidationError{FieldErrors: fields} + return fields } // tierConflict explains why routes keep a server's 4K switch where it is: they @@ -536,6 +546,22 @@ func tierMismatch(in Integration, uhd bool) string { return "" } +// destinationMismatch explains why a server can't take a route's HD or 4K +// version: it is the other type of server, doesn't take the media type, or is +// on the other side of the 4K switch. It is empty when the server fits. +func destinationMismatch(in Integration, mediaType MediaType, uhd bool) string { + if kind, _ := in.PluginConfig[configServiceKind].(string); kind != "" { + wantKind := map[MediaType]string{MediaTypeMovie: kindRadarr, MediaTypeSeries: kindSonarr}[mediaType] + if wantKind != "" && kind != wantKind { + return fmt.Sprintf("%s is a %s server; %s need %s.", in.Name, kindLabel(kind), mediaTypePlural(mediaType), kindLabel(wantKind)) + } + } + if mediaType != "" && !integrationSupportsMediaType(in, mediaType) { + return fmt.Sprintf("%s does not take %s.", in.Name, mediaTypePlural(mediaType)) + } + return tierMismatch(in, uhd) +} + func validateDestination(field string, dest *RouteDestination, mediaType MediaType, integrations []Integration, fields map[string]string) { dest.IntegrationID = strings.TrimSpace(dest.IntegrationID) if dest.IntegrationID == "" { @@ -552,19 +578,8 @@ func validateDestination(field string, dest *RouteDestination, mediaType MediaTy fields[field+".integration_id"] = "That server no longer exists." return } - if kind, _ := in.PluginConfig[configServiceKind].(string); kind != "" { - wantKind := map[MediaType]string{MediaTypeMovie: kindRadarr, MediaTypeSeries: kindSonarr}[mediaType] - if wantKind != "" && kind != wantKind { - fields[field+".integration_id"] = fmt.Sprintf("%s is a %s server; %s need %s.", in.Name, kindLabel(kind), mediaTypePlural(mediaType), kindLabel(wantKind)) - } - } - if mediaType != "" && fields[field+".integration_id"] == "" && !integrationSupportsMediaType(*in, mediaType) { - fields[field+".integration_id"] = fmt.Sprintf("%s does not take %s.", in.Name, mediaTypePlural(mediaType)) - } - if fields[field+".integration_id"] == "" { - if msg := tierMismatch(*in, field == fieldUHD); msg != "" { - fields[field+".integration_id"] = msg - } + if msg := destinationMismatch(*in, mediaType, field == fieldUHD); msg != "" { + fields[field+".integration_id"] = msg } for _, key := range routingOwnedConfigKeys { if _, ok := dest.Overrides[key]; ok { @@ -732,7 +747,7 @@ func (r *Repository) SaveRouteConditional(ctx context.Context, route Route, expe allowMissing := expected == 0 || route.IsFallback // Servers before the route row, the order a server save that turns // Advanced on takes them in. - if err := ensureDestinationTiers(ctx, tx, route); err != nil { + if err := ensureDestinationsFit(ctx, tx, route); err != nil { return nil, err } if err := lockRevision(ctx, tx, `SELECT revision FROM request_routes WHERE id = $1 FOR UPDATE`, []any{route.ID}, expected, allowMissing); err != nil { @@ -762,11 +777,11 @@ func (r *Repository) SaveRouteConditional(ctx context.Context, route Route, expe return &saved, tx.Commit(ctx) } -// ensureDestinationTiers checks the route's servers against the tier rule -// again inside the save, holding their rows FOR SHARE: a 4K switch change -// committed since validateRoute is seen here, and one still in flight waits -// for this save and then finds the route (ensureTierKeptUnderAdvanced). -func ensureDestinationTiers(ctx context.Context, tx pgx.Tx, route Route) error { +// ensureDestinationsFit checks the route's servers again inside the save, +// holding their rows FOR SHARE: a change to a server's type, media types or 4K +// switch committed since validateRoute is seen here, and one still in flight +// waits for this save and then finds the route (ensureRoutesStillFit). +func ensureDestinationsFit(ctx context.Context, tx pgx.Tx, route Route) error { ids := make([]string, 0, 2) for _, id := range []string{route.HD.IntegrationID, route.UHD.IntegrationID} { if id != "" { @@ -776,7 +791,7 @@ func ensureDestinationTiers(ctx context.Context, tx pgx.Tx, route Route) error { if len(ids) == 0 { return nil } - rows, err := tx.Query(ctx, `SELECT id, name, plugin_config FROM request_integrations WHERE id = ANY($1) ORDER BY id FOR SHARE`, ids) + rows, err := tx.Query(ctx, `SELECT id, name, supported_media_types, plugin_config FROM request_integrations WHERE id = ANY($1) ORDER BY id FOR SHARE`, ids) if err != nil { return fmt.Errorf("lock route servers: %w", err) } @@ -785,7 +800,7 @@ func ensureDestinationTiers(ctx context.Context, tx pgx.Tx, route Route) error { for rows.Next() { var in Integration var raw []byte - if err := rows.Scan(&in.ID, &in.Name, &raw); err != nil { + if err := rows.Scan(&in.ID, &in.Name, &in.SupportedMediaTypes, &raw); err != nil { return fmt.Errorf("scan route server: %w", err) } if len(raw) > 0 { @@ -799,7 +814,7 @@ func ensureDestinationTiers(ctx context.Context, tx pgx.Tx, route Route) error { dest = route.UHD } if dest.IntegrationID == in.ID { - if msg := tierMismatch(in, uhd); msg != "" { + if msg := destinationMismatch(in, route.MediaType, uhd); msg != "" { fields[field+".integration_id"] = msg } } diff --git a/internal/requests/routes_admin_test.go b/internal/requests/routes_admin_test.go index 4ad754a612..213e95eec2 100644 --- a/internal/requests/routes_admin_test.go +++ b/internal/requests/routes_admin_test.go @@ -586,3 +586,26 @@ func TestSaveRouteRechecksTierDatabase(t *testing.T) { t.Fatalf("4K only: %v", err) } } + +// A server's type changed after the editor validated the route is caught by +// the same locked check: movies can't go to what is now a Sonarr. +func TestSaveRouteRechecksServerKindDatabase(t *testing.T) { + ctx := t.Context() + repo, pool := routingModeRepository(t) + in := arrServer("radarr", kindRadarr, nil) + in.APIKeyRef = "" + if _, err := repo.SaveIntegrationWithDefaults(ctx, in, true); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `UPDATE request_integrations + SET plugin_config = plugin_config || '{"service_kind": "sonarr"}', supported_media_types = '{series}' + WHERE id = 'radarr'`); err != nil { + t.Fatal(err) + } + route := Route{ID: "anime", MediaType: MediaTypeMovie, Name: "Anime", Enabled: true, + Conditions: RouteConditions{Anime: boolPtr(true)}, HD: RouteDestination{IntegrationID: "radarr"}, SkipUHD: true} + _, err := repo.SaveRouteConditional(ctx, route, 0) + if msg := fieldErrors(t, err)["hd.integration_id"]; !strings.Contains(msg, "is a Sonarr server") { + t.Fatalf("hd.integration_id = %q, want the server type refused", msg) + } +} diff --git a/internal/requests/routing_mode.go b/internal/requests/routing_mode.go index eb822e1634..ba40b5d387 100644 --- a/internal/requests/routing_mode.go +++ b/internal/requests/routing_mode.go @@ -442,14 +442,15 @@ func (r *Repository) standardBeforeSave(ctx context.Context, tx pgx.Tx) (layout return layout, true, nil } -// ensureTierKeptUnderAdvanced refuses, under the routing-mode lock taken by -// standardBeforeSave, a server update whose 4K switch leaves an active route -// sending it the other version. The service checks the same before the save; -// this catches Advanced turned on in between. -func ensureTierKeptUnderAdvanced(ctx context.Context, tx pgx.Tx, in Integration) error { +// ensureRoutesStillFit refuses, under the routing-mode lock taken by +// standardBeforeSave and the server's row lock, a server update that leaves a +// route sending it requests it would no longer take: the other media type, or +// under Advanced the other version. The service checks the same before the +// save; this catches a route saved, or Advanced turned on, in between. +func ensureRoutesStillFit(ctx context.Context, tx pgx.Tx, in Integration, advanced bool) error { var raw []byte // FOR UPDATE before reading the routes: a route save holds its servers - // FOR SHARE (ensureDestinationTiers), so one in flight commits first and + // FOR SHARE (ensureDestinationsFit), so one in flight commits first and // its route is read below. err := tx.QueryRow(ctx, `SELECT plugin_config FROM request_integrations WHERE id = $1 FOR UPDATE`, in.ID).Scan(&raw) if errors.Is(err, pgx.ErrNoRows) { @@ -464,15 +465,18 @@ func ensureTierKeptUnderAdvanced(ctx context.Context, tx pgx.Tx, in Integration) return fmt.Errorf("decode request integration %s config: %w", in.ID, err) } } - if is4KServer(stored) == is4KServer(in) { - return nil - } routes, err := listRoutes(ctx, tx) if err != nil { return err } - if msg := tierConflict(in, routes); msg != "" { - return &ValidationError{FieldErrors: map[string]string{"plugin_config." + configIs4K: msg}} + fields := routeKindConflicts(in, routes) + if advanced && is4KServer(stored) != is4KServer(in) { + if msg := tierConflict(in, routes); msg != "" { + fields["plugin_config."+configIs4K] = msg + } + } + if len(fields) > 0 { + return &ValidationError{FieldErrors: fields} } return nil } diff --git a/internal/requests/routing_mode_test.go b/internal/requests/routing_mode_test.go index aa7b9fcda5..3be91bd76f 100644 --- a/internal/requests/routing_mode_test.go +++ b/internal/requests/routing_mode_test.go @@ -823,3 +823,34 @@ func TestSaveRechecksTierUnderAdvancedDatabase(t *testing.T) { t.Fatalf("rename the 4K server an older rule sends HD to: %v", err) } } + +// A route saved after the service checked the server's routes is found by the +// save's locked check, under Standard too: the server can't become a Sonarr +// while a rule sends it movies. +func TestSaveRechecksServerKindUnderLockDatabase(t *testing.T) { + ctx := t.Context() + repo, pool := routingModeRepository(t) + in := arrServer("radarr", kindRadarr, nil) + in.APIKeyRef = "" + if _, err := repo.SaveIntegrationWithDefaults(ctx, in, true); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `INSERT INTO request_routes (id, media_type, position, name, enabled, conditions, hd_integration_id) + VALUES ('anime', 'movie', 0, 'Anime', true, '{"anime":true}', 'radarr')`); err != nil { + t.Fatal(err) + } + stored, err := repo.GetIntegration(ctx, "radarr") + if err != nil { + t.Fatal(err) + } + sonarr := *stored + sonarr.PluginConfig["service_kind"] = kindSonarr + sonarr.SupportedMediaTypes = []string{string(MediaTypeSeries)} + _, err = repo.UpdateIntegrationConditional(ctx, sonarr, sonarr.Revision) + if fields := fieldErrors(t, err); !strings.Contains(fields["plugin_config.service_kind"], "Anime") || !strings.Contains(fields["supported_media_types"], "Anime") { + t.Fatalf("conditional save making radarr a Sonarr: %v, want the type and media type refused naming Anime", fields) + } + if _, err := repo.SaveIntegrationWithDefaults(ctx, sonarr, false); !strings.Contains(fieldErrors(t, err)["plugin_config.service_kind"], "Anime") { + t.Fatalf("save making radarr a Sonarr: %v, want a service_kind field error", err) + } +} diff --git a/internal/requests/service_test.go b/internal/requests/service_test.go index 31fbe773c1..ca1c2bc6f6 100644 --- a/internal/requests/service_test.go +++ b/internal/requests/service_test.go @@ -1991,7 +1991,8 @@ type fakeStore struct { notified []string reconciled []string follows map[string]Follower // key: media_type/tmdb_id/user_id/profile_id - clearErr error // returned by ClearTitleFollowers when set + followFor map[string]string // follow key -> the request it waits for + clearErr error // returned by ClearRequestFollowers when set routes []Route factsSet map[string]RoutingFacts groupLimits map[int64]*GroupLimit @@ -2386,9 +2387,8 @@ func (f *fakeStore) SetOutcome(_ context.Context, id string, from StateGuard, ou req.LastError = message if outcome == OutcomeDeclined || outcome == OutcomeCancelled { req.OutcomeReason = message - prefix := fmt.Sprintf("%s/%d/", req.MediaType, req.TMDBID) for key := range f.follows { - if strings.HasPrefix(key, prefix) { + if f.followFor[key] == req.ID { delete(f.follows, key) } } @@ -2598,8 +2598,9 @@ func (f *fakeStore) DeleteIntegration(_ context.Context, id string) error { return ErrNotFound } -func followKey(mediaType MediaType, tmdbID int, userID int, profileID string) string { - return fmt.Sprintf("%s/%d/%d/%s", mediaType, tmdbID, userID, profileID) +// followKey names one profile's follow of one request of a title. +func followKey(mediaType MediaType, tmdbID int, userID int, profileID, requestID string) string { + return fmt.Sprintf("%s/%d/%d/%s/%s", mediaType, tmdbID, userID, profileID, requestID) } func (f *fakeStore) FollowTitle(_ context.Context, mediaType MediaType, tmdbID int, viewer Viewer) error { @@ -2620,39 +2621,71 @@ func (f *fakeStore) seedFollow(mediaType MediaType, tmdbID int, viewer Viewer) { f.seedFollowLocked(mediaType, tmdbID, viewer) } +// seedFollowLocked records a follow for the title's open request, or else +// for the title's request in f.requests with the lowest id. func (f *fakeStore) seedFollowLocked(mediaType MediaType, tmdbID int, viewer Viewer) { + requestID := "" + if req := f.active[mediaType][tmdbID]; req != nil { + requestID = req.ID + } else { + for id, req := range f.requests { + if req.MediaType == mediaType && req.TMDBID == tmdbID && (requestID == "" || id < requestID) { + requestID = id + } + } + } + f.seedFollowForLocked(mediaType, tmdbID, viewer, requestID) +} + +// seedFollowFor records a follow waiting for a given request. +func (f *fakeStore) seedFollowFor(mediaType MediaType, tmdbID int, viewer Viewer, requestID string) { + f.mu.Lock() + defer f.mu.Unlock() + f.seedFollowForLocked(mediaType, tmdbID, viewer, requestID) +} + +func (f *fakeStore) seedFollowForLocked(mediaType MediaType, tmdbID int, viewer Viewer, requestID string) { if f.follows == nil { f.follows = map[string]Follower{} } - f.follows[followKey(mediaType, tmdbID, viewer.UserID, viewer.ProfileID)] = Follower{UserID: viewer.UserID, ProfileID: viewer.ProfileID} + if f.followFor == nil { + f.followFor = map[string]string{} + } + key := followKey(mediaType, tmdbID, viewer.UserID, viewer.ProfileID, requestID) + f.follows[key] = Follower{UserID: viewer.UserID, ProfileID: viewer.ProfileID} + f.followFor[key] = requestID } func (f *fakeStore) UnfollowTitle(_ context.Context, mediaType MediaType, tmdbID int, viewer Viewer) error { f.mu.Lock() defer f.mu.Unlock() - delete(f.follows, followKey(mediaType, tmdbID, viewer.UserID, viewer.ProfileID)) + prefix := fmt.Sprintf("%s/%d/%d/%s/", mediaType, tmdbID, viewer.UserID, viewer.ProfileID) + for key := range f.follows { + if strings.HasPrefix(key, prefix) { + delete(f.follows, key) + } + } return nil } -func (f *fakeStore) FollowedTitles(_ context.Context, mediaType MediaType, tmdbIDs []int, viewer Viewer) (map[int]bool, error) { +func (f *fakeStore) FollowedRequests(_ context.Context, requestIDs []string, viewer Viewer) (map[string]bool, error) { f.mu.Lock() defer f.mu.Unlock() - out := map[int]bool{} - for _, id := range tmdbIDs { - if _, ok := f.follows[followKey(mediaType, id, viewer.UserID, viewer.ProfileID)]; ok { - out[id] = true + out := map[string]bool{} + for key, follower := range f.follows { + if follower.UserID == viewer.UserID && follower.ProfileID == viewer.ProfileID && slices.Contains(requestIDs, f.followFor[key]) { + out[f.followFor[key]] = true } } return out, nil } -func (f *fakeStore) ListTitleFollowers(_ context.Context, mediaType MediaType, tmdbID int) ([]Follower, error) { +func (f *fakeStore) ListRequestFollowers(_ context.Context, req Request) ([]Follower, error) { f.mu.Lock() defer f.mu.Unlock() - prefix := fmt.Sprintf("%s/%d/", mediaType, tmdbID) var out []Follower for key, follower := range f.follows { - if strings.HasPrefix(key, prefix) { + if f.followFor[key] == req.ID { out = append(out, follower) } } @@ -2665,14 +2698,28 @@ func (f *fakeStore) ListTitleFollowers(_ context.Context, mediaType MediaType, t return out, nil } -func (f *fakeStore) ClearTitleFollowers(_ context.Context, mediaType MediaType, tmdbID int, followers []Follower) error { +// titleFollowers lists every follow on a title, whichever request it waits for. +func (f *fakeStore) titleFollowers(mediaType MediaType, tmdbID int) ([]Follower, error) { + f.mu.Lock() + defer f.mu.Unlock() + prefix := fmt.Sprintf("%s/%d/", mediaType, tmdbID) + var out []Follower + for key, follower := range f.follows { + if strings.HasPrefix(key, prefix) { + out = append(out, follower) + } + } + return out, nil +} + +func (f *fakeStore) ClearRequestFollowers(_ context.Context, req Request, followers []Follower) error { f.mu.Lock() defer f.mu.Unlock() if f.clearErr != nil { return f.clearErr } for _, follower := range followers { - delete(f.follows, followKey(mediaType, tmdbID, follower.UserID, follower.ProfileID)) + delete(f.follows, followKey(req.MediaType, req.TMDBID, follower.UserID, follower.ProfileID, req.ID)) } return nil } diff --git a/internal/requests/store.go b/internal/requests/store.go index 395fb71abf..ca6ccb0c1d 100644 --- a/internal/requests/store.go +++ b/internal/requests/store.go @@ -80,15 +80,16 @@ type Store interface { // RecomputeStatus re-derives an approved request's status and outcome from // its targets, for a submission that found nothing left to send. RecomputeStatus(ctx context.Context, id string, actor Viewer) (*Request, error) - // FollowTitle, UnfollowTitle and FollowedTitles manage a profile's follows - // on titles, keyed by account and profile; ListTitleFollowers and ClearTitleFollowers serve the - // fulfilled notification. All are idempotent. FollowTitle answers - // ErrNotRequested when the title has no open request. + // FollowTitle, UnfollowTitle and FollowedRequests manage a profile's + // follows, keyed by account, profile and request; ListRequestFollowers and + // ClearRequestFollowers serve a request's fulfilled notification. All are + // idempotent. FollowTitle answers ErrNotRequested when the title has no + // open request. FollowTitle(ctx context.Context, mediaType MediaType, tmdbID int, viewer Viewer) error UnfollowTitle(ctx context.Context, mediaType MediaType, tmdbID int, viewer Viewer) error - FollowedTitles(ctx context.Context, mediaType MediaType, tmdbIDs []int, viewer Viewer) (map[int]bool, error) - ListTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int) ([]Follower, error) - ClearTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int, followers []Follower) error + FollowedRequests(ctx context.Context, requestIDs []string, viewer Viewer) (map[string]bool, error) + ListRequestFollowers(ctx context.Context, req Request) ([]Follower, error) + ClearRequestFollowers(ctx context.Context, req Request, followers []Follower) error // ListRoutes returns every routing rule, in no particular order; // decideRoutes orders them. ListRoutes(ctx context.Context) ([]Route, error) diff --git a/migrations/sql/20260929001953_request_follows_request_key.sql b/migrations/sql/20260929001953_request_follows_request_key.sql new file mode 100644 index 0000000000..1d18fc09af --- /dev/null +++ b/migrations/sql/20260929001953_request_follows_request_key.sql @@ -0,0 +1,87 @@ +-- +goose Up +-- +goose StatementBegin +-- A series can have completed requests still waiting for the library beside a +-- newer open request for other seasons, so a follow belongs to the request +-- that was open when it was made, and a request's notification goes to its +-- own follows. A profile can follow each of a title's requests. +ALTER TABLE public.media_request_follows ADD COLUMN request_id text; + +-- A title has one open request at a time, so the request open when a follow +-- was made is the title's latest request created before it, provided that +-- request could still have been open then: still active, or completed no +-- earlier than the follow. A completed request that has not notified also +-- keeps a follow stamped just after its completion when the title has no +-- request since: a completion's timestamp is taken when its transaction +-- begins, so a follow that committed during it can look later, and a deleted +-- request always has a replacement created after the follow. Otherwise that +-- request had closed (a failed request, or one its requester replaced and +-- deleted), and the follow goes to the title's first request since that still +-- has a notification to send, as a new request takes such follows over. With +-- none, it stays with the title's failed request, the one it was made for or +-- else the latest, for the title's next request to take over. +UPDATE public.media_request_follows f +SET request_id = CASE + WHEN made_for.outcome = 'active' + AND (made_for.status <> 'completed' OR made_for.completed_at >= f.created_at + OR (made_for.fulfilled_notified_at IS NULL AND NOT EXISTS ( + SELECT 1 FROM public.media_requests r + WHERE r.media_type = f.media_type AND r.provider = 'tmdb' AND r.tmdb_id = f.tmdb_id + AND r.created_at > f.created_at))) + THEN made_for.id + ELSE coalesce( + (SELECT r.id FROM public.media_requests r + WHERE r.media_type = f.media_type AND r.provider = 'tmdb' AND r.tmdb_id = f.tmdb_id + AND r.created_at > f.created_at AND r.outcome = 'active' + AND (r.status <> 'completed' OR r.fulfilled_notified_at IS NULL) + ORDER BY r.created_at, r.id + LIMIT 1), + CASE WHEN made_for.outcome = 'failed' THEN made_for.id END, + (SELECT r.id FROM public.media_requests r + WHERE r.media_type = f.media_type AND r.provider = 'tmdb' AND r.tmdb_id = f.tmdb_id + AND r.outcome = 'failed' + ORDER BY r.created_at DESC, r.id DESC + LIMIT 1)) + END +FROM ( + SELECT DISTINCT ON (f2.media_type, f2.tmdb_id, f2.user_id, f2.profile_id) + f2.media_type, f2.tmdb_id, f2.user_id, f2.profile_id, + r.id, r.outcome, r.status, r.completed_at, r.fulfilled_notified_at + FROM public.media_request_follows f2 + LEFT JOIN public.media_requests r + ON r.media_type = f2.media_type AND r.provider = 'tmdb' AND r.tmdb_id = f2.tmdb_id + AND r.created_at <= f2.created_at + ORDER BY f2.media_type, f2.tmdb_id, f2.user_id, f2.profile_id, r.created_at DESC NULLS LAST, r.id DESC +) made_for +WHERE made_for.media_type = f.media_type AND made_for.tmdb_id = f.tmdb_id + AND made_for.user_id = f.user_id AND made_for.profile_id = f.profile_id; + +-- A follow no request is left to tell has nothing to wait for. +DELETE FROM public.media_request_follows WHERE request_id IS NULL; + +ALTER TABLE public.media_request_follows + ALTER COLUMN request_id SET NOT NULL, + ADD CONSTRAINT media_request_follows_request_fkey FOREIGN KEY (request_id) + REFERENCES public.media_requests (id) ON DELETE CASCADE, + DROP CONSTRAINT media_request_follows_pkey, + ADD PRIMARY KEY (user_id, profile_id, request_id); +DROP INDEX public.media_request_follows_profile_idx; +CREATE INDEX media_request_follows_request_idx ON public.media_request_follows (request_id); +CREATE INDEX media_request_follows_title_idx ON public.media_request_follows (media_type, tmdb_id); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +-- The title-wide key holds one follow per profile and title; keep the earliest. +DELETE FROM public.media_request_follows f +USING public.media_request_follows keep +WHERE keep.media_type = f.media_type AND keep.tmdb_id = f.tmdb_id + AND keep.user_id = f.user_id AND keep.profile_id = f.profile_id + AND (keep.created_at, keep.request_id) < (f.created_at, f.request_id); +DROP INDEX public.media_request_follows_title_idx; +DROP INDEX public.media_request_follows_request_idx; +ALTER TABLE public.media_request_follows + DROP CONSTRAINT media_request_follows_pkey, + ADD PRIMARY KEY (media_type, tmdb_id, user_id, profile_id), + DROP COLUMN request_id; +CREATE INDEX media_request_follows_profile_idx ON public.media_request_follows (user_id, profile_id); +-- +goose StatementEnd