Merge branch 'fix/concurrent-sync-duplicate-items'
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012CHcxDSPbHhe7sLJQBRVC9
This commit is contained in:
@@ -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
@@ -4,8 +4,7 @@ use crate::error::AppError;
|
|||||||
use crate::json_serialization::user::JsonUser;
|
use crate::json_serialization::user::JsonUser;
|
||||||
use crate::models::feed::rss_feed::Feed;
|
use crate::models::feed::rss_feed::Feed;
|
||||||
use crate::models::feed_item::new_feed_item::NewFeedItem;
|
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, url};
|
||||||
use crate::schema::feed_item::{feed_id, title};
|
|
||||||
use crate::{
|
use crate::{
|
||||||
database::establish_connection,
|
database::establish_connection,
|
||||||
schema::{
|
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
|
let new_feed_item = NewFeedItem::new(
|
||||||
.filter(feed_id.eq(feed.id))
|
feed.id,
|
||||||
.filter(title.eq(&item_title))
|
content.clone(),
|
||||||
.load(connection)?;
|
item_title.clone(),
|
||||||
|
item.link.expect("checked above"),
|
||||||
|
Some(time),
|
||||||
|
);
|
||||||
|
|
||||||
if existing_item.is_empty() {
|
// `on_conflict` on the (feed_id, url) unique constraint makes this
|
||||||
let new_feed_item = NewFeedItem::new(
|
// insert idempotent, so two syncs for the same feed running concurrently
|
||||||
feed.id,
|
// (e.g. a sync button press racing a page-reload sync) can't both pass a
|
||||||
content.clone(),
|
// check-then-insert race and create duplicate items. Keying on the
|
||||||
item_title.clone(),
|
// article's link rather than its title also handles feeds that edit a
|
||||||
item.link.expect("checked above"),
|
// headline after publishing while keeping the same link — the link is
|
||||||
Some(time),
|
// the stable identity.
|
||||||
);
|
let inserted_rows = diesel::insert_into(feed_item::table)
|
||||||
let insert_result = diesel::insert_into(feed_item::table)
|
.values(&new_feed_item)
|
||||||
.values(&new_feed_item)
|
.on_conflict((feed_id, url))
|
||||||
.execute(connection);
|
.do_nothing()
|
||||||
|
.execute(connection)?;
|
||||||
|
|
||||||
log::info!("Insert Result: {:?}", insert_result);
|
if inserted_rows > 0 {
|
||||||
|
log::info!("Inserted item: {}", item_title);
|
||||||
} else {
|
} else {
|
||||||
log::info!("Item {} already exists.", item_title);
|
log::info!("Item {} already exists.", item_title);
|
||||||
}
|
}
|
||||||
@@ -210,10 +214,11 @@ pub async fn sync(
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use crate::models::feed::new_feed::NewFeed;
|
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::new_user::NewUser;
|
||||||
use crate::models::user::rss_user::User;
|
use crate::models::user::rss_user::User;
|
||||||
use crate::schema::users;
|
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 chrono::Duration;
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
@@ -440,6 +445,58 @@ mod tests {
|
|||||||
.ok();
|
.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]
|
#[actix_web::test]
|
||||||
async fn create_feed_item_strips_onerror_from_feed_image() {
|
async fn create_feed_item_strips_onerror_from_feed_image() {
|
||||||
let mut connection = establish_connection();
|
let mut connection = establish_connection();
|
||||||
|
|||||||
Reference in New Issue
Block a user