fix(watchsync): preserve provider event semantics
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -672,14 +672,16 @@ function WatchProviderCard({ providerKey }: { providerKey: string }) {
|
||||
disabled={isBusy}
|
||||
onChange={(checked) => updateConnection.mutate({ export_watched_enabled: checked })}
|
||||
/>
|
||||
<ToggleRow
|
||||
id={`watch-provider-${providerKey}-export-unwatched`}
|
||||
label="Send unwatched changes"
|
||||
description="When you mark something unwatched, remove matching history from this provider."
|
||||
checked={connection.export_unwatched_enabled}
|
||||
disabled={isBusy}
|
||||
onChange={(checked) => updateConnection.mutate({ export_unwatched_enabled: checked })}
|
||||
/>
|
||||
{connection.capabilities.export_unwatched ? (
|
||||
<ToggleRow
|
||||
id={`watch-provider-${providerKey}-export-unwatched`}
|
||||
label="Send unwatched changes"
|
||||
description="When you mark something unwatched, remove matching history from this provider."
|
||||
checked={connection.export_unwatched_enabled}
|
||||
disabled={isBusy}
|
||||
onChange={(checked) => updateConnection.mutate({ export_unwatched_enabled: checked })}
|
||||
/>
|
||||
) : null}
|
||||
{connection.capabilities.import_favorites ||
|
||||
connection.capabilities.export_favorites ? (
|
||||
<ToggleRow
|
||||
|
||||
Reference in New Issue
Block a user