Fix duplicate feed items from concurrent syncs

A sync button press racing a page-reload sync could run two sync
requests concurrently. create_feed_item used a check-then-insert
pattern (query for an existing item, then insert if none found), so
both requests could pass the check before either inserted, creating
duplicate items.

- Add a UNIQUE (feed_id, url) constraint on feed_item, keyed on the
  article's link rather than its title since that's RSS's stable
  article identity (a feed editing a headline after publishing, same
  link, would otherwise dupe under title-based matching).
- The migration first collapses any duplicates already created by the
  race, propagating read=true onto the surviving row when any
  duplicate in its group was already read, so the cleanup can't make
  an already-read article look unread.
- create_feed_item now does a single atomic
  INSERT ... ON CONFLICT (feed_id, url) DO NOTHING instead of
  select-then-insert, closing the race entirely.
- Add a test that races two threads on separate connections inserting
  the same item and asserts only one row survives.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012CHcxDSPbHhe7sLJQBRVC9
This commit is contained in:
2026-09-10 16:04:34 +02:00
co-authored by Claude Sonnet 5
parent e2fb32f898
commit 6c8b4e9e32
3 changed files with 115 additions and 19 deletions
@@ -0,0 +1,3 @@
-- This file should undo anything in `up.sql`
ALTER TABLE feed_item
DROP CONSTRAINT feed_item_feed_id_url_unique;
@@ -0,0 +1,36 @@
-- Your SQL goes here
-- Concurrent syncs for the same feed could previously race past the
-- application-level "does this article already exist" check and both
-- insert, producing duplicate feed items. Before enforcing uniqueness,
-- collapse any duplicates that already exist, keyed on the article's link
-- (rather than its title) since that's the stable identity RSS items are
-- built around — some feeds edit a headline after publishing while keeping
-- the same link, which title-based matching would treat as a new article.
-- If any duplicate in a group was already read, propagate that onto the
-- surviving (lowest-id) row first, so this one-time cleanup can't make an
-- already-read article look unread again.
UPDATE feed_item AS keep
SET read = true
FROM feed_item AS dupe
WHERE keep.feed_id = dupe.feed_id
AND keep.url = dupe.url
AND keep.id <> dupe.id
AND keep.id = (
SELECT MIN(candidate.id)
FROM feed_item AS candidate
WHERE candidate.feed_id = dupe.feed_id
AND candidate.url = dupe.url
)
AND dupe.read = true;
-- Keep the oldest (lowest id) row of each duplicate set.
DELETE FROM feed_item a
USING feed_item b
WHERE a.feed_id = b.feed_id
AND a.url = b.url
AND a.id > b.id;
ALTER TABLE feed_item
ADD CONSTRAINT feed_item_feed_id_url_unique UNIQUE (feed_id, url);
+76 -19
View File
@@ -4,8 +4,7 @@ use crate::error::AppError;
use crate::json_serialization::user::JsonUser;
use crate::models::feed::rss_feed::Feed;
use crate::models::feed_item::new_feed_item::NewFeedItem;
use crate::models::feed_item::rss_feed_item::FeedItem;
use crate::schema::feed_item::{feed_id, title};
use crate::schema::feed_item::{feed_id, url};
use crate::{
database::establish_connection,
schema::{
@@ -143,24 +142,29 @@ fn create_feed_item(item: Item, feed: &Feed, connection: &mut PgConnection) -> a
}
}
let existing_item: Vec<FeedItem> = feed_item::table
.filter(feed_id.eq(feed.id))
.filter(title.eq(&item_title))
.load(connection)?;
let new_feed_item = NewFeedItem::new(
feed.id,
content.clone(),
item_title.clone(),
item.link.expect("checked above"),
Some(time),
);
if existing_item.is_empty() {
let new_feed_item = NewFeedItem::new(
feed.id,
content.clone(),
item_title.clone(),
item.link.expect("checked above"),
Some(time),
);
let insert_result = diesel::insert_into(feed_item::table)
.values(&new_feed_item)
.execute(connection);
// `on_conflict` on the (feed_id, url) unique constraint makes this
// insert idempotent, so two syncs for the same feed running concurrently
// (e.g. a sync button press racing a page-reload sync) can't both pass a
// check-then-insert race and create duplicate items. Keying on the
// article's link rather than its title also handles feeds that edit a
// headline after publishing while keeping the same link — the link is
// the stable identity.
let inserted_rows = diesel::insert_into(feed_item::table)
.values(&new_feed_item)
.on_conflict((feed_id, url))
.do_nothing()
.execute(connection)?;
log::info!("Insert Result: {:?}", insert_result);
if inserted_rows > 0 {
log::info!("Inserted item: {}", item_title);
} else {
log::info!("Item {} already exists.", item_title);
}
@@ -210,10 +214,11 @@ pub async fn sync(
#[cfg(test)]
mod tests {
use crate::models::feed::new_feed::NewFeed;
use crate::models::feed_item::rss_feed_item::FeedItem;
use crate::models::user::new_user::NewUser;
use crate::models::user::rss_user::User;
use crate::schema::users;
use crate::test_helpers::unique_suffix;
use crate::test_helpers::{delete_feed, delete_user, insert_feed, insert_user, unique_suffix};
use chrono::Duration;
use super::*;
@@ -440,6 +445,58 @@ mod tests {
.ok();
}
#[actix_web::test]
async fn create_feed_item_is_race_safe_under_concurrent_syncs() {
let mut connection = establish_connection();
let suffix = unique_suffix();
let user = insert_user(&mut connection, "secret");
let feed = insert_feed(&mut connection, user.id);
let mut item = Item::default();
item.set_title(Some(format!("Race test article {suffix}")));
item.set_link(Some(format!("https://example.test/race/{suffix}")));
item.set_content(Some("<p>Hello world</p>".to_string()));
// Simulate two concurrent syncs for the same feed (e.g. a sync
// button press racing a page-reload sync) — each on its own
// connection/thread, released together so they genuinely overlap —
// both inserting the same item.
let barrier = std::sync::Arc::new(std::sync::Barrier::new(2));
let handles: Vec<_> = (0..2)
.map(|_| {
let feed = feed.clone();
let item = item.clone();
let barrier = barrier.clone();
std::thread::spawn(move || {
let mut connection = establish_connection();
barrier.wait();
create_feed_item(item, &feed, &mut connection)
})
})
.collect();
for handle in handles {
handle.join().unwrap().unwrap();
}
let items: Vec<FeedItem> = feed_item::table
.filter(feed_id.eq(feed.id))
.load(&mut connection)
.unwrap();
assert_eq!(
1,
items.len(),
"concurrent syncs must not create duplicate feed items"
);
diesel::delete(feed_item::table.filter(feed_id.eq(feed.id)))
.execute(&mut connection)
.ok();
delete_feed(&mut connection, feed.id);
delete_user(&mut connection, user.id);
}
#[actix_web::test]
async fn create_feed_item_strips_onerror_from_feed_image() {
let mut connection = establish_connection();