From 1b5b6492175f8a85a62b9aea7ff5b07b1b00b3d2 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Mon, 28 Sep 2026 23:32:32 +0000 Subject: [PATCH 01/10] fix(requests): tell only a request's own followers when it arrives A series can have a completed request still waiting for the library beside a newer open request for other seasons. Follows belong to the title, so the completed request's notification went to, and cleared, follows made for the open request, and declining the open request deleted follows still waiting for the completed one. The notification now goes to the follows made before its request completed, and a decline keeps the follows a completed, not yet notified request of the title is waiting to tell. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/architecture/media-requests.md | 9 ++- internal/requests/follows.go | 22 ++++-- internal/requests/follows_test.go | 106 +++++++++++++++++++++++----- internal/requests/notify.go | 17 +++-- internal/requests/service_test.go | 25 +++++-- internal/requests/store.go | 5 +- 6 files changed, 149 insertions(+), 35 deletions(-) diff --git a/docs/architecture/media-requests.md b/docs/architecture/media-requests.md index 75e2bb1468..b4b131858f 100644 --- a/docs/architecture/media-requests.md +++ b/docs/architecture/media-requests.md @@ -311,8 +311,13 @@ 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 +to the title's followers, marked `follower` so its wording does not say +"your request", and those follows are then cleared. A series can have a +completed request still waiting for the library beside a newer open request +for other seasons, so the notification goes only to the follows made before +its request completed; later follows were made for the open request and wait +for it. Declining that open request leaves the earlier follows for the +completed one. 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/internal/requests/follows.go b/internal/requests/follows.go index 8dc8210f9f..edbc0fbbbb 100644 --- a/internal/requests/follows.go +++ b/internal/requests/follows.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "strings" + "time" ) // Following a title: a profile that finds a title someone else has already @@ -14,6 +15,12 @@ import ( // 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. +// +// A series can have a completed request still waiting for the library and a +// newer open one for other seasons. A follow made after the first completed +// was made for the open one, so the first request's notification goes only to +// the follows made before it completed, and closing the open one leaves those +// follows for the completed request. // Follower is a profile waiting to hear that a requested title is available. type Follower struct { @@ -159,15 +166,21 @@ func (r *Repository) FollowTitle(ctx context.Context, mediaType MediaType, tmdbI // 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. +// followers are still waiting for it. A follow made before a completed request +// of the title completed stays until that request's notification goes out. func forgetTitleFollows(ctx context.Context, exec requestExecutor, closed *Request) error { if _, err := exec.Exec(ctx, ` - DELETE FROM media_request_follows + DELETE FROM media_request_follows f 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) + 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 fulfilled_notified_at IS NULL AND completed_at >= f.created_at) `, closed.MediaType, closed.TMDBID, closed.ID); err != nil { return fmt.Errorf("forget title follows: %w", err) } @@ -207,12 +220,13 @@ func (r *Repository) FollowedTitles(ctx context.Context, mediaType MediaType, tm return out, rows.Err() } -func (r *Repository) ListTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int) ([]Follower, error) { +func (r *Repository) ListTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int, followedBy *time.Time) ([]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 + AND ($3::timestamptz IS NULL OR created_at <= $3) ORDER BY created_at, user_id, profile_id - `, mediaType, tmdbID) + `, mediaType, tmdbID, followedBy) if err != nil { return nil, fmt.Errorf("list title followers: %w", err) } diff --git a/internal/requests/follows_test.go b/internal/requests/follows_test.go index ba9470d675..8a93bcc35b 100644 --- a/internal/requests/follows_test.go +++ b/internal/requests/follows_test.go @@ -3,6 +3,7 @@ package requests import ( "context" "errors" + "slices" "testing" "time" @@ -86,7 +87,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.ListTitleFollowers(context.Background(), MediaTypeMovie, 949, nil) 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 +160,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.ListTitleFollowers(context.Background(), MediaTypeMovie, 949, nil); len(followers) != 0 { t.Fatalf("followers after %s = %+v, want none", tc.name, followers) } }) @@ -213,11 +214,40 @@ 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.ListTitleFollowers(context.Background(), MediaTypeMovie, 42, nil); 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 the follows made +// before it completed; one made since was made for the open request. +func TestNotifyFulfilledLeavesFollowsOfNewerRequest(t *testing.T) { + store := newFakeStore() + completed := time.Now().Add(-time.Hour) + req := completedRequestFixture("req1", 42) + req.CompletedAt = &completed + store.requests["req1"] = req + store.unnotified = []string{"req1"} + early := Viewer{UserID: 3, ProfileID: "early"} + late := Viewer{UserID: 4, ProfileID: "late"} + store.seedFollowAt(MediaTypeMovie, 42, early, completed.Add(-time.Minute)) + store.seedFollowAt(MediaTypeMovie, 42, late, completed.Add(time.Minute)) + 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 the follow made before completion", notifier.followers) + } + left, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42, nil) + if !slices.Equal(left, []Follower{{UserID: 4, ProfileID: "late"}}) { + t.Fatalf("followers left = %+v, want the newer request's follower", left) + } +} + func TestNotifyFulfilledKeepsFollowersWhenDispatchFails(t *testing.T) { store := newFakeStore() store.requests["req1"] = completedRequestFixture("req1", 42) @@ -228,7 +258,7 @@ func TestNotifyFulfilledKeepsFollowersWhenDispatchFails(t *testing.T) { svc.notifyFulfilledPending(context.Background()) - if followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42); len(followers) != 1 { + if followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42, nil); len(followers) != 1 { t.Fatalf("followers after a failed dispatch = %+v, want kept for the retry", followers) } } @@ -254,7 +284,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.ListTitleFollowers(context.Background(), MediaTypeMovie, 42, nil); len(followers) != 0 { t.Fatalf("followers after the retry = %+v, want cleared", followers) } } @@ -286,7 +316,7 @@ func TestFollowsDatabase(t *testing.T) { t.Fatal(err) } - followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949) + followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil) if err != nil { t.Fatal(err) } @@ -307,10 +337,10 @@ func TestFollowsDatabase(t *testing.T) { 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.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil); 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.ListTitleFollowers(ctx, MediaTypeSeries, 949, nil); len(followers) != 1 { t.Fatalf("series followers = %+v, want the one untouched follow", followers) } @@ -322,7 +352,7 @@ 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.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil); 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] { @@ -331,7 +361,7 @@ func TestFollowsDatabase(t *testing.T) { if err := repo.UnfollowTitle(ctx, MediaTypeMovie, 949, mine); err != nil { t.Fatal(err) } - followers, err = repo.ListTitleFollowers(ctx, MediaTypeMovie, 949) + followers, err = repo.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil) if err != nil { t.Fatal(err) } @@ -341,17 +371,17 @@ func TestFollowsDatabase(t *testing.T) { if err := repo.ClearTitleFollowers(ctx, MediaTypeMovie, 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.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil); 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.ListTitleFollowers(ctx, MediaTypeSeries, 949, nil); 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.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil); len(followers) != 1 { t.Fatalf("movie followers after the series withdrawal = %+v, want the one left", followers) } } @@ -378,7 +408,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 := repo.ListTitleFollowers(ctx, MediaTypeMovie, tc.tmdbID, nil); err != nil || len(followers) != 0 { t.Fatalf("followers once %s committed = %+v, err = %v; want none", tc.id, followers, err) } } @@ -396,11 +426,55 @@ 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 := repo.ListTitleFollowers(ctx, MediaTypeMovie, 973, nil); err != nil || len(followers) != 1 { t.Fatalf("followers of the open request after closing the failed one = %+v, err = %v; want one", followers, err) } } +// A completed request still waiting for the library keeps the follows made +// before it completed when a newer request for the title is declined; the +// follows made for the declined request go. +func TestDecliningNewerRequestKeepsCompletedRequestFollowsDatabase(t *testing.T) { + repo, pool := lifecycleTestRepository(t) + ctx := t.Context() + early := Viewer{UserID: 1, ProfileID: "profile-early"} + late := Viewer{UserID: 2, ProfileID: "profile-late"} + insertLifecycleRequest(t, repo, "done", 5, 974, StatusPending) + if err := repo.FollowTitle(ctx, MediaTypeMovie, 974, early); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() WHERE id = 'done'`); err != nil { + t.Fatal(err) + } + insertLifecycleRequest(t, repo, "next", 6, 974, StatusPending) + if err := repo.FollowTitle(ctx, MediaTypeMovie, 974, late); err != nil { + t.Fatal(err) + } + // Pin the follow times either side of the completion. + if _, err := pool.Exec(ctx, ` + UPDATE media_request_follows f SET created_at = r.completed_at + + CASE WHEN f.profile_id = 'profile-early' THEN -interval '1 minute' ELSE interval '1 minute' END + FROM media_requests r WHERE r.id = 'done' AND f.tmdb_id = 974`); err != nil { + t.Fatal(err) + } + done, err := repo.GetRequest(ctx, "done") + if err != nil { + t.Fatal(err) + } + followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 974, done.CompletedAt) + if err != nil || !slices.Equal(followers, []Follower{{UserID: 1, ProfileID: "profile-early"}}) { + t.Fatalf("followers of the completed request = %+v, err = %v; want the early follow only", followers, err) + } + + if _, err := repo.SetOutcome(ctx, "next", guardWithdrawable, OutcomeDeclined, Viewer{}, ""); err != nil { + t.Fatal(err) + } + followers, err = repo.ListTitleFollowers(ctx, MediaTypeMovie, 974, nil) + if err != nil || !slices.Equal(followers, []Follower{{UserID: 1, ProfileID: "profile-early"}}) { + t.Fatalf("followers after declining the newer request = %+v, err = %v; want the completed request's follower kept", 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 +525,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 := repo.ListTitleFollowers(ctx, MediaTypeMovie, 959, nil); 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..e0ec19e7b2 100644 --- a/internal/requests/notify.go +++ b/internal/requests/notify.go @@ -129,7 +129,10 @@ func (s *Service) notifyFulfilledPending(ctx context.Context) { s.markNotifyChecked(ctx, req.ID) continue } - followers, err := s.store.ListTitleFollowers(ctx, req.MediaType, req.TMDBID) + // A follow made after this request completed was made for another + // open request of the title (other seasons of a series) and waits + // for that one. + followers, err := s.store.ListTitleFollowers(ctx, req.MediaType, req.TMDBID, req.CompletedAt) if err != nil { slog.WarnContext(ctx, "request fulfill-notify: list followers failed", "component", "requests", "request_id", req.ID, "err", err) @@ -141,12 +144,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. + // The followers have been told. Only the listed rows are cleared, + // leaving the follows made since for the request they were made + // for. 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 { slog.WarnContext(ctx, "request fulfill-notify: clear followers failed", "component", "requests", "request_id", req.ID, "err", err) diff --git a/internal/requests/service_test.go b/internal/requests/service_test.go index 45d5a88831..bb45746504 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 + followedAt map[string]time.Time + clearErr error // returned by ClearTitleFollowers when set routes []Route factsSet map[string]RoutingFacts groupLimits map[int64]*GroupLimit @@ -2553,10 +2554,26 @@ func (f *fakeStore) seedFollow(mediaType MediaType, tmdbID int, viewer Viewer) { } func (f *fakeStore) seedFollowLocked(mediaType MediaType, tmdbID int, viewer Viewer) { + f.seedFollowAtLocked(mediaType, tmdbID, viewer, time.Now()) +} + +// seedFollowAt records a follow made at a given time. +func (f *fakeStore) seedFollowAt(mediaType MediaType, tmdbID int, viewer Viewer, at time.Time) { + f.mu.Lock() + defer f.mu.Unlock() + f.seedFollowAtLocked(mediaType, tmdbID, viewer, at) +} + +func (f *fakeStore) seedFollowAtLocked(mediaType MediaType, tmdbID int, viewer Viewer, at time.Time) { 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.followedAt == nil { + f.followedAt = map[string]time.Time{} + } + key := followKey(mediaType, tmdbID, viewer.UserID, viewer.ProfileID) + f.follows[key] = Follower{UserID: viewer.UserID, ProfileID: viewer.ProfileID} + f.followedAt[key] = at } func (f *fakeStore) UnfollowTitle(_ context.Context, mediaType MediaType, tmdbID int, viewer Viewer) error { @@ -2578,13 +2595,13 @@ func (f *fakeStore) FollowedTitles(_ context.Context, mediaType MediaType, tmdbI return out, nil } -func (f *fakeStore) ListTitleFollowers(_ context.Context, mediaType MediaType, tmdbID int) ([]Follower, error) { +func (f *fakeStore) ListTitleFollowers(_ context.Context, mediaType MediaType, tmdbID int, followedBy *time.Time) ([]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 strings.HasPrefix(key, prefix) && (followedBy == nil || !f.followedAt[key].After(*followedBy)) { out = append(out, follower) } } diff --git a/internal/requests/store.go b/internal/requests/store.go index fcab42db4a..e2857bebd2 100644 --- a/internal/requests/store.go +++ b/internal/requests/store.go @@ -79,11 +79,12 @@ type Store interface { // 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. + // ErrNotRequested when the title has no open request. ListTitleFollowers + // lists the follows made at or before followedBy, or all when it is nil. 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) + ListTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int, followedBy *time.Time) ([]Follower, error) ClearTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int, followers []Follower) error // ListRoutes returns every routing rule, in no particular order; // decideRoutes orders them. From 2d1103e10fc69f0ee9fa7679734de48ef4c53b0f Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Mon, 28 Sep 2026 23:32:38 +0000 Subject: [PATCH 02/10] fix(requests): recheck a server's type against its routes under the row lock A route save and a server save each validated the other before their transactions and rechecked only the 4K switch under the server's row lock. Two admins saving at once could leave a movie route pointing at a server that had just become a Sonarr, or one that no longer takes movies, and every request it routed would fail when sent. Both saves now recheck the server's type and media types with the row locked, and a server save does so under Standard too. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/architecture/media-requests.md | 4 ++ internal/requests/editor_concurrency.go | 6 +- internal/requests/repository.go | 6 +- internal/requests/routes_admin.go | 89 +++++++++++++++---------- internal/requests/routes_admin_test.go | 23 +++++++ internal/requests/routing_mode.go | 26 +++++--- internal/requests/routing_mode_test.go | 31 +++++++++ 7 files changed, 128 insertions(+), 57 deletions(-) diff --git a/docs/architecture/media-requests.md b/docs/architecture/media-requests.md index b4b131858f..7a4010a935 100644 --- a/docs/architecture/media-requests.md +++ b/docs/architecture/media-requests.md @@ -118,6 +118,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 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/repository.go b/internal/requests/repository.go index 7e64e9d350..61cf2d4ac3 100644 --- a/internal/requests/repository.go +++ b/internal/requests/repository.go @@ -1016,11 +1016,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) + } +} From 9e329f4a6d1585e368efb318a0a068e0679eed81 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Mon, 28 Sep 2026 23:55:49 +0000 Subject: [PATCH 03/10] fix(requests): bound a request's followers by the title's previous completion Two completed requests for one series can both wait for the library, and the newer can arrive first. Bounding followers only by the request's own completion time told the newer request's notification about the older request's followers too, and cleared them. A title has one open request at a time, so a request's followers are the follows made after the title's previous request completed and no later than it did. Clearing them is bounded the same way, so a profile that unfollowed and followed again during the dispatch keeps its new follow. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/architecture/media-requests.md | 14 +++-- internal/requests/follows.go | 38 +++++++---- internal/requests/follows_test.go | 97 ++++++++++++++++++++++------- internal/requests/notify.go | 18 +++--- internal/requests/service_test.go | 28 +++++++-- internal/requests/store.go | 12 ++-- 6 files changed, 145 insertions(+), 62 deletions(-) diff --git a/docs/architecture/media-requests.md b/docs/architecture/media-requests.md index 7a4010a935..87c90e6c22 100644 --- a/docs/architecture/media-requests.md +++ b/docs/architecture/media-requests.md @@ -316,12 +316,14 @@ 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 title's followers, marked `follower` so its wording does not say -"your request", and those follows are then cleared. A series can have a -completed request still waiting for the library beside a newer open request -for other seasons, so the notification goes only to the follows made before -its request completed; later follows were made for the open request and wait -for it. Declining that open request leaves the earlier follows for the -completed one. A dispatch failure leaves +"your request", and those follows are then cleared. A series can have +completed requests still waiting for the library beside a newer open request +for other seasons. A title has one open request at a time, so a request's +notification goes to the follows made after the title's previous request +completed and no later than it did; the others wait for their own request. +Declining the open request leaves the follows the completed ones are waiting +to tell, and clearing a request's follows spares a profile that followed again +since. 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/internal/requests/follows.go b/internal/requests/follows.go index edbc0fbbbb..abf18945f0 100644 --- a/internal/requests/follows.go +++ b/internal/requests/follows.go @@ -4,7 +4,6 @@ import ( "context" "fmt" "strings" - "time" ) // Following a title: a profile that finds a title someone else has already @@ -16,11 +15,12 @@ import ( // follower can request it themselves. The requester is always notified and // never needs a follow. // -// A series can have a completed request still waiting for the library and a -// newer open one for other seasons. A follow made after the first completed -// was made for the open one, so the first request's notification goes only to -// the follows made before it completed, and closing the open one leaves those -// follows for the completed request. +// A series can have completed requests still waiting for the library and a +// newer open one for other seasons. A title has one open request at a time, so +// each follow was made for the request open then: a request's notification +// goes to the follows made after the title's previous request completed and no +// later than it did, and closing the open request leaves the follows the +// completed ones are waiting to tell. // Follower is a profile waiting to hear that a requested title is available. type Follower struct { @@ -220,15 +220,23 @@ func (r *Repository) FollowedTitles(ctx context.Context, mediaType MediaType, tm return out, rows.Err() } -func (r *Repository) ListTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int, followedBy *time.Time) ([]Follower, error) { +// ListRequestFollowers lists the follows a completed request's notification +// goes to. A title has at most one open request at a time, so its follows were +// made after the title's previous request completed and no later than this one +// did. A request without a completion time reaches every follow on the title. +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 AND ($3::timestamptz IS NULL OR created_at <= $3) + AND created_at > coalesce(( + SELECT max(completed_at) FROM media_requests + WHERE media_type = $1 AND provider = 'tmdb' AND tmdb_id = $2 + AND status = 'completed' AND id <> $4 AND completed_at < $3), '-infinity') ORDER BY created_at, user_id, profile_id - `, mediaType, tmdbID, followedBy) + `, req.MediaType, req.TMDBID, req.CompletedAt, 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 @@ -242,7 +250,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 } @@ -255,9 +266,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) + AND ($3::timestamptz IS NULL OR created_at <= $3) + AND (user_id, profile_id) IN (SELECT * FROM unnest($4::int[], $5::text[])) + `, req.MediaType, req.TMDBID, req.CompletedAt, 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 8a93bcc35b..6b47b0fba1 100644 --- a/internal/requests/follows_test.go +++ b/internal/requests/follows_test.go @@ -87,7 +87,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, nil) + followers, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 949}) if len(followers) != 1 || followers[0] != (Follower{UserID: 1, ProfileID: "default"}) { t.Fatalf("followers = %+v, want the other account's default profile", followers) } @@ -160,7 +160,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, nil); len(followers) != 0 { + if followers, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 949}); len(followers) != 0 { t.Fatalf("followers after %s = %+v, want none", tc.name, followers) } }) @@ -214,7 +214,7 @@ 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, nil); len(followers) != 0 { + if followers, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 42}); len(followers) != 0 { t.Fatalf("followers after notifying = %+v, want cleared", followers) } } @@ -242,7 +242,7 @@ func TestNotifyFulfilledLeavesFollowsOfNewerRequest(t *testing.T) { 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 the follow made before completion", notifier.followers) } - left, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42, nil) + left, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 42}) if !slices.Equal(left, []Follower{{UserID: 4, ProfileID: "late"}}) { t.Fatalf("followers left = %+v, want the newer request's follower", left) } @@ -258,7 +258,7 @@ func TestNotifyFulfilledKeepsFollowersWhenDispatchFails(t *testing.T) { svc.notifyFulfilledPending(context.Background()) - if followers, _ := store.ListTitleFollowers(context.Background(), MediaTypeMovie, 42, nil); len(followers) != 1 { + if followers, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 42}); len(followers) != 1 { t.Fatalf("followers after a failed dispatch = %+v, want kept for the retry", followers) } } @@ -284,7 +284,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, nil); len(followers) != 0 { + if followers, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 42}); len(followers) != 0 { t.Fatalf("followers after the retry = %+v, want cleared", followers) } } @@ -316,7 +316,7 @@ func TestFollowsDatabase(t *testing.T) { t.Fatal(err) } - followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil) + followers, err := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 949}) if err != nil { t.Fatal(err) } @@ -331,16 +331,16 @@ func TestFollowsDatabase(t *testing.T) { t.Fatalf("followed = %v, want only 949", followed) } - if err := repo.ClearTitleFollowers(ctx, MediaTypeMovie, 949, []Follower{{UserID: a.UserID, ProfileID: a.ProfileID}}); err != nil { + if err := repo.ClearRequestFollowers(ctx, Request{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, nil); len(followers) != 0 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{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, nil); len(followers) != 1 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeSeries, TMDBID: 949}); len(followers) != 1 { t.Fatalf("series followers = %+v, want the one untouched follow", followers) } @@ -352,7 +352,7 @@ func TestFollowsDatabase(t *testing.T) { t.Fatalf("follow as account %d: %v", v.UserID, err) } } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil); len(followers) != 2 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{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] { @@ -361,27 +361,27 @@ func TestFollowsDatabase(t *testing.T) { if err := repo.UnfollowTitle(ctx, MediaTypeMovie, 949, mine); err != nil { t.Fatal(err) } - followers, err = repo.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil) + followers, err = repo.ListRequestFollowers(ctx, Request{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{MediaType: MediaTypeMovie, TMDBID: 949}, []Follower{{UserID: mine.UserID, ProfileID: mine.ProfileID}}); err != nil { t.Fatal(err) } - if followers, _ := repo.ListTitleFollowers(ctx, MediaTypeMovie, 949, nil); len(followers) != 1 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{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, nil); len(followers) != 0 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{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, nil); len(followers) != 1 { + if followers, _ := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 949}); len(followers) != 1 { t.Fatalf("movie followers after the series withdrawal = %+v, want the one left", followers) } } @@ -408,7 +408,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, nil); err != nil || len(followers) != 0 { + if followers, err := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: tc.tmdbID}); err != nil || len(followers) != 0 { t.Fatalf("followers once %s committed = %+v, err = %v; want none", tc.id, followers, err) } } @@ -426,7 +426,7 @@ 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, nil); err != nil || len(followers) != 1 { + if followers, err := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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) } } @@ -461,7 +461,7 @@ func TestDecliningNewerRequestKeepsCompletedRequestFollowsDatabase(t *testing.T) if err != nil { t.Fatal(err) } - followers, err := repo.ListTitleFollowers(ctx, MediaTypeMovie, 974, done.CompletedAt) + followers, err := repo.ListRequestFollowers(ctx, *done) if err != nil || !slices.Equal(followers, []Follower{{UserID: 1, ProfileID: "profile-early"}}) { t.Fatalf("followers of the completed request = %+v, err = %v; want the early follow only", followers, err) } @@ -469,12 +469,67 @@ func TestDecliningNewerRequestKeepsCompletedRequestFollowsDatabase(t *testing.T) if _, err := repo.SetOutcome(ctx, "next", guardWithdrawable, OutcomeDeclined, Viewer{}, ""); err != nil { t.Fatal(err) } - followers, err = repo.ListTitleFollowers(ctx, MediaTypeMovie, 974, nil) + followers, err = repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 974}) if err != nil || !slices.Equal(followers, []Follower{{UserID: 1, ProfileID: "profile-early"}}) { t.Fatalf("followers after declining the newer request = %+v, err = %v; want the completed request's follower kept", followers, err) } } +// Two completed requests for one title can both wait for the library, and the +// newer one can arrive first. Each tells only the follows made while it was +// the title's open request, and a clear spares a follow made again since. +func TestRequestFollowersWindowDatabase(t *testing.T) { + repo, pool := lifecycleTestRepository(t) + ctx := t.Context() + insertLifecycleRequest(t, repo, "older", 5, 975, StatusPending) + if _, err := pool.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() - interval '2 hours' WHERE id = 'older'`); err != nil { + t.Fatal(err) + } + insertLifecycleRequest(t, repo, "newer", 6, 975, StatusPending) + if _, err := pool.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() - interval '1 hour' WHERE id = 'newer'`); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `INSERT INTO media_request_follows (media_type, tmdb_id, user_id, profile_id, created_at) VALUES + ('movie', 975, 1, 'profile-older', now() - interval '3 hours'), + ('movie', 975, 2, 'profile-newer', now() - interval '90 minutes')`); err != nil { + t.Fatal(err) + } + get := func(id string) Request { + t.Helper() + req, err := repo.GetRequest(ctx, id) + if err != nil { + t.Fatal(err) + } + return *req + } + older, newer := get("older"), get("newer") + for _, tc := range []struct { + req Request + want Follower + }{ + {newer, Follower{UserID: 2, ProfileID: "profile-newer"}}, + {older, Follower{UserID: 1, ProfileID: "profile-older"}}, + } { + followers, err := repo.ListRequestFollowers(ctx, tc.req) + if err != nil || !slices.Equal(followers, []Follower{tc.want}) { + t.Fatalf("followers of %s = %+v, err = %v; want %+v", tc.req.ID, followers, err, tc.want) + } + } + + // profile-newer unfollowed and followed again while the notification was + // going out; the new follow is for a later request. + if _, err := pool.Exec(ctx, `UPDATE media_request_follows SET created_at = now() WHERE profile_id = 'profile-newer'`); err != nil { + t.Fatal(err) + } + if err := repo.ClearRequestFollowers(ctx, newer, []Follower{{UserID: 2, ProfileID: "profile-newer"}}); err != nil { + t.Fatal(err) + } + left, err := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 975}) + if err != nil || len(left) != 2 { + t.Fatalf("follows after the clear = %+v, err = %v; want both kept", left, 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. @@ -525,7 +580,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, nil); err != nil || len(followers) != 0 { + if followers, err := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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 e0ec19e7b2..872db2a54a 100644 --- a/internal/requests/notify.go +++ b/internal/requests/notify.go @@ -129,10 +129,9 @@ func (s *Service) notifyFulfilledPending(ctx context.Context) { s.markNotifyChecked(ctx, req.ID) continue } - // A follow made after this request completed was made for another - // open request of the title (other seasons of a series) and waits - // for that one. - followers, err := s.store.ListTitleFollowers(ctx, req.MediaType, req.TMDBID, req.CompletedAt) + // Only the follows made while this request was the title's open + // one: a series can have several requests for different seasons. + 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) @@ -145,12 +144,11 @@ func (s *Service) notifyFulfilledPending(ctx context.Context) { continue } // The followers have been told. Only the listed rows are cleared, - // leaving the follows made since for the request they were made - // for. 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 { + // 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/service_test.go b/internal/requests/service_test.go index bb45746504..5c92632897 100644 --- a/internal/requests/service_test.go +++ b/internal/requests/service_test.go @@ -1992,7 +1992,7 @@ type fakeStore struct { reconciled []string follows map[string]Follower // key: media_type/tmdb_id/user_id/profile_id followedAt map[string]time.Time - clearErr error // returned by ClearTitleFollowers when set + clearErr error // returned by ClearRequestFollowers when set routes []Route factsSet map[string]RoutingFacts groupLimits map[int64]*GroupLimit @@ -2595,13 +2595,26 @@ func (f *fakeStore) FollowedTitles(_ context.Context, mediaType MediaType, tmdbI return out, nil } -func (f *fakeStore) ListTitleFollowers(_ context.Context, mediaType MediaType, tmdbID int, followedBy *time.Time) ([]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) + // The window the repository uses: after the title's previous completed + // request, no later than this one. + var after time.Time + if req.CompletedAt != nil { + for _, other := range f.requests { + if other.ID != req.ID && other.MediaType == req.MediaType && other.TMDBID == req.TMDBID && + other.Status == StatusCompleted && other.CompletedAt != nil && + other.CompletedAt.Before(*req.CompletedAt) && other.CompletedAt.After(after) { + after = *other.CompletedAt + } + } + } + prefix := fmt.Sprintf("%s/%d/", req.MediaType, req.TMDBID) var out []Follower for key, follower := range f.follows { - if strings.HasPrefix(key, prefix) && (followedBy == nil || !f.followedAt[key].After(*followedBy)) { + at := f.followedAt[key] + if strings.HasPrefix(key, prefix) && at.After(after) && (req.CompletedAt == nil || !at.After(*req.CompletedAt)) { out = append(out, follower) } } @@ -2614,14 +2627,17 @@ 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 { +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)) + key := followKey(req.MediaType, req.TMDBID, follower.UserID, follower.ProfileID) + if req.CompletedAt == nil || !f.followedAt[key].After(*req.CompletedAt) { + delete(f.follows, key) + } } return nil } diff --git a/internal/requests/store.go b/internal/requests/store.go index e2857bebd2..79c4f64c2b 100644 --- a/internal/requests/store.go +++ b/internal/requests/store.go @@ -77,15 +77,15 @@ type Store interface { // 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. ListTitleFollowers - // lists the follows made at or before followedBy, or all when it is nil. + // on titles, keyed by account and profile; 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, followedBy *time.Time) ([]Follower, error) - ClearTitleFollowers(ctx context.Context, mediaType MediaType, tmdbID int, followers []Follower) 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) From 5bebde23f2acc7ced8d3caa4cc9ad9d4e8a39225 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 00:27:12 +0000 Subject: [PATCH 04/10] fix(requests): record the open request on each follow Bounding a request's followers by timestamps broke when a follow committed while a completion was under way: the completion's timestamp comes from the start of its transaction, so it could be earlier than the follow's, and the follow was left out of the request it had read. A follow now records the request that was open when it was made, and a request's notification, its follow cleanup and a decline all go by that request. A new request takes over the follows of a failed request, or of one its requester replaced, so a follow still survives its request failing. Existing follows go to the title's open request, or else to its latest completed request that has not been notified. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/architecture/media-requests.md | 29 +-- internal/requests/follows.go | 85 +++--- internal/requests/follows_test.go | 246 +++++++++++------- internal/requests/notify.go | 2 - internal/requests/repository.go | 3 + internal/requests/service_test.go | 67 +++-- ...0929001953_request_follows_request_key.sql | 28 ++ 7 files changed, 267 insertions(+), 193 deletions(-) create mode 100644 migrations/sql/20260929001953_request_follows_request_key.sql diff --git a/docs/architecture/media-requests.md b/docs/architecture/media-requests.md index 87c90e6c22..02028eda3f 100644 --- a/docs/architecture/media-requests.md +++ b/docs/architecture/media-requests.md @@ -308,22 +308,19 @@ 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 the title's followers, marked `follower` so its wording does not say -"your request", and those follows are then cleared. A series can have -completed requests still waiting for the library beside a newer open request -for other seasons. A title has one open request at a time, so a request's -notification goes to the follows made after the title's previous request -completed and no later than it did; the others wait for their own request. -Declining the open request leaves the follows the completed ones are waiting -to tell, and clearing a request's follows spares a profile that followed again -since. A dispatch failure leaves +A follow belongs to a title and a profile (`media_request_follows`, keyed +by account and profile id, since profile ids repeat across accounts) and records +the request that was open when it was made. A series can have completed +requests still waiting for the library beside a newer open request for other +seasons, and each request's notification goes to its own follows. 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/internal/requests/follows.go b/internal/requests/follows.go index abf18945f0..b1f17d13c1 100644 --- a/internal/requests/follows.go +++ b/internal/requests/follows.go @@ -8,19 +8,15 @@ import ( // 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. -// -// A series can have completed requests still waiting for the library and a -// newer open one for other seasons. A title has one open request at a time, so -// each follow was made for the request open then: a request's notification -// goes to the follows made after the title's previous request completed and no -// later than it did, and closing the open request leaves the follows the -// completed ones are waiting to tell. +// requesting it again. A follow records 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. 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 { @@ -134,7 +130,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 @@ -144,14 +141,14 @@ 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 + 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 (media_type, tmdb_id, user_id, profile_id) DO NOTHING ) SELECT EXISTS (SELECT 1 FROM open_request) @@ -164,25 +161,26 @@ 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. A follow made before a completed request -// of the title completed stays until that request's notification goes out. +// 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 request_id = $1`, closed.ID); err != nil { + return fmt.Errorf("forget title follows: %w", err) + } + return nil +} + +// adoptTitleFollows gives a new request the title's follows that have no +// request waiting to tell them: their request failed, or was replaced by its +// requester. The caller creates the request in the same transaction. +func adoptTitleFollows(ctx context.Context, exec requestExecutor, req *Request) error { if _, err := exec.Exec(ctx, ` - DELETE FROM media_request_follows f - WHERE media_type = $1 AND tmdb_id = $2 + UPDATE media_request_follows f SET request_id = $3 + WHERE f.media_type = $1 AND f.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) - 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 fulfilled_notified_at IS NULL AND completed_at >= f.created_at) - `, closed.MediaType, closed.TMDBID, closed.ID); err != nil { - return fmt.Errorf("forget title follows: %w", err) + SELECT 1 FROM media_requests r WHERE r.id = f.request_id AND r.outcome <> 'failed') + `, req.MediaType, req.TMDBID, req.ID); err != nil { + return fmt.Errorf("adopt title follows: %w", err) } return nil } @@ -220,21 +218,13 @@ func (r *Repository) FollowedTitles(ctx context.Context, mediaType MediaType, tm return out, rows.Err() } -// ListRequestFollowers lists the follows a completed request's notification -// goes to. A title has at most one open request at a time, so its follows were -// made after the title's previous request completed and no later than this one -// did. A request without a completion time reaches every follow on the title. +// 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 - AND ($3::timestamptz IS NULL OR created_at <= $3) - AND created_at > coalesce(( - SELECT max(completed_at) FROM media_requests - WHERE media_type = $1 AND provider = 'tmdb' AND tmdb_id = $2 - AND status = 'completed' AND id <> $4 AND completed_at < $3), '-infinity') + WHERE request_id = $1 ORDER BY created_at, user_id, profile_id - `, req.MediaType, req.TMDBID, req.CompletedAt, req.ID) + `, req.ID) if err != nil { return nil, fmt.Errorf("list request followers: %w", err) } @@ -265,10 +255,9 @@ func (r *Repository) ClearRequestFollowers(ctx context.Context, req Request, fol } if _, err := r.pool.Exec(ctx, ` DELETE FROM media_request_follows - WHERE media_type = $1 AND tmdb_id = $2 - AND ($3::timestamptz IS NULL OR created_at <= $3) - AND (user_id, profile_id) IN (SELECT * FROM unnest($4::int[], $5::text[])) - `, req.MediaType, req.TMDBID, req.CompletedAt, userIDs, profileIDs); err != nil { + 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 6b47b0fba1..b3516a5b7e 100644 --- a/internal/requests/follows_test.go +++ b/internal/requests/follows_test.go @@ -3,10 +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" ) @@ -87,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.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 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) } @@ -160,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.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 949}); len(followers) != 0 { + if followers, _ := store.titleFollowers(MediaTypeMovie, 949); len(followers) != 0 { t.Fatalf("followers after %s = %+v, want none", tc.name, followers) } }) @@ -214,25 +218,20 @@ 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.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 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 the follows made -// before it completed; one made since was made for the open request. +// 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() - completed := time.Now().Add(-time.Hour) - req := completedRequestFixture("req1", 42) - req.CompletedAt = &completed - store.requests["req1"] = req + store.requests["req1"] = completedRequestFixture("req1", 42) store.unnotified = []string{"req1"} - early := Viewer{UserID: 3, ProfileID: "early"} - late := Viewer{UserID: 4, ProfileID: "late"} - store.seedFollowAt(MediaTypeMovie, 42, early, completed.Add(-time.Minute)) - store.seedFollowAt(MediaTypeMovie, 42, late, completed.Add(time.Minute)) + 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) @@ -240,11 +239,11 @@ func TestNotifyFulfilledLeavesFollowsOfNewerRequest(t *testing.T) { 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 the follow made before completion", notifier.followers) + t.Fatalf("followers handed to the notifier = %+v, want only req1's follower", notifier.followers) } - left, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 42}) + left, _ := store.ListRequestFollowers(context.Background(), Request{ID: "req2", MediaType: MediaTypeMovie, TMDBID: 42}) if !slices.Equal(left, []Follower{{UserID: 4, ProfileID: "late"}}) { - t.Fatalf("followers left = %+v, want the newer request's follower", left) + t.Fatalf("req2's followers = %+v, want kept", left) } } @@ -258,7 +257,7 @@ func TestNotifyFulfilledKeepsFollowersWhenDispatchFails(t *testing.T) { svc.notifyFulfilledPending(context.Background()) - if followers, _ := store.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 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) } } @@ -284,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.ListRequestFollowers(context.Background(), Request{MediaType: MediaTypeMovie, TMDBID: 42}); len(followers) != 0 { + if followers, _ := store.titleFollowers(MediaTypeMovie, 42); len(followers) != 0 { t.Fatalf("followers after the retry = %+v, want cleared", followers) } } @@ -316,7 +315,7 @@ func TestFollowsDatabase(t *testing.T) { t.Fatal(err) } - followers, err := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 949}) + followers, err := repo.ListRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}) if err != nil { t.Fatal(err) } @@ -331,16 +330,16 @@ func TestFollowsDatabase(t *testing.T) { t.Fatalf("followed = %v, want only 949", followed) } - if err := repo.ClearRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeSeries, TMDBID: 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) } @@ -352,7 +351,7 @@ func TestFollowsDatabase(t *testing.T) { t.Fatalf("follow as account %d: %v", v.UserID, err) } } - if followers, _ := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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] { @@ -361,27 +360,27 @@ func TestFollowsDatabase(t *testing.T) { if err := repo.UnfollowTitle(ctx, MediaTypeMovie, 949, mine); err != nil { t.Fatal(err) } - followers, err = repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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.ClearRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeSeries, TMDBID: 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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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) } } @@ -408,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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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) } } @@ -426,107 +425,156 @@ 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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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) } } -// A completed request still waiting for the library keeps the follows made -// before it completed when a newer request for the title is declined; the -// follows made for the declined request go. -func TestDecliningNewerRequestKeepsCompletedRequestFollowsDatabase(t *testing.T) { +// 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() - early := Viewer{UserID: 1, ProfileID: "profile-early"} - late := Viewer{UserID: 2, ProfileID: "profile-late"} - insertLifecycleRequest(t, repo, "done", 5, 974, StatusPending) - if err := repo.FollowTitle(ctx, MediaTypeMovie, 974, early); err != nil { - t.Fatal(err) + 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) + } } - if _, err := pool.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() WHERE id = 'done'`); 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, "next", 6, 974, StatusPending) - if err := repo.FollowTitle(ctx, MediaTypeMovie, 974, late); err != nil { - t.Fatal(err) + 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) } - // Pin the follow times either side of the completion. - if _, err := pool.Exec(ctx, ` - UPDATE media_request_follows f SET created_at = r.completed_at - + CASE WHEN f.profile_id = 'profile-early' THEN -interval '1 minute' ELSE interval '1 minute' END - FROM media_requests r WHERE r.id = 'done' AND f.tmdb_id = 974`); err != nil { + 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) } - done, err := repo.GetRequest(ctx, "done") - if err != nil { + 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) } - followers, err := repo.ListRequestFollowers(ctx, *done) - if err != nil || !slices.Equal(followers, []Follower{{UserID: 1, ProfileID: "profile-early"}}) { - t.Fatalf("followers of the completed request = %+v, err = %v; want the early follow only", followers, err) + if got := list("third"); len(got) != 1 { + t.Fatalf("third's followers after newer's clear = %+v, want profile-b kept", got) } - if _, err := repo.SetOutcome(ctx, "next", guardWithdrawable, OutcomeDeclined, Viewer{}, ""); err != nil { + if _, err := repo.SetOutcome(ctx, "third", guardWithdrawable, OutcomeDeclined, Viewer{}, ""); err != nil { t.Fatal(err) } - followers, err = repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 974}) - if err != nil || !slices.Equal(followers, []Follower{{UserID: 1, ProfileID: "profile-early"}}) { - t.Fatalf("followers after declining the newer request = %+v, err = %v; want the completed request's follower kept", followers, 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) } } -// Two completed requests for one title can both wait for the library, and the -// newer one can arrive first. Each tells only the follows made while it was -// the title's open request, and a clear spares a follow made again since. -func TestRequestFollowersWindowDatabase(t *testing.T) { +// 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, "older", 5, 975, StatusPending) - if _, err := pool.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() - interval '2 hours' WHERE id = 'older'`); err != nil { + insertLifecycleRequest(t, repo, "req", 5, 976, StatusPending) + tx, err := pool.Begin(ctx) + if err != nil { t.Fatal(err) } - insertLifecycleRequest(t, repo, "newer", 6, 975, StatusPending) - if _, err := pool.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() - interval '1 hour' WHERE id = 'newer'`); err != nil { + defer func() { _ = tx.Rollback(ctx) }() + if _, err := tx.Exec(ctx, `SELECT now()`); err != nil { t.Fatal(err) } - if _, err := pool.Exec(ctx, `INSERT INTO media_request_follows (media_type, tmdb_id, user_id, profile_id, created_at) VALUES - ('movie', 975, 1, 'profile-older', now() - interval '3 hours'), - ('movie', 975, 2, 'profile-newer', now() - interval '90 minutes')`); err != nil { + if err := repo.FollowTitle(ctx, MediaTypeMovie, 976, Viewer{UserID: 1, ProfileID: "profile-a"}); err != nil { t.Fatal(err) } - get := func(id string) Request { - t.Helper() - req, err := repo.GetRequest(ctx, id) - if err != nil { - t.Fatal(err) - } - return *req - } - older, newer := get("older"), get("newer") - for _, tc := range []struct { - req Request - want Follower - }{ - {newer, Follower{UserID: 2, ProfileID: "profile-newer"}}, - {older, Follower{UserID: 1, ProfileID: "profile-older"}}, - } { - followers, err := repo.ListRequestFollowers(ctx, tc.req) - if err != nil || !slices.Equal(followers, []Follower{tc.want}) { - t.Fatalf("followers of %s = %+v, err = %v; want %+v", tc.req.ID, followers, err, tc.want) - } - } - - // profile-newer unfollowed and followed again while the notification was - // going out; the new follow is for a later request. - if _, err := pool.Exec(ctx, `UPDATE media_request_follows SET created_at = now() WHERE profile_id = 'profile-newer'`); err != nil { + if _, err := tx.Exec(ctx, `UPDATE media_requests SET status = 'completed', completed_at = now() WHERE id = 'req'`); err != nil { t.Fatal(err) } - if err := repo.ClearRequestFollowers(ctx, newer, []Follower{{UserID: 2, ProfileID: "profile-newer"}}); err != nil { + if err := tx.Commit(ctx); err != nil { t.Fatal(err) } - left, err := repo.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 975}) - if err != nil || len(left) != 2 { - t.Fatalf("follows after the clear = %+v, err = %v; want both kept", left, 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) + } } } @@ -580,7 +628,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.ListRequestFollowers(ctx, Request{MediaType: MediaTypeMovie, TMDBID: 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 872db2a54a..5f40eca612 100644 --- a/internal/requests/notify.go +++ b/internal/requests/notify.go @@ -129,8 +129,6 @@ func (s *Service) notifyFulfilledPending(ctx context.Context) { s.markNotifyChecked(ctx, req.ID) continue } - // Only the follows made while this request was the title's open - // one: a series can have several requests for different seasons. followers, err := s.store.ListRequestFollowers(ctx, *req) if err != nil { slog.WarnContext(ctx, "request fulfill-notify: list followers failed", "component", "requests", diff --git a/internal/requests/repository.go b/internal/requests/repository.go index 61cf2d4ac3..a81fdc2489 100644 --- a/internal/requests/repository.go +++ b/internal/requests/repository.go @@ -292,6 +292,9 @@ func (r *Repository) CreateRequest(ctx context.Context, input CreateRequestRecor } return nil, err } + if err := adoptTitleFollows(ctx, tx, req); err != nil { + return nil, err + } if err := r.recordEvent(ctx, tx, req.ID, "created", input.Requester, ""); err != nil { return nil, err } diff --git a/internal/requests/service_test.go b/internal/requests/service_test.go index 5c92632897..e9ae852e05 100644 --- a/internal/requests/service_test.go +++ b/internal/requests/service_test.go @@ -1991,8 +1991,8 @@ type fakeStore struct { notified []string reconciled []string follows map[string]Follower // key: media_type/tmdb_id/user_id/profile_id - followedAt map[string]time.Time - clearErr error // returned by ClearRequestFollowers 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 @@ -2319,9 +2319,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) } } @@ -2553,27 +2552,39 @@ 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) { - f.seedFollowAtLocked(mediaType, tmdbID, viewer, time.Now()) + 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) } -// seedFollowAt records a follow made at a given time. -func (f *fakeStore) seedFollowAt(mediaType MediaType, tmdbID int, viewer Viewer, at time.Time) { +// 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.seedFollowAtLocked(mediaType, tmdbID, viewer, at) + f.seedFollowForLocked(mediaType, tmdbID, viewer, requestID) } -func (f *fakeStore) seedFollowAtLocked(mediaType MediaType, tmdbID int, viewer Viewer, at time.Time) { +func (f *fakeStore) seedFollowForLocked(mediaType MediaType, tmdbID int, viewer Viewer, requestID string) { if f.follows == nil { f.follows = map[string]Follower{} } - if f.followedAt == nil { - f.followedAt = map[string]time.Time{} + if f.followFor == nil { + f.followFor = map[string]string{} } key := followKey(mediaType, tmdbID, viewer.UserID, viewer.ProfileID) f.follows[key] = Follower{UserID: viewer.UserID, ProfileID: viewer.ProfileID} - f.followedAt[key] = at + f.followFor[key] = requestID } func (f *fakeStore) UnfollowTitle(_ context.Context, mediaType MediaType, tmdbID int, viewer Viewer) error { @@ -2598,23 +2609,9 @@ func (f *fakeStore) FollowedTitles(_ context.Context, mediaType MediaType, tmdbI func (f *fakeStore) ListRequestFollowers(_ context.Context, req Request) ([]Follower, error) { f.mu.Lock() defer f.mu.Unlock() - // The window the repository uses: after the title's previous completed - // request, no later than this one. - var after time.Time - if req.CompletedAt != nil { - for _, other := range f.requests { - if other.ID != req.ID && other.MediaType == req.MediaType && other.TMDBID == req.TMDBID && - other.Status == StatusCompleted && other.CompletedAt != nil && - other.CompletedAt.Before(*req.CompletedAt) && other.CompletedAt.After(after) { - after = *other.CompletedAt - } - } - } - prefix := fmt.Sprintf("%s/%d/", req.MediaType, req.TMDBID) var out []Follower for key, follower := range f.follows { - at := f.followedAt[key] - if strings.HasPrefix(key, prefix) && at.After(after) && (req.CompletedAt == nil || !at.After(*req.CompletedAt)) { + if f.followFor[key] == req.ID { out = append(out, follower) } } @@ -2627,6 +2624,20 @@ func (f *fakeStore) ListRequestFollowers(_ context.Context, req Request) ([]Foll return out, nil } +// 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() @@ -2635,7 +2646,7 @@ func (f *fakeStore) ClearRequestFollowers(_ context.Context, req Request, follow } for _, follower := range followers { key := followKey(req.MediaType, req.TMDBID, follower.UserID, follower.ProfileID) - if req.CompletedAt == nil || !f.followedAt[key].After(*req.CompletedAt) { + if f.followFor[key] == req.ID { delete(f.follows, key) } } 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..217e2cf4e3 --- /dev/null +++ b/migrations/sql/20260929001953_request_follows_request_key.sql @@ -0,0 +1,28 @@ +-- +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 records the request that +-- was open when it was made, and a request's notification goes to its own +-- follows. A follow whose request is gone (a failed request replaced by its +-- requester) has none until the title's next request takes it. +ALTER TABLE public.media_request_follows + ADD COLUMN request_id text REFERENCES public.media_requests (id) ON DELETE SET NULL; +CREATE INDEX media_request_follows_request_idx ON public.media_request_follows (request_id); + +-- Existing follows go to the title's open request, or else to its latest +-- completed request that has not sent its notification. +UPDATE public.media_request_follows f +SET request_id = ( + 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 = 'active' + AND (r.status <> 'completed' OR r.fulfilled_notified_at IS NULL) + ORDER BY r.status = 'completed', r.completed_at DESC NULLS LAST, r.id + LIMIT 1); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP INDEX IF EXISTS public.media_request_follows_request_idx; +ALTER TABLE public.media_request_follows DROP COLUMN IF EXISTS request_id; +-- +goose StatementEnd From 9d29a0a88dcbf8de540a5e31a98c165765e60f14 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 00:52:13 +0000 Subject: [PATCH 05/10] fix(requests): key follows by request so a profile can follow each one A profile following an older completed request of a series could not follow a newer request for other seasons: the key was per title, the insert kept the old row, and the title showed as followed. The backfill also gave every existing follow to the open request, even ones made for an older request. Follows are now keyed by account, profile and request, and the "following" state reads the follow on the title's open request. A new request takes over the follows of the title's failed requests before any it replaces are deleted. The backfill gives each follow the title's latest request created before it, which is the request that was open then. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/architecture/media-requests.md | 11 +-- internal/requests/follows.go | 67 +++++++++++-------- internal/requests/follows_test.go | 24 ++++--- internal/requests/repository.go | 30 +++++---- internal/requests/service_test.go | 29 ++++---- internal/requests/store.go | 6 +- ...0929001953_request_follows_request_key.sql | 61 ++++++++++++----- 7 files changed, 140 insertions(+), 88 deletions(-) diff --git a/docs/architecture/media-requests.md b/docs/architecture/media-requests.md index 02028eda3f..b994c9f9eb 100644 --- a/docs/architecture/media-requests.md +++ b/docs/architecture/media-requests.md @@ -308,11 +308,12 @@ 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 a title and a profile (`media_request_follows`, keyed -by account and profile id, since profile ids repeat across accounts) and records -the request that was open when it was made. A series can have completed -requests still waiting for the library beside a newer open request for other -seasons, and each request's notification goes to its own follows. A follow +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 diff --git a/internal/requests/follows.go b/internal/requests/follows.go index b1f17d13c1..9da55d9d45 100644 --- a/internal/requests/follows.go +++ b/internal/requests/follows.go @@ -3,16 +3,19 @@ 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 records the request that was open when it was -// made, since a series can have completed requests still waiting for the +// 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. A follow survives its request failing: -// the title's next request takes over the follows of a failed or replaced one. +// 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 @@ -100,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 @@ -109,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 @@ -149,7 +152,7 @@ func (r *Repository) FollowTitle(ctx context.Context, mediaType MediaType, tmdbI ), inserted AS ( 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 (media_type, tmdb_id, user_id, profile_id) DO NOTHING + 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 { @@ -170,21 +173,30 @@ func forgetTitleFollows(ctx context.Context, exec requestExecutor, closed *Reque return nil } -// adoptTitleFollows gives a new request the title's follows that have no -// request waiting to tell them: their request failed, or was replaced by its -// requester. The caller creates the request in the same transaction. +// 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. func adoptTitleFollows(ctx context.Context, exec requestExecutor, req *Request) error { if _, err := exec.Exec(ctx, ` - UPDATE media_request_follows f SET request_id = $3 - WHERE f.media_type = $1 AND f.tmdb_id = $2 - AND NOT EXISTS ( - SELECT 1 FROM media_requests r WHERE r.id = f.request_id AND r.outcome <> 'failed') + WITH failed AS ( + SELECT id FROM media_requests + WHERE media_type = $1 AND provider = 'tmdb' AND tmdb_id = $2 AND outcome = 'failed' + ), moved AS ( + DELETE FROM media_request_follows WHERE request_id IN (SELECT id FROM failed) + RETURNING user_id, profile_id, created_at + ) + INSERT INTO media_request_follows (media_type, tmdb_id, user_id, profile_id, request_id, created_at) + SELECT $1, $2, user_id, profile_id, $3, min(created_at) FROM moved + GROUP BY user_id, profile_id + ON CONFLICT (user_id, profile_id, request_id) DO NOTHING `, 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 @@ -195,25 +207,26 @@ 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() } diff --git a/internal/requests/follows_test.go b/internal/requests/follows_test.go index b3516a5b7e..a0afa3266e 100644 --- a/internal/requests/follows_test.go +++ b/internal/requests/follows_test.go @@ -39,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 { @@ -50,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") } } @@ -322,12 +322,12 @@ func TestFollowsDatabase(t *testing.T) { 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.ClearRequestFollowers(ctx, Request{ID: "movie-949", MediaType: MediaTypeMovie, TMDBID: 949}, []Follower{{UserID: a.UserID, ProfileID: a.ProfileID}}); err != nil { @@ -354,7 +354,7 @@ func TestFollowsDatabase(t *testing.T) { 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 { @@ -500,6 +500,14 @@ func TestRequestFollowersDatabase(t *testing.T) { 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) diff --git a/internal/requests/repository.go b/internal/requests/repository.go index a81fdc2489..92e245b25b 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, ` @@ -295,6 +281,22 @@ func (r *Repository) CreateRequest(ctx context.Context, input CreateRequestRecor 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 } diff --git a/internal/requests/service_test.go b/internal/requests/service_test.go index e9ae852e05..7d3224b674 100644 --- a/internal/requests/service_test.go +++ b/internal/requests/service_test.go @@ -2530,8 +2530,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 { @@ -2582,7 +2583,7 @@ func (f *fakeStore) seedFollowForLocked(mediaType MediaType, tmdbID int, viewer if f.followFor == nil { f.followFor = map[string]string{} } - key := followKey(mediaType, tmdbID, viewer.UserID, viewer.ProfileID) + key := followKey(mediaType, tmdbID, viewer.UserID, viewer.ProfileID, requestID) f.follows[key] = Follower{UserID: viewer.UserID, ProfileID: viewer.ProfileID} f.followFor[key] = requestID } @@ -2590,17 +2591,22 @@ func (f *fakeStore) seedFollowForLocked(mediaType MediaType, tmdbID int, viewer 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 @@ -2645,10 +2651,7 @@ func (f *fakeStore) ClearRequestFollowers(_ context.Context, req Request, follow return f.clearErr } for _, follower := range followers { - key := followKey(req.MediaType, req.TMDBID, follower.UserID, follower.ProfileID) - if f.followFor[key] == req.ID { - delete(f.follows, key) - } + 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 79c4f64c2b..d96f1e3330 100644 --- a/internal/requests/store.go +++ b/internal/requests/store.go @@ -76,14 +76,14 @@ 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; ListRequestFollowers and + // 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) + 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; diff --git a/migrations/sql/20260929001953_request_follows_request_key.sql b/migrations/sql/20260929001953_request_follows_request_key.sql index 217e2cf4e3..bc77f4078b 100644 --- a/migrations/sql/20260929001953_request_follows_request_key.sql +++ b/migrations/sql/20260929001953_request_follows_request_key.sql @@ -1,28 +1,53 @@ -- +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 records the request that --- was open when it was made, and a request's notification goes to its own --- follows. A follow whose request is gone (a failed request replaced by its --- requester) has none until the title's next request takes it. -ALTER TABLE public.media_request_follows - ADD COLUMN request_id text REFERENCES public.media_requests (id) ON DELETE SET NULL; -CREATE INDEX media_request_follows_request_idx ON public.media_request_follows (request_id); +-- 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; --- Existing follows go to the title's open request, or else to its latest --- completed request that has not sent its notification. +-- 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. A follow whose +-- request is gone (a failed request its requester replaced) goes to the +-- title's open request, as a new request takes such follows over. UPDATE public.media_request_follows f -SET request_id = ( - 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 = 'active' - AND (r.status <> 'completed' OR r.fulfilled_notified_at IS NULL) - ORDER BY r.status = 'completed', r.completed_at DESC NULLS LAST, r.id - LIMIT 1); +SET request_id = 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 + ORDER BY r.created_at DESC, r.id DESC + LIMIT 1), + (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 = 'active' AND r.status <> 'completed' + LIMIT 1)); +-- 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 -DROP INDEX IF EXISTS public.media_request_follows_request_idx; -ALTER TABLE public.media_request_follows DROP COLUMN IF EXISTS request_id; +-- 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 From 1b46795c29cc4fd983bbe4da497bb9e5d524e37e Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:06:44 +0000 Subject: [PATCH 06/10] fix(requests): move adopted follows in place and backfill replaced requests' follows A new request took over a failed request's follows by deleting and re-inserting them. An unfollow running at the same moment waited on the deleted row, could not see the new one, and returned while the profile still followed the new request. The follows now move with an UPDATE, which the waiting unfollow re-checks and deletes. The backfill gave a follow the title's latest request created before it even when that request was already closed by then, as when the followed request was deleted by its requester's replacement. It now takes that request only if it could still have been open, and otherwise the title's open request. Co-Authored-By: Claude Opus 5.5 (1M context) --- internal/requests/follows.go | 29 +++++---- internal/requests/follows_test.go | 62 +++++++++++++++++++ ...0929001953_request_follows_request_key.sql | 43 +++++++++---- 3 files changed, 110 insertions(+), 24 deletions(-) diff --git a/internal/requests/follows.go b/internal/requests/follows.go index 9da55d9d45..fec62ad96c 100644 --- a/internal/requests/follows.go +++ b/internal/requests/follows.go @@ -177,19 +177,26 @@ func forgetTitleFollows(ctx context.Context, exec requestExecutor, closed *Reque // 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, ` - WITH failed AS ( - SELECT id FROM media_requests - WHERE media_type = $1 AND provider = 'tmdb' AND tmdb_id = $2 AND outcome = 'failed' - ), moved AS ( - DELETE FROM media_request_follows WHERE request_id IN (SELECT id FROM failed) - RETURNING user_id, profile_id, created_at - ) - INSERT INTO media_request_follows (media_type, tmdb_id, user_id, profile_id, request_id, created_at) - SELECT $1, $2, user_id, profile_id, $3, min(created_at) FROM moved - GROUP BY user_id, profile_id - ON CONFLICT (user_id, profile_id, request_id) DO NOTHING + 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) } diff --git a/internal/requests/follows_test.go b/internal/requests/follows_test.go index a0afa3266e..1ab8a0fb08 100644 --- a/internal/requests/follows_test.go +++ b/internal/requests/follows_test.go @@ -586,6 +586,68 @@ func TestNewRequestAdoptsFollowsOfFailedRequestDatabase(t *testing.T) { } } +// 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. diff --git a/migrations/sql/20260929001953_request_follows_request_key.sql b/migrations/sql/20260929001953_request_follows_request_key.sql index bc77f4078b..a5f5d18a11 100644 --- a/migrations/sql/20260929001953_request_follows_request_key.sql +++ b/migrations/sql/20260929001953_request_follows_request_key.sql @@ -7,20 +7,37 @@ 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. A follow whose --- request is gone (a failed request its requester replaced) goes to the --- title's open request, as a new request takes such follows over. +-- 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. Otherwise the follow's own request is gone (a failed +-- request its requester replaced) and the follow goes to the title's open +-- request, as a new request takes such follows over; with none open, it stays +-- with a failed request for the title's next request to take over. UPDATE public.media_request_follows f -SET request_id = 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 - ORDER BY r.created_at DESC, r.id DESC - LIMIT 1), - (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 = 'active' AND r.status <> 'completed' - LIMIT 1)); +SET request_id = CASE + WHEN made_for.outcome = 'active' + AND (made_for.status <> 'completed' OR made_for.completed_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.outcome = 'active' AND r.status <> 'completed' + LIMIT 1), + CASE WHEN made_for.outcome = 'failed' THEN made_for.id END) + 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 + 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; From 8d7a4415d4eb76cdbaaf0a1a028f8c0665c78ac1 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:17:30 +0000 Subject: [PATCH 07/10] build: pin the plugin SDK to v0.19.0 The download-progress change pinned the SDK to its pull request commit, which was squash-merged as v0.19.0 and is no longer fetchable, so the Go jobs could not download the module. v0.19.0 has the same contents. Co-Authored-By: Claude Opus 5.5 (1M context) --- go.mod | 2 +- go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) 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= From 68627132370e093061e07cd006dd554aeb60cb94 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:23:47 +0000 Subject: [PATCH 08/10] fix(requests): backfill a failed request's follows to its replacement A follow made for a request that later failed stayed with the failed request when its replacement had already completed and was still waiting to notify, since the backfill only looked for an open request. It now gives such a follow the title's first request since the follow that still has a notification to send. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../20260929001953_request_follows_request_key.sql | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/migrations/sql/20260929001953_request_follows_request_key.sql b/migrations/sql/20260929001953_request_follows_request_key.sql index a5f5d18a11..1a130c24b4 100644 --- a/migrations/sql/20260929001953_request_follows_request_key.sql +++ b/migrations/sql/20260929001953_request_follows_request_key.sql @@ -9,10 +9,11 @@ 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. Otherwise the follow's own request is gone (a failed --- request its requester replaced) and the follow goes to the title's open --- request, as a new request takes such follows over; with none open, it stays --- with a failed request for the title's next request to take over. +-- earlier than 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 a failed request +-- for the title's next request to take over. UPDATE public.media_request_follows f SET request_id = CASE WHEN made_for.outcome = 'active' @@ -21,7 +22,9 @@ SET request_id = CASE 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.outcome = 'active' AND r.status <> 'completed' + 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) END From ae723e05fbb9cefbccde6cd621d069efafb4761f Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:42:10 +0000 Subject: [PATCH 09/10] fix(requests): keep a follow stamped during its request's completion A completion's timestamp is taken when its transaction begins, so a follow that committed while it ran can carry a later created_at than the request's completed_at. The backfill took that as the request having closed before the follow and moved or dropped the follow. A completed request that has not notified now keeps such a follow when the title has no request created after it; a follow whose request was replaced and deleted always has that replacement after it, so the two cases stay apart. Co-Authored-By: Claude Opus 5.5 (1M context) --- ...0929001953_request_follows_request_key.sql | 23 +++++++++++++------ 1 file changed, 16 insertions(+), 7 deletions(-) diff --git a/migrations/sql/20260929001953_request_follows_request_key.sql b/migrations/sql/20260929001953_request_follows_request_key.sql index 1a130c24b4..972e1e09d0 100644 --- a/migrations/sql/20260929001953_request_follows_request_key.sql +++ b/migrations/sql/20260929001953_request_follows_request_key.sql @@ -9,15 +9,24 @@ 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. 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 a failed request --- for the title's next request to take over. +-- 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 a failed request 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) + 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 @@ -31,7 +40,7 @@ SET request_id = CASE 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.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 From ae616e059e47ec8f88c259352135f1dade0f006f Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:52:03 +0000 Subject: [PATCH 10/10] fix(requests): backfill a follow to a failed replacement request A follow whose request was replaced and deleted, where the replacement has since failed too, was dropped: the fallback looked only at active requests. Such a follow now stays with the title's latest failed request, for the next request to take over. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../20260929001953_request_follows_request_key.sql | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/migrations/sql/20260929001953_request_follows_request_key.sql b/migrations/sql/20260929001953_request_follows_request_key.sql index 972e1e09d0..1d18fc09af 100644 --- a/migrations/sql/20260929001953_request_follows_request_key.sql +++ b/migrations/sql/20260929001953_request_follows_request_key.sql @@ -17,8 +17,8 @@ ALTER TABLE public.media_request_follows ADD COLUMN request_id text; -- 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 a failed request for the title's next request to take --- over. +-- 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' @@ -35,7 +35,12 @@ SET request_id = CASE 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) + 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)