diff --git a/commafeed-server/src/main/java/com/commafeed/backend/dao/FeedEntryDAO.java b/commafeed-server/src/main/java/com/commafeed/backend/dao/FeedEntryDAO.java index 4712e957..0cf2962b 100644 --- a/commafeed-server/src/main/java/com/commafeed/backend/dao/FeedEntryDAO.java +++ b/commafeed-server/src/main/java/com/commafeed/backend/dao/FeedEntryDAO.java @@ -1,6 +1,7 @@ package com.commafeed.backend.dao; import java.time.Instant; +import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; @@ -11,6 +12,7 @@ import jakarta.persistence.EntityManager; import com.commafeed.backend.model.Feed; import com.commafeed.backend.model.FeedEntry; import com.commafeed.backend.model.QFeedEntry; +import com.google.common.collect.Lists; import com.querydsl.core.Tuple; import com.querydsl.core.types.dsl.NumberExpression; import com.querydsl.jpa.impl.JPAQuery; @@ -19,6 +21,7 @@ import com.querydsl.jpa.impl.JPAQuery; public class FeedEntryDAO extends GenericDAO { private static final QFeedEntry ENTRY = QFeedEntry.feedEntry; + private static final int IN_CLAUSE_BATCH_SIZE = 1000; public FeedEntryDAO(EntityManager entityManager) { super(entityManager, FeedEntry.class); @@ -28,8 +31,16 @@ public class FeedEntryDAO extends GenericDAO { return query().select(ENTRY).from(ENTRY).where(ENTRY.guidHash.eq(guidHash), ENTRY.feed.eq(feed)).limit(1).fetchOne(); } - public Set findExistingGuids(Feed feed) { - return new HashSet<>(query().select(ENTRY.guidHash).from(ENTRY).where(ENTRY.feed.eq(feed)).fetch()); + public Set findExistingGuids(Feed feed, Set guidHashes) { + if (guidHashes.isEmpty()) { + return Set.of(); + } + + Set result = new HashSet<>(); + for (List batch : Lists.partition(new ArrayList<>(guidHashes), IN_CLAUSE_BATCH_SIZE)) { + result.addAll(query().select(ENTRY.guidHash).from(ENTRY).where(ENTRY.feed.eq(feed), ENTRY.guidHash.in(batch)).fetch()); + } + return result; } public List findFeedsExceedingCapacity(long maxCapacity, long max, boolean keepStarredEntries) { diff --git a/commafeed-server/src/main/java/com/commafeed/backend/feed/FeedRefreshUpdater.java b/commafeed-server/src/main/java/com/commafeed/backend/feed/FeedRefreshUpdater.java index 8bd1e42b..7440c3fa 100644 --- a/commafeed-server/src/main/java/com/commafeed/backend/feed/FeedRefreshUpdater.java +++ b/commafeed-server/src/main/java/com/commafeed/backend/feed/FeedRefreshUpdater.java @@ -129,8 +129,16 @@ public class FeedRefreshUpdater { Map> insertedUnreadEntriesBySubscription = new HashMap<>(); if (!entries.isEmpty()) { - Set existingGuids = unitOfWork.call(() -> feedEntryDAO.findExistingGuids(feed)); - List newEntries = entries.stream().filter(e -> !existingGuids.contains(Digests.sha1Hex(e.guid()))).toList(); + Map entriesByGuidHash = new HashMap<>(); + for (Entry entry : entries) { + entriesByGuidHash.put(Digests.sha1Hex(entry.guid()), entry); + } + Set existingGuids = unitOfWork.call(() -> feedEntryDAO.findExistingGuids(feed, entriesByGuidHash.keySet())); + List newEntries = entriesByGuidHash.entrySet() + .stream() + .filter(e -> !existingGuids.contains(e.getKey())) + .map(Map.Entry::getValue) + .toList(); List subscriptions = null; for (Entry entry : newEntries) { diff --git a/commafeed-server/src/test/java/com/commafeed/integration/rest/LargeDatasetIT.java b/commafeed-server/src/test/java/com/commafeed/integration/rest/LargeDatasetIT.java index ead46bb3..40105f53 100644 --- a/commafeed-server/src/test/java/com/commafeed/integration/rest/LargeDatasetIT.java +++ b/commafeed-server/src/test/java/com/commafeed/integration/rest/LargeDatasetIT.java @@ -30,6 +30,8 @@ class LargeDatasetIT extends BaseIT { private static final int ENTRIES_PER_FEED = 20; private static final int TOTAL_ENTRIES = FEED_COUNT * ENTRIES_PER_FEED; + private Long firstSubscriptionId; + @BeforeEach void setup() { initialSetup(TestConstants.ADMIN_USERNAME, TestConstants.ADMIN_PASSWORD); @@ -39,7 +41,10 @@ class LargeDatasetIT extends BaseIT { String path = "/feed/" + i; getMockServerClient().when(HttpRequest.request().withMethod("GET").withPath(path)) .respond(HttpResponse.response().withBody(generateFeed(i)).withContentType(MediaType.APPLICATION_XML)); - subscribe("http://localhost:" + getMockServerClient().getPort() + path); + Long subscriptionId = subscribe("http://localhost:" + getMockServerClient().getPort() + path); + if (i == 0) { + firstSubscriptionId = subscriptionId; + } } Awaitility.await().atMost(Duration.ofSeconds(60)).until(() -> getAllEntries().getEntries().size(), count -> count >= TOTAL_ENTRIES); @@ -66,6 +71,18 @@ class LargeDatasetIT extends BaseIT { Assertions.assertTrue(after.getEntries().stream().allMatch(Entry::isRead)); } + @Test + void refreshDoesNotCreateDuplicateEntries() { + Assertions.assertEquals(TOTAL_ENTRIES, getAllEntries().getEntries().size()); + Instant threshold = Instant.now().minus(Duration.ofSeconds(1)); + forceRefreshAllFeeds(); + + Awaitility.await() + .atMost(Duration.ofSeconds(15)) + .until(() -> getSubscription(firstSubscriptionId), f -> f.getLastRefresh().isAfter(threshold)); + Assertions.assertEquals(TOTAL_ENTRIES, getAllEntries().getEntries().size()); + } + @Test void paginationHasMore() { Entries firstPage = RestAssured.given()