Skip to content

Commit 42d2b71

Browse files
dannyvfilmsclaude
andcommitted
Dedupe duplicate plays on Trakt import
- Skip a Trakt-imported movie/episode watch if it falls within 15 minutes of an existing play for the same item, so double-watch entries no longer occur when Plex webhook scrobbling and Trakt import are both enabled (Plex fires at ~90% progress, Trakt waits for playback stop) - Dedupe checks against both pre-existing DB rows and other entries in the same import run Fixes #854 Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
1 parent d7a0c64 commit 42d2b71

2 files changed

Lines changed: 315 additions & 1 deletion

File tree

src/integrations/imports/trakt.py

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import json
22
import logging
33
from collections import defaultdict
4+
from datetime import timedelta
45

56
import requests
67
from django.conf import settings
@@ -24,6 +25,12 @@
2425
BULK_PAGE_SIZE = 1000
2526
TRAKT_UNKNOWN_DATE = "1970-01-01T00:00:00.000Z"
2627

28+
# Plex's webhook fires at ~90% progress while Trakt's scrobbler waits for
29+
# playback to stop, so a webhook-recorded play and its later Trakt-imported
30+
# counterpart can land several minutes apart. Treat plays for the same item
31+
# within this window as the same watch instead of double-counting it.
32+
_DUPLICATE_PLAY_WINDOW = timedelta(minutes=15)
33+
2734

2835
def _parse_watched_at(watched_at: str):
2936
if watched_at == TRAKT_UNKNOWN_DATE:
@@ -383,6 +390,12 @@ def __init__(self, username, user, mode, refresh_token=None):
383390
# does not create duplicate episode history rows.
384391
self.existing_episode_watch_keys = self._get_existing_episode_watch_keys()
385392

393+
# Track existing play timestamps so a play already recorded (e.g. via
394+
# Plex webhook) isn't duplicated by a nearby Trakt-imported play for
395+
# the same item. See _DUPLICATE_PLAY_WINDOW.
396+
self.existing_episode_play_times = self._get_existing_episode_play_times()
397+
self.existing_movie_play_times = self._get_existing_movie_play_times()
398+
386399
# Track media IDs to delete in overwrite mode
387400
self.to_delete = defaultdict(lambda: defaultdict(set))
388401

@@ -431,6 +444,43 @@ def _get_existing_episode_watch_keys(self):
431444
),
432445
)
433446

447+
def _get_existing_episode_play_times(self):
448+
"""Return existing episode play end_dates keyed by (tmdb_id, season, episode)."""
449+
play_times = defaultdict(list)
450+
rows = app.models.Episode.objects.filter(
451+
related_season__user=self.user,
452+
end_date__isnull=False,
453+
).values_list(
454+
"item__media_id",
455+
"item__season_number",
456+
"item__episode_number",
457+
"end_date",
458+
)
459+
for media_id, season_number, episode_number, end_date in rows:
460+
play_times[(media_id, season_number, episode_number)].append(end_date)
461+
return play_times
462+
463+
def _get_existing_movie_play_times(self):
464+
"""Return existing movie play end_dates keyed by tmdb_id."""
465+
play_times = defaultdict(list)
466+
rows = app.models.Movie.objects.filter(
467+
user=self.user,
468+
end_date__isnull=False,
469+
).values_list("item__media_id", "end_date")
470+
for media_id, end_date in rows:
471+
play_times[media_id].append(end_date)
472+
return play_times
473+
474+
@staticmethod
475+
def _is_duplicate_play(existing_times, candidate_dt):
476+
"""Return True if candidate_dt is within _DUPLICATE_PLAY_WINDOW of an existing play."""
477+
if candidate_dt is None:
478+
return False
479+
return any(
480+
abs(candidate_dt - existing_dt) <= _DUPLICATE_PLAY_WINDOW
481+
for existing_dt in existing_times
482+
)
483+
434484
def _raise_for_user_error(self, error):
435485
"""Translate a Trakt HTTP error about the user into a MediaImportError."""
436486
if error.response.status_code == requests.codes.not_found:
@@ -686,6 +736,17 @@ def process_watched_movie(self, entry):
686736
watched_at = entry["watched_at"]
687737
watched_at_dt = _parse_watched_at(watched_at)
688738

739+
if self._is_duplicate_play(
740+
self.existing_movie_play_times[tmdb_id],
741+
watched_at_dt,
742+
):
743+
logger.debug(
744+
"Skipping Trakt movie watch for %s at %s: duplicate of an existing play",
745+
movie["title"],
746+
watched_at,
747+
)
748+
return
749+
689750
key = f"{tmdb_id}"
690751

691752
movie_obj = app.models.Movie(
@@ -700,6 +761,8 @@ def process_watched_movie(self, entry):
700761

701762
self.media_instances[MediaTypes.MOVIE.value][key].append(movie_obj)
702763
self.bulk_media[MediaTypes.MOVIE.value].append(movie_obj)
764+
if watched_at_dt is not None:
765+
self.existing_movie_play_times[tmdb_id].append(watched_at_dt)
703766

704767
def process_watched_episode(self, entry):
705768
"""Process a single episode watch event."""
@@ -732,6 +795,21 @@ def process_watched_episode(self, entry):
732795
)
733796
return
734797

798+
episode_key = (tmdb_id, season_number, episode_number)
799+
if self._is_duplicate_play(
800+
self.existing_episode_play_times[episode_key],
801+
watched_at_dt,
802+
):
803+
logger.debug(
804+
"Skipping Trakt episode watch for %s S%sE%s at %s: "
805+
"duplicate of an existing play",
806+
show["title"],
807+
season_number,
808+
episode_number,
809+
watched_at,
810+
)
811+
return
812+
735813
tv_exists = (
736814
tmdb_id in self.existing_media[MediaTypes.TV.value][Sources.TMDB.value]
737815
)
@@ -875,6 +953,8 @@ def process_watched_episode(self, entry):
875953
self.media_instances[MediaTypes.EPISODE.value][ep_key].append(episode_obj)
876954
self.bulk_media[MediaTypes.EPISODE.value].append(episode_obj)
877955
self.existing_episode_watch_keys.add(episode_watch_key)
956+
if watched_at_dt is not None:
957+
self.existing_episode_play_times[episode_key].append(watched_at_dt)
878958

879959
# Update status if this is the last episode, but only for rows Floppy
880960
# just created (or an explicit overwrite re-sync) — never clobber the

src/integrations/tests/imports/test_trakt.py

Lines changed: 235 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -62,10 +62,98 @@ def test_process_watched_movie(self, mock_get_metadata):
6262
movie_obj = trakt_importer.bulk_media[MediaTypes.MOVIE.value][0]
6363
self.assertEqual(movie_obj.progress, 1)
6464

65-
# Process the same movie again to test repeat handling
65+
# Reprocessing the exact same entry is a duplicate play (issue #854)
66+
# and must not create a second row.
6667
trakt_importer.process_watched_movie(movie_entry)
68+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.MOVIE.value]), 1)
69+
70+
@patch("integrations.imports.trakt.TraktImporter._get_metadata")
71+
def test_process_watched_movie_dedupes_nearby_play(self, mock_get_metadata):
72+
"""A movie watch within the dedupe window of an existing play is skipped."""
73+
mock_get_metadata.return_value = {
74+
"title": "Test Movie",
75+
"image": "movie_image.jpg",
76+
}
77+
78+
trakt_importer = TraktImporter("test", self.user, "new")
79+
trakt_importer.process_watched_movie(
80+
{
81+
"type": "movie",
82+
"movie": {"title": "Test Movie", "ids": {"tmdb": 67890}},
83+
"watched_at": "2023-01-02T00:00:00.000Z",
84+
},
85+
)
86+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.MOVIE.value]), 1)
87+
88+
# 10 minutes later, same movie: within the 15 minute dedupe window.
89+
trakt_importer.process_watched_movie(
90+
{
91+
"type": "movie",
92+
"movie": {"title": "Test Movie", "ids": {"tmdb": 67890}},
93+
"watched_at": "2023-01-02T00:10:00.000Z",
94+
},
95+
)
96+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.MOVIE.value]), 1)
97+
98+
# 1 day later, same movie: a legitimate rewatch outside the window.
99+
trakt_importer.process_watched_movie(
100+
{
101+
"type": "movie",
102+
"movie": {"title": "Test Movie", "ids": {"tmdb": 67890}},
103+
"watched_at": "2023-01-03T00:00:00.000Z",
104+
},
105+
)
67106
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.MOVIE.value]), 2)
68107

108+
# A different movie at a nearby time is not a duplicate.
109+
mock_get_metadata.return_value = {
110+
"title": "Other Movie",
111+
"image": "movie_image.jpg",
112+
}
113+
trakt_importer.process_watched_movie(
114+
{
115+
"type": "movie",
116+
"movie": {"title": "Other Movie", "ids": {"tmdb": 11111}},
117+
"watched_at": "2023-01-03T00:05:00.000Z",
118+
},
119+
)
120+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.MOVIE.value]), 3)
121+
122+
@patch("integrations.imports.trakt.TraktImporter._get_metadata")
123+
def test_process_watched_movie_dedupes_against_existing_db_play(
124+
self,
125+
mock_get_metadata,
126+
):
127+
"""A Trakt-imported play is skipped if it's near an existing DB play (e.g. webhook)."""
128+
item = Item.objects.get_or_create(
129+
media_id="67890",
130+
source=Sources.TMDB.value,
131+
media_type=MediaTypes.MOVIE.value,
132+
defaults={"title": "Test Movie"},
133+
)[0]
134+
Movie.objects.create(
135+
item=item,
136+
user=self.user,
137+
end_date="2023-01-02T00:00:00Z",
138+
status=Status.COMPLETED.value,
139+
progress=1,
140+
)
141+
142+
mock_get_metadata.return_value = {
143+
"title": "Test Movie",
144+
"image": "movie_image.jpg",
145+
}
146+
trakt_importer = TraktImporter("test", self.user, "new")
147+
trakt_importer.process_watched_movie(
148+
{
149+
"type": "movie",
150+
"movie": {"title": "Test Movie", "ids": {"tmdb": 67890}},
151+
"watched_at": "2023-01-02T00:12:00.000Z",
152+
},
153+
)
154+
155+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.MOVIE.value]), 0)
156+
69157
@patch("integrations.imports.trakt.TraktImporter._get_metadata")
70158
def test_process_watched_episode(self, mock_get_metadata):
71159
"""Test processing an episode entry."""
@@ -132,6 +220,152 @@ def mock_metadata_side_effect(media_type, _, __, ___=None):
132220
)
133221
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.EPISODE.value]), 2)
134222

223+
@patch("integrations.imports.trakt.TraktImporter._get_metadata")
224+
def test_process_watched_episode_dedupes_nearby_play(self, mock_get_metadata):
225+
"""An episode watch within the dedupe window of an existing play is skipped."""
226+
227+
def mock_metadata_side_effect(media_type, _, __, ___=None):
228+
if media_type == MediaTypes.TV.value:
229+
return {
230+
"title": "Test Show",
231+
"image": "tv_image.jpg",
232+
"last_episode_season": 1,
233+
"max_progress": 1,
234+
}
235+
if media_type == MediaTypes.SEASON.value:
236+
return {
237+
"title": "Season 1",
238+
"image": "season_image.jpg",
239+
"episodes": [
240+
{
241+
"episode_number": 1,
242+
"still_path": "/still.jpg",
243+
"title": "Pilot Episode Title",
244+
},
245+
{
246+
"episode_number": 2,
247+
"still_path": "/still2.jpg",
248+
"title": "Episode 2 Title",
249+
},
250+
],
251+
"max_progress": 2,
252+
}
253+
return None
254+
255+
mock_get_metadata.side_effect = mock_metadata_side_effect
256+
257+
episode_entry = {
258+
"type": "episode",
259+
"episode": {"season": 1, "number": 1, "title": "Pilot"},
260+
"show": {"title": "Test Show", "ids": {"tmdb": 12345}},
261+
"watched_at": "2023-01-01T00:00:00.000Z",
262+
}
263+
264+
trakt_importer = TraktImporter("testuser", self.user, "new")
265+
trakt_importer.process_watched_episode(episode_entry)
266+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.EPISODE.value]), 1)
267+
268+
# 10 minutes later, same episode: within the 15 minute dedupe window.
269+
trakt_importer.process_watched_episode(
270+
{
271+
**episode_entry,
272+
"watched_at": "2023-01-01T00:10:00.000Z",
273+
},
274+
)
275+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.EPISODE.value]), 1)
276+
277+
# A different episode of the same show at a nearby time is not a duplicate.
278+
trakt_importer.process_watched_episode(
279+
{
280+
"type": "episode",
281+
"episode": {"season": 1, "number": 2, "title": "Episode 2"},
282+
"show": {"title": "Test Show", "ids": {"tmdb": 12345}},
283+
"watched_at": "2023-01-01T00:12:00.000Z",
284+
},
285+
)
286+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.EPISODE.value]), 2)
287+
288+
@patch("integrations.imports.trakt.TraktImporter._get_metadata")
289+
def test_process_watched_episode_dedupes_against_existing_db_play(
290+
self,
291+
mock_get_metadata,
292+
):
293+
"""A Trakt-imported episode play is skipped if it's near an existing DB play."""
294+
tv_item = Item.objects.get_or_create(
295+
media_id="12345",
296+
source=Sources.TMDB.value,
297+
media_type=MediaTypes.TV.value,
298+
defaults={"title": "Test Show"},
299+
)[0]
300+
tv_obj = TV.objects.create(
301+
item=tv_item,
302+
user=self.user,
303+
status=Status.IN_PROGRESS.value,
304+
)
305+
season_item = Item.objects.get_or_create(
306+
media_id="12345",
307+
source=Sources.TMDB.value,
308+
media_type=MediaTypes.SEASON.value,
309+
season_number=1,
310+
defaults={"title": "Season 1"},
311+
)[0]
312+
season_obj = Season.objects.create(
313+
item=season_item,
314+
related_tv=tv_obj,
315+
user=self.user,
316+
status=Status.IN_PROGRESS.value,
317+
)
318+
episode_item = Item.objects.get_or_create(
319+
media_id="12345",
320+
source=Sources.TMDB.value,
321+
media_type=MediaTypes.EPISODE.value,
322+
season_number=1,
323+
episode_number=1,
324+
defaults={"title": "Pilot"},
325+
)[0]
326+
Episode.objects.create(
327+
item=episode_item,
328+
related_season=season_obj,
329+
end_date="2023-01-01T00:00:00Z",
330+
)
331+
332+
def mock_metadata_side_effect(media_type, _, __, ___=None):
333+
if media_type == MediaTypes.TV.value:
334+
return {
335+
"title": "Test Show",
336+
"image": "tv_image.jpg",
337+
"last_episode_season": 1,
338+
"max_progress": 1,
339+
}
340+
if media_type == MediaTypes.SEASON.value:
341+
return {
342+
"title": "Season 1",
343+
"image": "season_image.jpg",
344+
"episodes": [
345+
{
346+
"episode_number": 1,
347+
"still_path": "/still.jpg",
348+
"title": "Pilot Episode Title",
349+
},
350+
],
351+
"max_progress": 1,
352+
}
353+
return None
354+
355+
mock_get_metadata.side_effect = mock_metadata_side_effect
356+
357+
trakt_importer = TraktImporter("testuser", self.user, "new")
358+
trakt_importer.process_watched_episode(
359+
{
360+
"type": "episode",
361+
"episode": {"season": 1, "number": 1, "title": "Pilot"},
362+
"show": {"title": "Test Show", "ids": {"tmdb": 12345}},
363+
"watched_at": "2023-01-01T00:12:00.000Z",
364+
},
365+
)
366+
367+
self.assertEqual(len(trakt_importer.bulk_media[MediaTypes.EPISODE.value]), 0)
368+
135369
@patch("integrations.imports.trakt.TraktImporter._get_metadata")
136370
def test_process_watched_episode_existing_show_imports_new_episode(
137371
self,

0 commit comments

Comments
 (0)