diff --git a/go.mod b/go.mod index 9ab91da8..39daa300 100644 --- a/go.mod +++ b/go.mod @@ -109,7 +109,7 @@ require ( ) require ( - github.com/Silo-Server/silo-plugin-sdk v0.13.1 + github.com/Silo-Server/silo-plugin-sdk v0.13.2 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 c927aadc..3fa83627 100644 --- a/go.sum +++ b/go.sum @@ -6,6 +6,8 @@ github.com/SherClockHolmes/webpush-go v1.4.0 h1:ocnzNKWN23T9nvHi6IfyrQjkIc0oJWv1 github.com/SherClockHolmes/webpush-go v1.4.0/go.mod h1:XSq8pKX11vNV8MJEMwjrlTkxhAj1zKfxmyhdV7Pd6UA= github.com/Silo-Server/silo-plugin-sdk v0.13.1 h1:3vMaV+aPT/vu47CPk/f3nODsjXCp1OseI0QIoLnOvOM= github.com/Silo-Server/silo-plugin-sdk v0.13.1/go.mod h1:etqmxLTwjxpFH9goAjBDfNDoqHMv2/sqUXu8yx3hNfA= +github.com/Silo-Server/silo-plugin-sdk v0.13.2 h1:w7U0mmljVPauKfzRLNKusuiYFpuoEhuCSX4Hp3s9eRw= +github.com/Silo-Server/silo-plugin-sdk v0.13.2/go.mod h1:etqmxLTwjxpFH9goAjBDfNDoqHMv2/sqUXu8yx3hNfA= github.com/abadojack/whatlanggo v1.0.1 h1:19N6YogDnf71CTHm3Mp2qhYfkRdyvbgwWdd2EPxJRG4= github.com/abadojack/whatlanggo v1.0.1/go.mod h1:66WiQbSbJBIlOZMsvbKe5m6pzQovxCH9B/K8tQB2uoc= github.com/agnivade/levenshtein v1.2.1 h1:EHBY3UOn1gwdy/VbFwgo4cxecRznFk7fKWN1KOX7eoM= diff --git a/internal/watchsync/plugin_provider.go b/internal/watchsync/plugin_provider.go index cf00a978..761f8e20 100644 --- a/internal/watchsync/plugin_provider.go +++ b/internal/watchsync/plugin_provider.go @@ -761,6 +761,7 @@ func watchEventFromScrobble(event ScrobbleEvent, operation pluginv1.WatchSyncOpe PositionSeconds: event.PositionSeconds, DurationSeconds: event.DurationSeconds, CompletionPercent: completion, + Completed: event.Completed, ProviderItemKey: event.ProviderItemKey, Media: mediaFromIdentity(event.MediaItemID, event.Kind, "", 0, event.IMDbID, event.TMDBID, event.TVDBID, "", 0, diff --git a/internal/watchsync/plugin_provider_test.go b/internal/watchsync/plugin_provider_test.go index 8c323a5c..b844fffb 100644 --- a/internal/watchsync/plugin_provider_test.go +++ b/internal/watchsync/plugin_provider_test.go @@ -1016,3 +1016,29 @@ func TestPluginProviderForwardsLiveScrobbleLifecycle(t *testing.T) { t.Fatalf("scrobble event = %#v", got) } } + +func TestPluginProviderForwardsAuthoritativeScrobbleCompletion(t *testing.T) { + completed := watchEventFromScrobble(ScrobbleEvent{ + PlaybackSessionID: testPlaybackSessionID, + MediaItemID: testMovieMediaID, + Kind: historyimport.KindMovie, + PositionSeconds: 90, + DurationSeconds: 100, + Completed: true, + }, pluginv1.WatchSyncOperation_WATCH_SYNC_OPERATION_SCROBBLE_STOP) + if !completed.GetCompleted() { + t.Fatal("completed event = false, want true") + } + + incomplete := watchEventFromScrobble(ScrobbleEvent{ + PlaybackSessionID: testPlaybackSessionID, + MediaItemID: testMovieMediaID, + Kind: historyimport.KindMovie, + PositionSeconds: 10, + DurationSeconds: 100, + Completed: false, + }, pluginv1.WatchSyncOperation_WATCH_SYNC_OPERATION_SCROBBLE_STOP) + if incomplete.GetCompleted() { + t.Fatal("incomplete event = true, want false") + } +} diff --git a/internal/watchsync/repository.go b/internal/watchsync/repository.go index 81f5439f..17e78d78 100644 --- a/internal/watchsync/repository.go +++ b/internal/watchsync/repository.go @@ -34,6 +34,7 @@ type Repository interface { ListListEventConnections(ctx context.Context, userID int, profileID string, list ListKind) ([]Connection, error) UpsertHistoryExports(ctx context.Context, exports []HistoryExport) error ListPendingHistoryExports(ctx context.Context, connectionID string, limit int) ([]HistoryExport, error) + ListPendingHistoryExportsByHistoryIDs(ctx context.Context, connectionID string, historyIDs []string) ([]HistoryExport, error) MarkHistoryExportStatus(ctx context.Context, id string, status string, lastError string) error MarkHistoryExportSatisfiedByScrobble(ctx context.Context, connectionID string, historyID string) error UpsertListItemStates(ctx context.Context, states []ListItemState) error @@ -765,6 +766,55 @@ func (r *PostgresRepository) ListPendingHistoryExports(ctx context.Context, conn return exports, nil } +func (r *PostgresRepository) ListPendingHistoryExportsByHistoryIDs( + ctx context.Context, + connectionID string, + historyIDs []string, +) ([]HistoryExport, error) { + if len(historyIDs) == 0 { + return nil, nil + } + rows, err := r.pool.Query(ctx, ` + SELECT id::text, connection_id::text, history_id, media_item_id, watched_at, + provider_item_key, status, attempt_count, last_attempt_at, last_error, created_at, updated_at + FROM watch_provider_history_exports + WHERE connection_id = $1::uuid + AND history_id = ANY($2::text[]) + AND status IN ('pending', 'failed') + AND attempt_count < 5 + ORDER BY watched_at ASC + `, connectionID, historyIDs) + if err != nil { + return nil, fmt.Errorf("list pending history exports by history ids: %w", err) + } + defer rows.Close() + var exports []HistoryExport + for rows.Next() { + var export HistoryExport + if err := rows.Scan( + &export.ID, + &export.ConnectionID, + &export.HistoryID, + &export.MediaItemID, + &export.WatchedAt, + &export.ProviderItemKey, + &export.Status, + &export.AttemptCount, + &export.LastAttemptAt, + &export.LastError, + &export.CreatedAt, + &export.UpdatedAt, + ); err != nil { + return nil, fmt.Errorf("scan pending history export by history ids: %w", err) + } + exports = append(exports, export) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate pending history exports by history ids: %w", err) + } + return exports, nil +} + func (r *PostgresRepository) MarkHistoryExportStatus(ctx context.Context, id string, status string, lastError string) error { _, err := r.pool.Exec(ctx, ` UPDATE watch_provider_history_exports diff --git a/internal/watchsync/service.go b/internal/watchsync/service.go index 0b2e4767..6b40816c 100644 --- a/internal/watchsync/service.go +++ b/internal/watchsync/service.go @@ -1534,7 +1534,14 @@ func (s *Service) exportLocalPlays( if err := s.repo.UpsertHistoryExports(ctx, exports); err != nil { return err } - pending, err := s.repo.ListPendingHistoryExports(ctx, conn.ID, 100) + historyIDs := make([]string, 0, len(exports)) + for _, export := range exports { + historyIDs = append(historyIDs, export.HistoryID) + } + // A live watch event must not wait behind a connection's historical + // backlog. Scheduled sync still drains that backlog oldest-first, while + // this path selects only the events the user just created. + pending, err := s.repo.ListPendingHistoryExportsByHistoryIDs(ctx, conn.ID, historyIDs) if err != nil { return err } diff --git a/internal/watchsync/service_test.go b/internal/watchsync/service_test.go index 0f7f6e11..55dcedfd 100644 --- a/internal/watchsync/service_test.go +++ b/internal/watchsync/service_test.go @@ -4,6 +4,7 @@ import ( "context" "database/sql" "errors" + "fmt" "strconv" "strings" "sync" @@ -346,6 +347,22 @@ func (r *serviceFakeRepo) ListPendingHistoryExports(_ context.Context, connectio return exports, nil } +func (r *serviceFakeRepo) ListPendingHistoryExportsByHistoryIDs(_ context.Context, connectionID string, historyIDs []string) ([]HistoryExport, error) { + wanted := make(map[string]struct{}, len(historyIDs)) + for _, historyID := range historyIDs { + wanted[historyID] = struct{}{} + } + var exports []HistoryExport + for _, export := range r.historyExports { + _, matches := wanted[export.HistoryID] + if export.ConnectionID == connectionID && matches && + (export.Status == historyExportStatusPending || export.Status == historyExportStatusFailed) && export.AttemptCount < 5 { + exports = append(exports, export) + } + } + return exports, nil +} + func (r *serviceFakeRepo) MarkHistoryExportStatus(_ context.Context, id string, status string, lastError string) error { if r.markHistoryStatusErr != nil { return r.markHistoryStatusErr @@ -789,6 +806,7 @@ func (p progressBatchImporterStub) FetchProgressBatch(context.Context, ServerCon type watchedExporterStub struct { exportErr error exportResult ExportResult + exported *[]LocalPlay key string source userstore.WatchHistorySource } @@ -812,7 +830,10 @@ func (p watchedExporterStub) FetchHistory(context.Context, ServerConfig, Connect return nil, nil } -func (p watchedExporterStub) ExportHistory(context.Context, ServerConfig, Connection, []LocalPlay) (ExportResult, error) { +func (p watchedExporterStub) ExportHistory(_ context.Context, _ ServerConfig, _ Connection, plays []LocalPlay) (ExportResult, error) { + if p.exported != nil { + *p.exported = append(*p.exported, plays...) + } return p.exportResult, p.exportErr } @@ -3130,6 +3151,40 @@ func TestServicePluginTransportFailureLeavesExportPending(t *testing.T) { } } +func TestServiceLocalWatchEventBypassesHistoricalExportBacklog(t *testing.T) { + repo := newServiceFakeRepo() + for i := 0; i < 100; i++ { + repo.historyExports = append(repo.historyExports, HistoryExport{ + ID: fmt.Sprintf("old-export-%d", i), + ConnectionID: "conn-1", + HistoryID: fmt.Sprintf("old-history-%d", i), + Status: historyExportStatusPending, + WatchedAt: time.Date(2025, time.January, 1, 0, i, 0, 0, time.UTC), + }) + } + var exported []LocalPlay + service := NewService(repo, NewRegistry()) + play := LocalPlay{ + HistoryID: "new-history", + MediaItemID: testMovieMediaID, + ProviderItemKey: testMovieProviderItemKey, + WatchedAt: time.Date(2026, time.August, 6, 12, 0, 0, 0, time.UTC), + } + err := service.exportLocalPlays(context.Background(), Connection{ID: "conn-1"}, ServerConfig{}, watchedExporterStub{ + exported: &exported, + exportResult: ExportResult{Sent: []string{play.HistoryID}}, + }, []LocalPlay{play}) + if err != nil { + t.Fatal(err) + } + if len(exported) != 1 || exported[0].HistoryID != play.HistoryID { + t.Fatalf("exported = %#v, want only the live event", exported) + } + if got := repo.historyExports[len(repo.historyExports)-1].Status; got != historyExportStatusSent { + t.Fatalf("new export status = %q, want %q", got, historyExportStatusSent) + } +} + func TestServicePluginInvalidCredentialLeavesExportPendingAndRecordsConnectionError(t *testing.T) { repo := newServiceFakeRepo() client := &fakeWatchSyncPluginClient{applyResponse: &pluginv1.WatchSyncApplyEventsResponse{ diff --git a/web/src/pages/settings/WatchProvidersSettings.tsx b/web/src/pages/settings/WatchProvidersSettings.tsx index 828579fb..e5d0af54 100644 --- a/web/src/pages/settings/WatchProvidersSettings.tsx +++ b/web/src/pages/settings/WatchProvidersSettings.tsx @@ -672,14 +672,16 @@ function WatchProviderCard({ providerKey }: { providerKey: string }) { disabled={isBusy} onChange={(checked) => updateConnection.mutate({ export_watched_enabled: checked })} /> - updateConnection.mutate({ export_unwatched_enabled: checked })} - /> + {connection.capabilities.export_unwatched ? ( + updateConnection.mutate({ export_unwatched_enabled: checked })} + /> + ) : null} {connection.capabilities.import_favorites || connection.capabilities.export_favorites ? (