| 1 |
use super::*; |
| 2 |
use crate::id_types::BusserId; |
| 3 |
use chrono::{DateTime, Duration}; |
| 4 |
use serde_json::json; |
| 5 |
|
| 6 |
|
| 7 |
async fn test_db() -> SqlitePool { |
| 8 |
let pool = SqlitePool::connect("sqlite::memory:").await.unwrap(); |
| 9 |
sqlx::migrate!("../../migrations/sqlite") |
| 10 |
.run(&pool) |
| 11 |
.await |
| 12 |
.unwrap(); |
| 13 |
pool |
| 14 |
} |
| 15 |
|
| 16 |
|
| 17 |
async fn make_feed(pool: &SqlitePool, busser_id: &str, name: &str) -> DbFeed { |
| 18 |
FeedsRepository::new(pool.clone()) |
| 19 |
.create(CreateFeed { |
| 20 |
busser_id: BusserId::new(busser_id), |
| 21 |
name: name.to_string(), |
| 22 |
config: json!({}), |
| 23 |
}) |
| 24 |
.await |
| 25 |
.unwrap() |
| 26 |
} |
| 27 |
|
| 28 |
|
| 29 |
async fn make_item(pool: &SqlitePool, feed: &DbFeed, external_id: &str) -> DbFeedItem { |
| 30 |
make_item_at(pool, feed, external_id, Utc::now()).await |
| 31 |
} |
| 32 |
|
| 33 |
|
| 34 |
async fn make_item_at( |
| 35 |
pool: &SqlitePool, |
| 36 |
feed: &DbFeed, |
| 37 |
external_id: &str, |
| 38 |
published_at: DateTime<Utc>, |
| 39 |
) -> DbFeedItem { |
| 40 |
ItemsRepository::new(pool.clone()) |
| 41 |
.upsert(CreateFeedItem { |
| 42 |
external_id: external_id.to_string(), |
| 43 |
feed_id: feed.id, |
| 44 |
busser_id: feed.busser_id.clone(), |
| 45 |
bite_author: "author".to_string(), |
| 46 |
bite_text: format!("Item {external_id}"), |
| 47 |
bite_secondary: None, |
| 48 |
bite_indicator: None, |
| 49 |
title: Some(format!("Title {external_id}")), |
| 50 |
body: None, |
| 51 |
url: None, |
| 52 |
media: vec![], |
| 53 |
published_at, |
| 54 |
source_name: "test".to_string(), |
| 55 |
score: None, |
| 56 |
tags: vec![], |
| 57 |
actions: vec![], |
| 58 |
}) |
| 59 |
.await |
| 60 |
.unwrap() |
| 61 |
} |
| 62 |
|
| 63 |
|
| 64 |
|
| 65 |
#[tokio::test] |
| 66 |
async fn feeds_create_and_get() { |
| 67 |
let pool = test_db().await; |
| 68 |
let feed = make_feed(&pool, "rss", "My Feed").await; |
| 69 |
|
| 70 |
assert_eq!(feed.busser_id, "rss"); |
| 71 |
assert_eq!(feed.name, "My Feed"); |
| 72 |
assert!(feed.enabled); |
| 73 |
assert!(feed.last_fetch.is_none()); |
| 74 |
|
| 75 |
let fetched = FeedsRepository::new(pool.clone()) |
| 76 |
.get(feed.id) |
| 77 |
.await |
| 78 |
.unwrap() |
| 79 |
.expect("feed should exist"); |
| 80 |
assert_eq!(fetched.id, feed.id); |
| 81 |
assert_eq!(fetched.name, "My Feed"); |
| 82 |
} |
| 83 |
|
| 84 |
#[tokio::test] |
| 85 |
async fn feeds_get_nonexistent_returns_none() { |
| 86 |
let pool = test_db().await; |
| 87 |
let result = FeedsRepository::new(pool.clone()) |
| 88 |
.get(FeedId::new()) |
| 89 |
.await |
| 90 |
.unwrap(); |
| 91 |
assert!(result.is_none()); |
| 92 |
} |
| 93 |
|
| 94 |
#[tokio::test] |
| 95 |
async fn feeds_list_all_returns_created() { |
| 96 |
let pool = test_db().await; |
| 97 |
make_feed(&pool, "rss", "Beta Feed").await; |
| 98 |
make_feed(&pool, "rss", "Alpha Feed").await; |
| 99 |
|
| 100 |
let all = FeedsRepository::new(pool.clone()).list_all().await.unwrap(); |
| 101 |
assert_eq!(all.len(), 2); |
| 102 |
assert_eq!(all[0].name, "Alpha Feed"); |
| 103 |
assert_eq!(all[1].name, "Beta Feed"); |
| 104 |
} |
| 105 |
|
| 106 |
#[tokio::test] |
| 107 |
async fn feeds_get_by_busser_filters() { |
| 108 |
let pool = test_db().await; |
| 109 |
make_feed(&pool, "rss", "RSS Feed").await; |
| 110 |
make_feed(&pool, "hn", "HN Feed").await; |
| 111 |
make_feed(&pool, "rss", "RSS Feed 2").await; |
| 112 |
|
| 113 |
let rss = FeedsRepository::new(pool.clone()) |
| 114 |
.get_by_busser("rss") |
| 115 |
.await |
| 116 |
.unwrap(); |
| 117 |
assert_eq!(rss.len(), 2); |
| 118 |
for f in &rss { |
| 119 |
assert_eq!(f.busser_id, "rss"); |
| 120 |
} |
| 121 |
|
| 122 |
let hn = FeedsRepository::new(pool.clone()) |
| 123 |
.get_by_busser("hn") |
| 124 |
.await |
| 125 |
.unwrap(); |
| 126 |
assert_eq!(hn.len(), 1); |
| 127 |
assert_eq!(hn[0].busser_id, "hn"); |
| 128 |
} |
| 129 |
|
| 130 |
#[tokio::test] |
| 131 |
async fn feeds_list_enabled_excludes_disabled() { |
| 132 |
let pool = test_db().await; |
| 133 |
let feeds_repo = FeedsRepository::new(pool.clone()); |
| 134 |
let feed = make_feed(&pool, "rss", "Disabled Feed").await; |
| 135 |
make_feed(&pool, "rss", "Enabled Feed").await; |
| 136 |
|
| 137 |
feeds_repo.set_enabled(feed.id, false).await.unwrap(); |
| 138 |
|
| 139 |
let enabled = feeds_repo.list_enabled().await.unwrap(); |
| 140 |
assert_eq!(enabled.len(), 1); |
| 141 |
assert_eq!(enabled[0].name, "Enabled Feed"); |
| 142 |
} |
| 143 |
|
| 144 |
#[tokio::test] |
| 145 |
async fn feeds_set_enabled_toggle() { |
| 146 |
let pool = test_db().await; |
| 147 |
let feeds_repo = FeedsRepository::new(pool.clone()); |
| 148 |
let feed = make_feed(&pool, "rss", "Toggle Feed").await; |
| 149 |
assert!(feed.enabled); |
| 150 |
|
| 151 |
feeds_repo.set_enabled(feed.id, false).await.unwrap(); |
| 152 |
let updated = feeds_repo.get(feed.id).await.unwrap().unwrap(); |
| 153 |
assert!(!updated.enabled); |
| 154 |
|
| 155 |
feeds_repo.set_enabled(feed.id, true).await.unwrap(); |
| 156 |
let updated = feeds_repo.get(feed.id).await.unwrap().unwrap(); |
| 157 |
assert!(updated.enabled); |
| 158 |
} |
| 159 |
|
| 160 |
#[tokio::test] |
| 161 |
async fn feeds_update_last_fetch_sets_timestamp() { |
| 162 |
let pool = test_db().await; |
| 163 |
let feeds_repo = FeedsRepository::new(pool.clone()); |
| 164 |
let feed = make_feed(&pool, "rss", "Fetch Feed").await; |
| 165 |
assert!(feed.last_fetch.is_none()); |
| 166 |
|
| 167 |
feeds_repo.update_last_fetch(feed.id).await.unwrap(); |
| 168 |
let updated = feeds_repo.get(feed.id).await.unwrap().unwrap(); |
| 169 |
assert!(updated.last_fetch.is_some()); |
| 170 |
} |
| 171 |
|
| 172 |
#[tokio::test] |
| 173 |
async fn feeds_delete_removes_feed() { |
| 174 |
let pool = test_db().await; |
| 175 |
let feeds_repo = FeedsRepository::new(pool.clone()); |
| 176 |
let feed = make_feed(&pool, "rss", "Doomed Feed").await; |
| 177 |
|
| 178 |
feeds_repo.delete(feed.id).await.unwrap(); |
| 179 |
let result = feeds_repo.get(feed.id).await.unwrap(); |
| 180 |
assert!(result.is_none()); |
| 181 |
} |
| 182 |
|
| 183 |
|
| 184 |
|
| 185 |
#[tokio::test] |
| 186 |
async fn items_upsert_and_get() { |
| 187 |
let pool = test_db().await; |
| 188 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 189 |
let item = make_item(&pool, &feed, "rss:1").await; |
| 190 |
|
| 191 |
assert_eq!(item.external_id, "rss:1"); |
| 192 |
assert_eq!(item.feed_id, feed.id); |
| 193 |
assert_eq!(item.busser_id, "rss"); |
| 194 |
assert_eq!(item.bite_author, "author"); |
| 195 |
assert!(!item.is_read); |
| 196 |
assert!(!item.is_starred); |
| 197 |
|
| 198 |
let fetched = ItemsRepository::new(pool.clone()) |
| 199 |
.get(item.id) |
| 200 |
.await |
| 201 |
.unwrap() |
| 202 |
.expect("item should exist"); |
| 203 |
assert_eq!(fetched.id, item.id); |
| 204 |
} |
| 205 |
|
| 206 |
#[tokio::test] |
| 207 |
async fn items_upsert_conflict_updates() { |
| 208 |
let pool = test_db().await; |
| 209 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 210 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 211 |
|
| 212 |
let first = make_item(&pool, &feed, "rss:dup").await; |
| 213 |
|
| 214 |
let second = items_repo |
| 215 |
.upsert(CreateFeedItem { |
| 216 |
external_id: "rss:dup".to_string(), |
| 217 |
feed_id: feed.id, |
| 218 |
busser_id: feed.busser_id.clone(), |
| 219 |
bite_author: "updated_author".to_string(), |
| 220 |
bite_text: "Updated text".to_string(), |
| 221 |
bite_secondary: None, |
| 222 |
bite_indicator: None, |
| 223 |
title: Some("Updated Title".to_string()), |
| 224 |
body: None, |
| 225 |
url: None, |
| 226 |
media: vec![], |
| 227 |
published_at: Utc::now(), |
| 228 |
source_name: "test".to_string(), |
| 229 |
score: None, |
| 230 |
tags: vec![], |
| 231 |
actions: vec![], |
| 232 |
}) |
| 233 |
.await |
| 234 |
.unwrap(); |
| 235 |
|
| 236 |
assert_eq!(first.id, second.id); |
| 237 |
assert_eq!(second.bite_author, "updated_author"); |
| 238 |
assert_eq!(second.bite_text, "Updated text"); |
| 239 |
assert_eq!(items_repo.count_all().await.unwrap(), 1); |
| 240 |
} |
| 241 |
|
| 242 |
#[tokio::test] |
| 243 |
async fn items_get_by_external_id() { |
| 244 |
let pool = test_db().await; |
| 245 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 246 |
let item = make_item(&pool, &feed, "rss:ext1").await; |
| 247 |
|
| 248 |
let found = ItemsRepository::new(pool.clone()) |
| 249 |
.get_by_external_id("rss:ext1") |
| 250 |
.await |
| 251 |
.unwrap() |
| 252 |
.expect("should find by external_id"); |
| 253 |
assert_eq!(found.id, item.id); |
| 254 |
} |
| 255 |
|
| 256 |
#[tokio::test] |
| 257 |
async fn items_list_all_pagination() { |
| 258 |
let pool = test_db().await; |
| 259 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 260 |
let now = Utc::now(); |
| 261 |
|
| 262 |
make_item_at(&pool, &feed, "p:1", now - Duration::hours(3)).await; |
| 263 |
make_item_at(&pool, &feed, "p:2", now - Duration::hours(2)).await; |
| 264 |
make_item_at(&pool, &feed, "p:3", now - Duration::hours(1)).await; |
| 265 |
|
| 266 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 267 |
|
| 268 |
let page1 = items_repo.list_all(2, 0).await.unwrap(); |
| 269 |
assert_eq!(page1.len(), 2); |
| 270 |
assert_eq!(page1[0].external_id, "p:3"); |
| 271 |
assert_eq!(page1[1].external_id, "p:2"); |
| 272 |
|
| 273 |
let page2 = items_repo.list_all(2, 2).await.unwrap(); |
| 274 |
assert_eq!(page2.len(), 1); |
| 275 |
assert_eq!(page2[0].external_id, "p:1"); |
| 276 |
} |
| 277 |
|
| 278 |
#[tokio::test] |
| 279 |
async fn items_list_by_feed_filters() { |
| 280 |
let pool = test_db().await; |
| 281 |
let feed_a = make_feed(&pool, "rss", "Feed A").await; |
| 282 |
let feed_b = make_feed(&pool, "rss", "Feed B").await; |
| 283 |
|
| 284 |
make_item(&pool, &feed_a, "a:1").await; |
| 285 |
make_item(&pool, &feed_a, "a:2").await; |
| 286 |
make_item(&pool, &feed_b, "b:1").await; |
| 287 |
|
| 288 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 289 |
|
| 290 |
let a_items = items_repo.list_by_feed(feed_a.id, 100, 0).await.unwrap(); |
| 291 |
assert_eq!(a_items.len(), 2); |
| 292 |
for i in &a_items { |
| 293 |
assert_eq!(i.feed_id, feed_a.id); |
| 294 |
} |
| 295 |
|
| 296 |
let b_items = items_repo.list_by_feed(feed_b.id, 100, 0).await.unwrap(); |
| 297 |
assert_eq!(b_items.len(), 1); |
| 298 |
assert_eq!(b_items[0].feed_id, feed_b.id); |
| 299 |
} |
| 300 |
|
| 301 |
#[tokio::test] |
| 302 |
async fn items_list_by_busser_filters() { |
| 303 |
let pool = test_db().await; |
| 304 |
let feed_rss = make_feed(&pool, "rss", "RSS Feed").await; |
| 305 |
let feed_hn = make_feed(&pool, "hn", "HN Feed").await; |
| 306 |
|
| 307 |
make_item(&pool, &feed_rss, "rss:1").await; |
| 308 |
make_item(&pool, &feed_hn, "hn:1").await; |
| 309 |
make_item(&pool, &feed_hn, "hn:2").await; |
| 310 |
|
| 311 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 312 |
|
| 313 |
let rss = items_repo.list_by_busser("rss", 100, 0).await.unwrap(); |
| 314 |
assert_eq!(rss.len(), 1); |
| 315 |
|
| 316 |
let hn = items_repo.list_by_busser("hn", 100, 0).await.unwrap(); |
| 317 |
assert_eq!(hn.len(), 2); |
| 318 |
} |
| 319 |
|
| 320 |
#[tokio::test] |
| 321 |
async fn items_list_unread_excludes_read() { |
| 322 |
let pool = test_db().await; |
| 323 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 324 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 325 |
|
| 326 |
let item1 = make_item(&pool, &feed, "u:1").await; |
| 327 |
make_item(&pool, &feed, "u:2").await; |
| 328 |
|
| 329 |
items_repo.mark_read(item1.id, true).await.unwrap(); |
| 330 |
|
| 331 |
let unread = items_repo.list_unread(100, 0).await.unwrap(); |
| 332 |
assert_eq!(unread.len(), 1); |
| 333 |
assert_eq!(unread[0].external_id, "u:2"); |
| 334 |
} |
| 335 |
|
| 336 |
#[tokio::test] |
| 337 |
async fn items_list_starred_only() { |
| 338 |
let pool = test_db().await; |
| 339 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 340 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 341 |
|
| 342 |
make_item(&pool, &feed, "s:1").await; |
| 343 |
let item2 = make_item(&pool, &feed, "s:2").await; |
| 344 |
|
| 345 |
items_repo.mark_starred(item2.id, true).await.unwrap(); |
| 346 |
|
| 347 |
let starred = items_repo.list_starred(100, 0).await.unwrap(); |
| 348 |
assert_eq!(starred.len(), 1); |
| 349 |
assert_eq!(starred[0].external_id, "s:2"); |
| 350 |
} |
| 351 |
|
| 352 |
#[tokio::test] |
| 353 |
async fn items_mark_read_and_unread() { |
| 354 |
let pool = test_db().await; |
| 355 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 356 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 357 |
let item = make_item(&pool, &feed, "r:1").await; |
| 358 |
|
| 359 |
assert!(!item.is_read); |
| 360 |
|
| 361 |
items_repo.mark_read(item.id, true).await.unwrap(); |
| 362 |
let updated = items_repo.get(item.id).await.unwrap().unwrap(); |
| 363 |
assert!(updated.is_read); |
| 364 |
|
| 365 |
items_repo.mark_read(item.id, false).await.unwrap(); |
| 366 |
let updated = items_repo.get(item.id).await.unwrap().unwrap(); |
| 367 |
assert!(!updated.is_read); |
| 368 |
} |
| 369 |
|
| 370 |
#[tokio::test] |
| 371 |
async fn items_mark_starred_and_unstarred() { |
| 372 |
let pool = test_db().await; |
| 373 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 374 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 375 |
let item = make_item(&pool, &feed, "st:1").await; |
| 376 |
|
| 377 |
assert!(!item.is_starred); |
| 378 |
|
| 379 |
items_repo.mark_starred(item.id, true).await.unwrap(); |
| 380 |
let updated = items_repo.get(item.id).await.unwrap().unwrap(); |
| 381 |
assert!(updated.is_starred); |
| 382 |
|
| 383 |
items_repo.mark_starred(item.id, false).await.unwrap(); |
| 384 |
let updated = items_repo.get(item.id).await.unwrap().unwrap(); |
| 385 |
assert!(!updated.is_starred); |
| 386 |
} |
| 387 |
|
| 388 |
#[tokio::test] |
| 389 |
async fn items_count_all_and_unread() { |
| 390 |
let pool = test_db().await; |
| 391 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 392 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 393 |
|
| 394 |
let item1 = make_item(&pool, &feed, "c:1").await; |
| 395 |
make_item(&pool, &feed, "c:2").await; |
| 396 |
make_item(&pool, &feed, "c:3").await; |
| 397 |
|
| 398 |
assert_eq!(items_repo.count_all().await.unwrap(), 3); |
| 399 |
assert_eq!(items_repo.count_unread().await.unwrap(), 3); |
| 400 |
|
| 401 |
items_repo.mark_read(item1.id, true).await.unwrap(); |
| 402 |
assert_eq!(items_repo.count_all().await.unwrap(), 3); |
| 403 |
assert_eq!(items_repo.count_unread().await.unwrap(), 2); |
| 404 |
} |
| 405 |
|
| 406 |
#[tokio::test] |
| 407 |
async fn items_delete_by_feed_removes_items() { |
| 408 |
let pool = test_db().await; |
| 409 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 410 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 411 |
|
| 412 |
make_item(&pool, &feed, "d:1").await; |
| 413 |
make_item(&pool, &feed, "d:2").await; |
| 414 |
assert_eq!(items_repo.count_all().await.unwrap(), 2); |
| 415 |
|
| 416 |
let removed = items_repo.delete_by_feed(feed.id).await.unwrap(); |
| 417 |
assert_eq!(removed, 2); |
| 418 |
assert_eq!(items_repo.count_all().await.unwrap(), 0); |
| 419 |
} |
| 420 |
|
| 421 |
#[tokio::test] |
| 422 |
async fn items_delete_stale_read() { |
| 423 |
let pool = test_db().await; |
| 424 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 425 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 426 |
let now = Utc::now(); |
| 427 |
|
| 428 |
|
| 429 |
let old_read = make_item_at(&pool, &feed, "stale:1", now - Duration::days(60)).await; |
| 430 |
items_repo.mark_read(old_read.id, true).await.unwrap(); |
| 431 |
|
| 432 |
|
| 433 |
let old_starred = make_item_at(&pool, &feed, "stale:2", now - Duration::days(60)).await; |
| 434 |
items_repo.mark_read(old_starred.id, true).await.unwrap(); |
| 435 |
items_repo.mark_starred(old_starred.id, true).await.unwrap(); |
| 436 |
|
| 437 |
|
| 438 |
make_item_at(&pool, &feed, "stale:3", now - Duration::days(60)).await; |
| 439 |
|
| 440 |
|
| 441 |
let recent_read = make_item_at(&pool, &feed, "stale:4", now - Duration::days(5)).await; |
| 442 |
items_repo.mark_read(recent_read.id, true).await.unwrap(); |
| 443 |
|
| 444 |
|
| 445 |
make_item_at(&pool, &feed, "stale:5", now - Duration::days(5)).await; |
| 446 |
|
| 447 |
assert_eq!(items_repo.count_all().await.unwrap(), 5); |
| 448 |
|
| 449 |
let cutoff = now - Duration::days(30); |
| 450 |
let deleted = items_repo.delete_stale_read(cutoff).await.unwrap(); |
| 451 |
assert_eq!(deleted, 1); |
| 452 |
|
| 453 |
assert_eq!(items_repo.count_all().await.unwrap(), 4); |
| 454 |
|
| 455 |
|
| 456 |
assert!(items_repo.get(old_starred.id).await.unwrap().is_some()); |
| 457 |
assert!(items_repo.get(recent_read.id).await.unwrap().is_some()); |
| 458 |
assert!(items_repo.get(old_read.id).await.unwrap().is_none()); |
| 459 |
} |
| 460 |
|
| 461 |
|
| 462 |
|
| 463 |
#[tokio::test] |
| 464 |
async fn health_success_resets_counter() { |
| 465 |
let pool = test_db().await; |
| 466 |
let feeds_repo = FeedsRepository::new(pool.clone()); |
| 467 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 468 |
|
| 469 |
|
| 470 |
feeds_repo |
| 471 |
.record_fetch_failure(feed.id, "timeout") |
| 472 |
.await |
| 473 |
.unwrap(); |
| 474 |
feeds_repo |
| 475 |
.record_fetch_failure(feed.id, "timeout") |
| 476 |
.await |
| 477 |
.unwrap(); |
| 478 |
|
| 479 |
let f = feeds_repo.get(feed.id).await.unwrap().unwrap(); |
| 480 |
assert_eq!(f.consecutive_failures, 2); |
| 481 |
assert_eq!(f.last_error.as_deref(), Some("timeout")); |
| 482 |
|
| 483 |
|
| 484 |
feeds_repo.record_fetch_success(feed.id).await.unwrap(); |
| 485 |
let f = feeds_repo.get(feed.id).await.unwrap().unwrap(); |
| 486 |
assert_eq!(f.consecutive_failures, 0); |
| 487 |
assert!(f.last_error.is_none()); |
| 488 |
assert!(f.last_success_at.is_some()); |
| 489 |
} |
| 490 |
|
| 491 |
#[tokio::test] |
| 492 |
async fn health_failure_increments() { |
| 493 |
let pool = test_db().await; |
| 494 |
let feeds_repo = FeedsRepository::new(pool.clone()); |
| 495 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 496 |
|
| 497 |
feeds_repo |
| 498 |
.record_fetch_failure(feed.id, "dns error") |
| 499 |
.await |
| 500 |
.unwrap(); |
| 501 |
let f = feeds_repo.get(feed.id).await.unwrap().unwrap(); |
| 502 |
assert_eq!(f.consecutive_failures, 1); |
| 503 |
assert_eq!(f.last_error.as_deref(), Some("dns error")); |
| 504 |
|
| 505 |
feeds_repo |
| 506 |
.record_fetch_failure(feed.id, "connection refused") |
| 507 |
.await |
| 508 |
.unwrap(); |
| 509 |
let f = feeds_repo.get(feed.id).await.unwrap().unwrap(); |
| 510 |
assert_eq!(f.consecutive_failures, 2); |
| 511 |
assert_eq!(f.last_error.as_deref(), Some("connection refused")); |
| 512 |
} |
| 513 |
|
| 514 |
#[tokio::test] |
| 515 |
async fn health_defaults_on_new_feed() { |
| 516 |
let pool = test_db().await; |
| 517 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 518 |
assert_eq!(feed.consecutive_failures, 0); |
| 519 |
assert!(feed.last_error.is_none()); |
| 520 |
assert!(feed.last_success_at.is_none()); |
| 521 |
} |
| 522 |
|
| 523 |
|
| 524 |
|
| 525 |
#[tokio::test] |
| 526 |
async fn fts5_search_matches_title() { |
| 527 |
let pool = test_db().await; |
| 528 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 529 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 530 |
|
| 531 |
items_repo |
| 532 |
.upsert(CreateFeedItem { |
| 533 |
external_id: "fts:1".to_string(), |
| 534 |
feed_id: feed.id, |
| 535 |
busser_id: BusserId::new("rss"), |
| 536 |
bite_author: "author".to_string(), |
| 537 |
bite_text: "Some bite text".to_string(), |
| 538 |
bite_secondary: None, |
| 539 |
bite_indicator: None, |
| 540 |
title: Some("Rust programming language".to_string()), |
| 541 |
body: Some("Body content here".to_string()), |
| 542 |
url: None, |
| 543 |
media: vec![], |
| 544 |
published_at: Utc::now(), |
| 545 |
source_name: "test".to_string(), |
| 546 |
score: None, |
| 547 |
tags: vec![], |
| 548 |
actions: vec![], |
| 549 |
}) |
| 550 |
.await |
| 551 |
.unwrap(); |
| 552 |
|
| 553 |
let results = items_repo |
| 554 |
.list_search("Rust", None, false, false, 10, 0) |
| 555 |
.await |
| 556 |
.unwrap(); |
| 557 |
assert_eq!(results.len(), 1); |
| 558 |
assert_eq!(results[0].external_id, "fts:1"); |
| 559 |
} |
| 560 |
|
| 561 |
#[tokio::test] |
| 562 |
async fn fts5_search_matches_body() { |
| 563 |
let pool = test_db().await; |
| 564 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 565 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 566 |
|
| 567 |
items_repo |
| 568 |
.upsert(CreateFeedItem { |
| 569 |
external_id: "fts:2".to_string(), |
| 570 |
feed_id: feed.id, |
| 571 |
busser_id: BusserId::new("rss"), |
| 572 |
bite_author: "author".to_string(), |
| 573 |
bite_text: "Bite".to_string(), |
| 574 |
bite_secondary: None, |
| 575 |
bite_indicator: None, |
| 576 |
title: Some("Title".to_string()), |
| 577 |
body: Some("SQLite full text search is powerful".to_string()), |
| 578 |
url: None, |
| 579 |
media: vec![], |
| 580 |
published_at: Utc::now(), |
| 581 |
source_name: "test".to_string(), |
| 582 |
score: None, |
| 583 |
tags: vec![], |
| 584 |
actions: vec![], |
| 585 |
}) |
| 586 |
.await |
| 587 |
.unwrap(); |
| 588 |
|
| 589 |
let results = items_repo |
| 590 |
.list_search("powerful", None, false, false, 10, 0) |
| 591 |
.await |
| 592 |
.unwrap(); |
| 593 |
assert_eq!(results.len(), 1); |
| 594 |
} |
| 595 |
|
| 596 |
#[tokio::test] |
| 597 |
async fn fts5_search_matches_bite_text() { |
| 598 |
let pool = test_db().await; |
| 599 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 600 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 601 |
|
| 602 |
items_repo |
| 603 |
.upsert(CreateFeedItem { |
| 604 |
external_id: "fts:3".to_string(), |
| 605 |
feed_id: feed.id, |
| 606 |
busser_id: BusserId::new("rss"), |
| 607 |
bite_author: "author".to_string(), |
| 608 |
bite_text: "Breaking news about databases".to_string(), |
| 609 |
bite_secondary: None, |
| 610 |
bite_indicator: None, |
| 611 |
title: None, |
| 612 |
body: None, |
| 613 |
url: None, |
| 614 |
media: vec![], |
| 615 |
published_at: Utc::now(), |
| 616 |
source_name: "test".to_string(), |
| 617 |
score: None, |
| 618 |
tags: vec![], |
| 619 |
actions: vec![], |
| 620 |
}) |
| 621 |
.await |
| 622 |
.unwrap(); |
| 623 |
|
| 624 |
let results = items_repo |
| 625 |
.list_search("databases", None, false, false, 10, 0) |
| 626 |
.await |
| 627 |
.unwrap(); |
| 628 |
assert_eq!(results.len(), 1); |
| 629 |
} |
| 630 |
|
| 631 |
#[tokio::test] |
| 632 |
async fn fts5_multi_word_query() { |
| 633 |
let pool = test_db().await; |
| 634 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 635 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 636 |
|
| 637 |
items_repo |
| 638 |
.upsert(CreateFeedItem { |
| 639 |
external_id: "fts:4".to_string(), |
| 640 |
feed_id: feed.id, |
| 641 |
busser_id: BusserId::new("rss"), |
| 642 |
bite_author: "author".to_string(), |
| 643 |
bite_text: "bite".to_string(), |
| 644 |
bite_secondary: None, |
| 645 |
bite_indicator: None, |
| 646 |
title: Some("Rust async programming guide".to_string()), |
| 647 |
body: None, |
| 648 |
url: None, |
| 649 |
media: vec![], |
| 650 |
published_at: Utc::now(), |
| 651 |
source_name: "test".to_string(), |
| 652 |
score: None, |
| 653 |
tags: vec![], |
| 654 |
actions: vec![], |
| 655 |
}) |
| 656 |
.await |
| 657 |
.unwrap(); |
| 658 |
|
| 659 |
let results = items_repo |
| 660 |
.list_search("Rust programming", None, false, false, 10, 0) |
| 661 |
.await |
| 662 |
.unwrap(); |
| 663 |
assert_eq!(results.len(), 1); |
| 664 |
} |
| 665 |
|
| 666 |
#[tokio::test] |
| 667 |
async fn fts5_special_characters_handled() { |
| 668 |
let pool = test_db().await; |
| 669 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 670 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 671 |
|
| 672 |
make_item(&pool, &feed, "fts:5").await; |
| 673 |
|
| 674 |
|
| 675 |
let results = items_repo |
| 676 |
.list_search("AND OR NOT", None, false, false, 10, 0) |
| 677 |
.await |
| 678 |
.unwrap(); |
| 679 |
assert_eq!(results.len(), 0); |
| 680 |
} |
| 681 |
|
| 682 |
|
| 683 |
|
| 684 |
#[test] |
| 685 |
fn fts5_sanitize_basic_words() { |
| 686 |
assert_eq!(sanitize_fts_query("hello world"), r#""hello" "world""#); |
| 687 |
} |
| 688 |
|
| 689 |
#[test] |
| 690 |
fn fts5_sanitize_operators_quoted() { |
| 691 |
|
| 692 |
assert_eq!(sanitize_fts_query("AND OR NOT"), r#""AND" "OR" "NOT""#,); |
| 693 |
} |
| 694 |
|
| 695 |
#[test] |
| 696 |
fn fts5_sanitize_near_operator() { |
| 697 |
assert_eq!(sanitize_fts_query("NEAR"), r#""NEAR""#); |
| 698 |
assert_eq!(sanitize_fts_query("NEAR/3"), r#""NEAR/3""#); |
| 699 |
assert_eq!( |
| 700 |
sanitize_fts_query("word NEAR/5 other"), |
| 701 |
r#""word" "NEAR/5" "other""# |
| 702 |
); |
| 703 |
} |
| 704 |
|
| 705 |
#[test] |
| 706 |
fn fts5_sanitize_column_prefix() { |
| 707 |
|
| 708 |
assert_eq!(sanitize_fts_query("title:rust"), r#""title:rust""#); |
| 709 |
assert_eq!( |
| 710 |
sanitize_fts_query("body:hello title:world"), |
| 711 |
r#""body:hello" "title:world""#, |
| 712 |
); |
| 713 |
} |
| 714 |
|
| 715 |
#[test] |
| 716 |
fn fts5_sanitize_caret_stripped() { |
| 717 |
|
| 718 |
assert_eq!(sanitize_fts_query("^hello"), r#""hello""#); |
| 719 |
assert_eq!(sanitize_fts_query("^hello world"), r#""hello" "world""#); |
| 720 |
|
| 721 |
assert_eq!(sanitize_fts_query("^^hello"), r#""hello""#); |
| 722 |
} |
| 723 |
|
| 724 |
#[test] |
| 725 |
fn fts5_sanitize_star_stripped() { |
| 726 |
|
| 727 |
assert_eq!(sanitize_fts_query("hello*"), r#""hello""#); |
| 728 |
assert_eq!(sanitize_fts_query("hel*"), r#""hel""#); |
| 729 |
|
| 730 |
assert_eq!(sanitize_fts_query("hello**"), r#""hello""#); |
| 731 |
} |
| 732 |
|
| 733 |
#[test] |
| 734 |
fn fts5_sanitize_caret_and_star_combined() { |
| 735 |
assert_eq!(sanitize_fts_query("^hello*"), r#""hello""#); |
| 736 |
assert_eq!(sanitize_fts_query("^*"), ""); |
| 737 |
} |
| 738 |
|
| 739 |
#[test] |
| 740 |
fn fts5_sanitize_bare_special_chars_dropped() { |
| 741 |
|
| 742 |
assert_eq!(sanitize_fts_query("^"), ""); |
| 743 |
assert_eq!(sanitize_fts_query("*"), ""); |
| 744 |
assert_eq!(sanitize_fts_query("^ word"), r#""word""#); |
| 745 |
assert_eq!(sanitize_fts_query("* word"), r#""word""#); |
| 746 |
} |
| 747 |
|
| 748 |
#[test] |
| 749 |
fn fts5_sanitize_embedded_quotes() { |
| 750 |
|
| 751 |
|
| 752 |
|
| 753 |
assert_eq!(sanitize_fts_query("say \"hi\""), "\"say\" \"\"\"hi\"\"\"",); |
| 754 |
} |
| 755 |
|
| 756 |
#[test] |
| 757 |
fn fts5_sanitize_empty_and_whitespace() { |
| 758 |
assert_eq!(sanitize_fts_query(""), ""); |
| 759 |
assert_eq!(sanitize_fts_query(" "), ""); |
| 760 |
} |
| 761 |
|
| 762 |
#[test] |
| 763 |
fn fts5_sanitize_mixed_special_syntax() { |
| 764 |
|
| 765 |
assert_eq!( |
| 766 |
sanitize_fts_query("^title:rust* NEAR/3 OR body:hello*"), |
| 767 |
r#""title:rust" "NEAR/3" "OR" "body:hello""#, |
| 768 |
); |
| 769 |
} |
| 770 |
|
| 771 |
|
| 772 |
|
| 773 |
#[tokio::test] |
| 774 |
async fn fts5_near_operator_does_not_crash() { |
| 775 |
let pool = test_db().await; |
| 776 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 777 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 778 |
make_item(&pool, &feed, "fts:near1").await; |
| 779 |
|
| 780 |
|
| 781 |
let results = items_repo |
| 782 |
.list_search("NEAR", None, false, false, 10, 0) |
| 783 |
.await |
| 784 |
.unwrap(); |
| 785 |
assert_eq!(results.len(), 0); |
| 786 |
|
| 787 |
let results = items_repo |
| 788 |
.list_search("word NEAR/3 other", None, false, false, 10, 0) |
| 789 |
.await |
| 790 |
.unwrap(); |
| 791 |
assert_eq!(results.len(), 0); |
| 792 |
} |
| 793 |
|
| 794 |
#[tokio::test] |
| 795 |
async fn fts5_column_prefix_does_not_target_column() { |
| 796 |
let pool = test_db().await; |
| 797 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 798 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 799 |
|
| 800 |
items_repo |
| 801 |
.upsert(CreateFeedItem { |
| 802 |
external_id: "fts:col1".to_string(), |
| 803 |
feed_id: feed.id, |
| 804 |
busser_id: BusserId::new("rss"), |
| 805 |
bite_author: "author".to_string(), |
| 806 |
bite_text: "bite".to_string(), |
| 807 |
bite_secondary: None, |
| 808 |
bite_indicator: None, |
| 809 |
title: Some("rust programming".to_string()), |
| 810 |
body: Some("unrelated body".to_string()), |
| 811 |
url: None, |
| 812 |
media: vec![], |
| 813 |
published_at: Utc::now(), |
| 814 |
source_name: "test".to_string(), |
| 815 |
score: None, |
| 816 |
tags: vec![], |
| 817 |
actions: vec![], |
| 818 |
}) |
| 819 |
.await |
| 820 |
.unwrap(); |
| 821 |
|
| 822 |
|
| 823 |
|
| 824 |
let results = items_repo |
| 825 |
.list_search("title:rust", None, false, false, 10, 0) |
| 826 |
.await |
| 827 |
.unwrap(); |
| 828 |
assert_eq!(results.len(), 0); |
| 829 |
} |
| 830 |
|
| 831 |
#[tokio::test] |
| 832 |
async fn fts5_caret_prefix_does_not_crash() { |
| 833 |
let pool = test_db().await; |
| 834 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 835 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 836 |
make_item(&pool, &feed, "fts:caret1").await; |
| 837 |
|
| 838 |
let results = items_repo |
| 839 |
.list_search("^Title", None, false, false, 10, 0) |
| 840 |
.await |
| 841 |
.unwrap(); |
| 842 |
|
| 843 |
assert!(!results.is_empty()); |
| 844 |
} |
| 845 |
|
| 846 |
#[tokio::test] |
| 847 |
async fn fts5_star_suffix_does_not_crash() { |
| 848 |
let pool = test_db().await; |
| 849 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 850 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 851 |
make_item(&pool, &feed, "fts:star1").await; |
| 852 |
|
| 853 |
let results = items_repo |
| 854 |
.list_search("Titl*", None, false, false, 10, 0) |
| 855 |
.await |
| 856 |
.unwrap(); |
| 857 |
|
| 858 |
|
| 859 |
assert_eq!(results.len(), 0); |
| 860 |
} |
| 861 |
|
| 862 |
#[tokio::test] |
| 863 |
async fn fts5_bare_star_and_caret_safe() { |
| 864 |
let pool = test_db().await; |
| 865 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 866 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 867 |
make_item(&pool, &feed, "fts:bare1").await; |
| 868 |
|
| 869 |
|
| 870 |
let results = items_repo |
| 871 |
.list_search("*", None, false, false, 10, 0) |
| 872 |
.await |
| 873 |
.unwrap(); |
| 874 |
assert_eq!(results.len(), 0); |
| 875 |
|
| 876 |
let results = items_repo |
| 877 |
.list_search("^", None, false, false, 10, 0) |
| 878 |
.await |
| 879 |
.unwrap(); |
| 880 |
assert_eq!(results.len(), 0); |
| 881 |
|
| 882 |
let results = items_repo |
| 883 |
.list_search("^*", None, false, false, 10, 0) |
| 884 |
.await |
| 885 |
.unwrap(); |
| 886 |
assert_eq!(results.len(), 0); |
| 887 |
} |
| 888 |
|
| 889 |
#[tokio::test] |
| 890 |
async fn fts5_pagination_works() { |
| 891 |
let pool = test_db().await; |
| 892 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 893 |
let items_repo = ItemsRepository::new(pool.clone()); |
| 894 |
|
| 895 |
for i in 0..5 { |
| 896 |
items_repo |
| 897 |
.upsert(CreateFeedItem { |
| 898 |
external_id: format!("fts:page:{i}"), |
| 899 |
feed_id: feed.id, |
| 900 |
busser_id: BusserId::new("rss"), |
| 901 |
bite_author: "author".to_string(), |
| 902 |
bite_text: "bite".to_string(), |
| 903 |
bite_secondary: None, |
| 904 |
bite_indicator: None, |
| 905 |
title: Some("Searchable title here".to_string()), |
| 906 |
body: None, |
| 907 |
url: None, |
| 908 |
media: vec![], |
| 909 |
published_at: Utc::now() - Duration::hours(i), |
| 910 |
source_name: "test".to_string(), |
| 911 |
score: None, |
| 912 |
tags: vec![], |
| 913 |
actions: vec![], |
| 914 |
}) |
| 915 |
.await |
| 916 |
.unwrap(); |
| 917 |
} |
| 918 |
|
| 919 |
let page1 = items_repo |
| 920 |
.list_search("Searchable", None, false, false, 2, 0) |
| 921 |
.await |
| 922 |
.unwrap(); |
| 923 |
assert_eq!(page1.len(), 2); |
| 924 |
|
| 925 |
let page2 = items_repo |
| 926 |
.list_search("Searchable", None, false, false, 2, 2) |
| 927 |
.await |
| 928 |
.unwrap(); |
| 929 |
assert_eq!(page2.len(), 2); |
| 930 |
|
| 931 |
let page3 = items_repo |
| 932 |
.list_search("Searchable", None, false, false, 2, 4) |
| 933 |
.await |
| 934 |
.unwrap(); |
| 935 |
assert_eq!(page3.len(), 1); |
| 936 |
} |
| 937 |
|
| 938 |
|
| 939 |
|
| 940 |
#[tokio::test] |
| 941 |
async fn tags_set_and_get() { |
| 942 |
let pool = test_db().await; |
| 943 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 944 |
let tags_repo = TagsRepository::new(pool.clone()); |
| 945 |
|
| 946 |
tags_repo |
| 947 |
.set_tags(feed.id, &["tech".into(), "rust".into()]) |
| 948 |
.await |
| 949 |
.unwrap(); |
| 950 |
|
| 951 |
let tags = tags_repo.get_tags(feed.id).await.unwrap(); |
| 952 |
assert_eq!(tags, vec!["rust", "tech"]); |
| 953 |
} |
| 954 |
|
| 955 |
#[tokio::test] |
| 956 |
async fn tags_set_idempotent() { |
| 957 |
let pool = test_db().await; |
| 958 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 959 |
let tags_repo = TagsRepository::new(pool.clone()); |
| 960 |
|
| 961 |
tags_repo |
| 962 |
.set_tags(feed.id, &["a".into(), "b".into()]) |
| 963 |
.await |
| 964 |
.unwrap(); |
| 965 |
tags_repo.set_tags(feed.id, &["c".into()]).await.unwrap(); |
| 966 |
|
| 967 |
let tags = tags_repo.get_tags(feed.id).await.unwrap(); |
| 968 |
assert_eq!(tags, vec!["c"]); |
| 969 |
} |
| 970 |
|
| 971 |
#[tokio::test] |
| 972 |
async fn tags_add_and_remove() { |
| 973 |
let pool = test_db().await; |
| 974 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 975 |
let tags_repo = TagsRepository::new(pool.clone()); |
| 976 |
|
| 977 |
tags_repo.add_tag(feed.id, "news").await.unwrap(); |
| 978 |
tags_repo.add_tag(feed.id, "tech").await.unwrap(); |
| 979 |
tags_repo.add_tag(feed.id, "news").await.unwrap(); |
| 980 |
|
| 981 |
let tags = tags_repo.get_tags(feed.id).await.unwrap(); |
| 982 |
assert_eq!(tags, vec!["news", "tech"]); |
| 983 |
|
| 984 |
tags_repo.remove_tag(feed.id, "news").await.unwrap(); |
| 985 |
let tags = tags_repo.get_tags(feed.id).await.unwrap(); |
| 986 |
assert_eq!(tags, vec!["tech"]); |
| 987 |
} |
| 988 |
|
| 989 |
#[tokio::test] |
| 990 |
async fn tags_list_all_tags() { |
| 991 |
let pool = test_db().await; |
| 992 |
let feed_a = make_feed(&pool, "rss", "A").await; |
| 993 |
let feed_b = make_feed(&pool, "hn", "B").await; |
| 994 |
let tags_repo = TagsRepository::new(pool.clone()); |
| 995 |
|
| 996 |
tags_repo |
| 997 |
.set_tags(feed_a.id, &["tech".into(), "news".into()]) |
| 998 |
.await |
| 999 |
.unwrap(); |
| 1000 |
tags_repo |
| 1001 |
.set_tags(feed_b.id, &["tech".into(), "fun".into()]) |
| 1002 |
.await |
| 1003 |
.unwrap(); |
| 1004 |
|
| 1005 |
let all = tags_repo.list_all_tags().await.unwrap(); |
| 1006 |
assert_eq!(all, vec!["fun", "news", "tech"]); |
| 1007 |
} |
| 1008 |
|
| 1009 |
#[tokio::test] |
| 1010 |
async fn tags_cascade_on_feed_delete() { |
| 1011 |
let pool = test_db().await; |
| 1012 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1013 |
let tags_repo = TagsRepository::new(pool.clone()); |
| 1014 |
let feeds_repo = FeedsRepository::new(pool.clone()); |
| 1015 |
|
| 1016 |
tags_repo.set_tags(feed.id, &["x".into()]).await.unwrap(); |
| 1017 |
feeds_repo.delete(feed.id).await.unwrap(); |
| 1018 |
|
| 1019 |
let tags = tags_repo.get_tags(feed.id).await.unwrap(); |
| 1020 |
assert!(tags.is_empty()); |
| 1021 |
} |
| 1022 |
|
| 1023 |
#[tokio::test] |
| 1024 |
async fn tags_feed_ids_with_tags() { |
| 1025 |
let pool = test_db().await; |
| 1026 |
let feed_a = make_feed(&pool, "rss", "A").await; |
| 1027 |
let feed_b = make_feed(&pool, "hn", "B").await; |
| 1028 |
let _feed_c = make_feed(&pool, "reddit", "C").await; |
| 1029 |
let tags_repo = TagsRepository::new(pool.clone()); |
| 1030 |
|
| 1031 |
tags_repo |
| 1032 |
.set_tags(feed_a.id, &["tech".into()]) |
| 1033 |
.await |
| 1034 |
.unwrap(); |
| 1035 |
tags_repo |
| 1036 |
.set_tags(feed_b.id, &["tech".into(), "news".into()]) |
| 1037 |
.await |
| 1038 |
.unwrap(); |
| 1039 |
|
| 1040 |
let ids = tags_repo |
| 1041 |
.feed_ids_with_tags(&["tech".into()]) |
| 1042 |
.await |
| 1043 |
.unwrap(); |
| 1044 |
assert_eq!(ids.len(), 2); |
| 1045 |
assert!(ids.contains(&feed_a.id)); |
| 1046 |
assert!(ids.contains(&feed_b.id)); |
| 1047 |
} |
| 1048 |
|
| 1049 |
|
| 1050 |
|
| 1051 |
#[tokio::test] |
| 1052 |
async fn state_set_and_get() { |
| 1053 |
let pool = test_db().await; |
| 1054 |
let state = StateRepository::new(pool.clone()); |
| 1055 |
|
| 1056 |
state.set("rss", "cursor", "abc123").await.unwrap(); |
| 1057 |
let val = state.get("rss", "cursor").await.unwrap(); |
| 1058 |
assert_eq!(val, Some("abc123".to_string())); |
| 1059 |
} |
| 1060 |
|
| 1061 |
#[tokio::test] |
| 1062 |
async fn state_set_overwrites_value() { |
| 1063 |
let pool = test_db().await; |
| 1064 |
let state = StateRepository::new(pool.clone()); |
| 1065 |
|
| 1066 |
state.set("rss", "cursor", "first").await.unwrap(); |
| 1067 |
state.set("rss", "cursor", "second").await.unwrap(); |
| 1068 |
|
| 1069 |
let val = state.get("rss", "cursor").await.unwrap(); |
| 1070 |
assert_eq!(val, Some("second".to_string())); |
| 1071 |
} |
| 1072 |
|
| 1073 |
#[tokio::test] |
| 1074 |
async fn state_get_missing_returns_none() { |
| 1075 |
let pool = test_db().await; |
| 1076 |
let state = StateRepository::new(pool.clone()); |
| 1077 |
|
| 1078 |
let val = state.get("rss", "nonexistent").await.unwrap(); |
| 1079 |
assert!(val.is_none()); |
| 1080 |
} |
| 1081 |
|
| 1082 |
#[tokio::test] |
| 1083 |
async fn state_delete_removes_key() { |
| 1084 |
let pool = test_db().await; |
| 1085 |
let state = StateRepository::new(pool.clone()); |
| 1086 |
|
| 1087 |
state.set("rss", "token", "secret").await.unwrap(); |
| 1088 |
state.delete("rss", "token").await.unwrap(); |
| 1089 |
|
| 1090 |
let val = state.get("rss", "token").await.unwrap(); |
| 1091 |
assert!(val.is_none()); |
| 1092 |
} |
| 1093 |
|
| 1094 |
#[tokio::test] |
| 1095 |
async fn state_delete_all_clears_busser() { |
| 1096 |
let pool = test_db().await; |
| 1097 |
let state = StateRepository::new(pool.clone()); |
| 1098 |
|
| 1099 |
state.set("rss", "key1", "val1").await.unwrap(); |
| 1100 |
state.set("rss", "key2", "val2").await.unwrap(); |
| 1101 |
state.set("rss", "key3", "val3").await.unwrap(); |
| 1102 |
|
| 1103 |
let removed = state.delete_all("rss").await.unwrap(); |
| 1104 |
assert_eq!(removed, 3); |
| 1105 |
|
| 1106 |
let remaining = state.list("rss").await.unwrap(); |
| 1107 |
assert!(remaining.is_empty()); |
| 1108 |
} |
| 1109 |
|
| 1110 |
|
| 1111 |
|
| 1112 |
#[tokio::test] |
| 1113 |
async fn config_set_and_get() { |
| 1114 |
let pool = test_db().await; |
| 1115 |
let config = ConfigRepository::new(pool); |
| 1116 |
config.set("theme", "dark").await.unwrap(); |
| 1117 |
let val = config.get("theme").await.unwrap(); |
| 1118 |
assert_eq!(val, Some("dark".to_string())); |
| 1119 |
} |
| 1120 |
|
| 1121 |
#[tokio::test] |
| 1122 |
async fn config_set_overwrites() { |
| 1123 |
let pool = test_db().await; |
| 1124 |
let config = ConfigRepository::new(pool); |
| 1125 |
config.set("lang", "en").await.unwrap(); |
| 1126 |
config.set("lang", "fr").await.unwrap(); |
| 1127 |
let val = config.get("lang").await.unwrap(); |
| 1128 |
assert_eq!(val, Some("fr".to_string())); |
| 1129 |
} |
| 1130 |
|
| 1131 |
#[tokio::test] |
| 1132 |
async fn config_get_missing_returns_none() { |
| 1133 |
let pool = test_db().await; |
| 1134 |
let config = ConfigRepository::new(pool); |
| 1135 |
let val = config.get("nonexistent").await.unwrap(); |
| 1136 |
assert!(val.is_none()); |
| 1137 |
} |
| 1138 |
|
| 1139 |
#[tokio::test] |
| 1140 |
async fn config_delete() { |
| 1141 |
let pool = test_db().await; |
| 1142 |
let config = ConfigRepository::new(pool); |
| 1143 |
config.set("key", "value").await.unwrap(); |
| 1144 |
config.delete("key").await.unwrap(); |
| 1145 |
let val = config.get("key").await.unwrap(); |
| 1146 |
assert!(val.is_none()); |
| 1147 |
} |
| 1148 |
|
| 1149 |
|
| 1150 |
|
| 1151 |
#[tokio::test] |
| 1152 |
async fn feeds_update_name() { |
| 1153 |
let pool = test_db().await; |
| 1154 |
let feed = make_feed(&pool, "rss", "Old Name").await; |
| 1155 |
let feeds = FeedsRepository::new(pool); |
| 1156 |
feeds.update_name(feed.id, "New Name").await.unwrap(); |
| 1157 |
let updated = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1158 |
assert_eq!(updated.name, "New Name"); |
| 1159 |
} |
| 1160 |
|
| 1161 |
#[tokio::test] |
| 1162 |
async fn feeds_update_config() { |
| 1163 |
let pool = test_db().await; |
| 1164 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1165 |
let feeds = FeedsRepository::new(pool); |
| 1166 |
feeds |
| 1167 |
.update_config(feed.id, r#"{"url":"https://example.com/rss"}"#) |
| 1168 |
.await |
| 1169 |
.unwrap(); |
| 1170 |
let updated = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1171 |
assert_eq!(updated.config, r#"{"url":"https://example.com/rss"}"#); |
| 1172 |
} |
| 1173 |
|
| 1174 |
#[tokio::test] |
| 1175 |
async fn feeds_record_fetch_success_resets_failures() { |
| 1176 |
let pool = test_db().await; |
| 1177 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1178 |
let feeds = FeedsRepository::new(pool); |
| 1179 |
|
| 1180 |
feeds |
| 1181 |
.record_fetch_failure(feed.id, "timeout") |
| 1182 |
.await |
| 1183 |
.unwrap(); |
| 1184 |
feeds |
| 1185 |
.record_fetch_failure(feed.id, "dns error") |
| 1186 |
.await |
| 1187 |
.unwrap(); |
| 1188 |
let failed = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1189 |
assert_eq!(failed.consecutive_failures, 2); |
| 1190 |
assert_eq!(failed.last_error.as_deref(), Some("dns error")); |
| 1191 |
|
| 1192 |
feeds.record_fetch_success(feed.id).await.unwrap(); |
| 1193 |
let ok = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1194 |
assert_eq!(ok.consecutive_failures, 0); |
| 1195 |
assert!(ok.last_error.is_none()); |
| 1196 |
assert!(ok.last_success_at.is_some()); |
| 1197 |
} |
| 1198 |
|
| 1199 |
#[tokio::test] |
| 1200 |
async fn feeds_record_fetch_failure_increments() { |
| 1201 |
let pool = test_db().await; |
| 1202 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1203 |
let feeds = FeedsRepository::new(pool); |
| 1204 |
feeds.record_fetch_failure(feed.id, "err1").await.unwrap(); |
| 1205 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1206 |
assert_eq!(f.consecutive_failures, 1); |
| 1207 |
assert_eq!(f.last_error.as_deref(), Some("err1")); |
| 1208 |
feeds.record_fetch_failure(feed.id, "err2").await.unwrap(); |
| 1209 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1210 |
assert_eq!(f.consecutive_failures, 2); |
| 1211 |
assert_eq!(f.last_error.as_deref(), Some("err2")); |
| 1212 |
} |
| 1213 |
|
| 1214 |
|
| 1215 |
|
| 1216 |
#[tokio::test] |
| 1217 |
async fn structured_failure_rate_limited_no_increment() { |
| 1218 |
let pool = test_db().await; |
| 1219 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1220 |
let feeds = FeedsRepository::new(pool); |
| 1221 |
let err = StructuredError::rate_limited("429 Too Many Requests", 120); |
| 1222 |
let tripped = feeds |
| 1223 |
.record_fetch_failure_structured(feed.id, &err) |
| 1224 |
.await |
| 1225 |
.unwrap(); |
| 1226 |
assert!(!tripped); |
| 1227 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1228 |
assert_eq!( |
| 1229 |
f.consecutive_failures, 0, |
| 1230 |
"rate_limited should not increment failures" |
| 1231 |
); |
| 1232 |
assert!(f.last_error.is_some(), "error should be stored"); |
| 1233 |
|
| 1234 |
let stored = StructuredError::from_last_error(f.last_error.as_ref().unwrap()); |
| 1235 |
assert_eq!(stored.category, ErrorCategory::RateLimited); |
| 1236 |
assert_eq!(stored.retry_after_secs, Some(120)); |
| 1237 |
} |
| 1238 |
|
| 1239 |
#[tokio::test] |
| 1240 |
async fn structured_failure_auth_immediate_circuit_break() { |
| 1241 |
let pool = test_db().await; |
| 1242 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1243 |
let feeds = FeedsRepository::new(pool); |
| 1244 |
let err = StructuredError::new(ErrorCategory::Auth, "401 Unauthorized"); |
| 1245 |
let tripped = feeds |
| 1246 |
.record_fetch_failure_structured(feed.id, &err) |
| 1247 |
.await |
| 1248 |
.unwrap(); |
| 1249 |
assert!( |
| 1250 |
tripped, |
| 1251 |
"auth error should immediately trip circuit breaker" |
| 1252 |
); |
| 1253 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1254 |
assert!(f.circuit_broken); |
| 1255 |
assert_eq!(f.consecutive_failures, 1); |
| 1256 |
} |
| 1257 |
|
| 1258 |
#[tokio::test] |
| 1259 |
async fn structured_failure_config_immediate_circuit_break() { |
| 1260 |
let pool = test_db().await; |
| 1261 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1262 |
let feeds = FeedsRepository::new(pool); |
| 1263 |
let err = StructuredError::new(ErrorCategory::Config, "404 Not Found"); |
| 1264 |
let tripped = feeds |
| 1265 |
.record_fetch_failure_structured(feed.id, &err) |
| 1266 |
.await |
| 1267 |
.unwrap(); |
| 1268 |
assert!( |
| 1269 |
tripped, |
| 1270 |
"config error should immediately trip circuit breaker" |
| 1271 |
); |
| 1272 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1273 |
assert!(f.circuit_broken); |
| 1274 |
} |
| 1275 |
|
| 1276 |
#[tokio::test] |
| 1277 |
async fn structured_failure_transient_increments_normally() { |
| 1278 |
let pool = test_db().await; |
| 1279 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1280 |
let feeds = FeedsRepository::new(pool); |
| 1281 |
let err = StructuredError::new(ErrorCategory::Transient, "HTTP 503"); |
| 1282 |
let tripped = feeds |
| 1283 |
.record_fetch_failure_structured(feed.id, &err) |
| 1284 |
.await |
| 1285 |
.unwrap(); |
| 1286 |
assert!(!tripped); |
| 1287 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1288 |
assert_eq!(f.consecutive_failures, 1); |
| 1289 |
assert!(!f.circuit_broken); |
| 1290 |
} |
| 1291 |
|
| 1292 |
#[tokio::test] |
| 1293 |
async fn structured_failure_transient_trips_at_threshold() { |
| 1294 |
let pool = test_db().await; |
| 1295 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1296 |
let feeds = FeedsRepository::new(pool); |
| 1297 |
|
| 1298 |
for i in 0..CIRCUIT_BREAKER_THRESHOLD - 1 { |
| 1299 |
let err = StructuredError::new(ErrorCategory::Transient, format!("error {i}")); |
| 1300 |
let tripped = feeds |
| 1301 |
.record_fetch_failure_structured(feed.id, &err) |
| 1302 |
.await |
| 1303 |
.unwrap(); |
| 1304 |
assert!(!tripped); |
| 1305 |
} |
| 1306 |
|
| 1307 |
let err = StructuredError::new(ErrorCategory::Transient, "final error"); |
| 1308 |
let tripped = feeds |
| 1309 |
.record_fetch_failure_structured(feed.id, &err) |
| 1310 |
.await |
| 1311 |
.unwrap(); |
| 1312 |
assert!(tripped); |
| 1313 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1314 |
assert!(f.circuit_broken); |
| 1315 |
} |
| 1316 |
|
| 1317 |
|
| 1318 |
|
| 1319 |
#[tokio::test] |
| 1320 |
async fn items_count_by_busser() { |
| 1321 |
let pool = test_db().await; |
| 1322 |
let feed_a = make_feed(&pool, "rss", "A").await; |
| 1323 |
let feed_b = make_feed(&pool, "hn", "B").await; |
| 1324 |
make_item(&pool, &feed_a, "a1").await; |
| 1325 |
make_item(&pool, &feed_a, "a2").await; |
| 1326 |
make_item(&pool, &feed_b, "b1").await; |
| 1327 |
let items = ItemsRepository::new(pool); |
| 1328 |
assert_eq!(items.count_by_busser("rss").await.unwrap(), 2); |
| 1329 |
assert_eq!(items.count_by_busser("hn").await.unwrap(), 1); |
| 1330 |
assert_eq!(items.count_by_busser("nonexistent").await.unwrap(), 0); |
| 1331 |
} |
| 1332 |
|
| 1333 |
#[tokio::test] |
| 1334 |
async fn items_count_unread_by_busser() { |
| 1335 |
let pool = test_db().await; |
| 1336 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1337 |
let item1 = make_item(&pool, &feed, "i1").await; |
| 1338 |
make_item(&pool, &feed, "i2").await; |
| 1339 |
let items = ItemsRepository::new(pool); |
| 1340 |
assert_eq!(items.count_unread_by_busser("rss").await.unwrap(), 2); |
| 1341 |
items.mark_read(item1.id, true).await.unwrap(); |
| 1342 |
assert_eq!(items.count_unread_by_busser("rss").await.unwrap(), 1); |
| 1343 |
} |
| 1344 |
|
| 1345 |
#[tokio::test] |
| 1346 |
async fn items_counts_by_busser_bulk() { |
| 1347 |
let pool = test_db().await; |
| 1348 |
let feed_a = make_feed(&pool, "rss", "A").await; |
| 1349 |
let feed_b = make_feed(&pool, "hn", "B").await; |
| 1350 |
let item_a1 = make_item(&pool, &feed_a, "a1").await; |
| 1351 |
make_item(&pool, &feed_a, "a2").await; |
| 1352 |
make_item(&pool, &feed_b, "b1").await; |
| 1353 |
let items = ItemsRepository::new(pool); |
| 1354 |
items.mark_read(item_a1.id, true).await.unwrap(); |
| 1355 |
let counts = items.counts_by_busser().await.unwrap(); |
| 1356 |
|
| 1357 |
let rss = counts.iter().find(|(b, _, _)| b == "rss").unwrap(); |
| 1358 |
assert_eq!(rss.1, 2); |
| 1359 |
assert_eq!(rss.2, 1); |
| 1360 |
let hn = counts.iter().find(|(b, _, _)| b == "hn").unwrap(); |
| 1361 |
assert_eq!(hn.1, 1); |
| 1362 |
assert_eq!(hn.2, 1); |
| 1363 |
} |
| 1364 |
|
| 1365 |
#[tokio::test] |
| 1366 |
async fn items_search_with_source_filter() { |
| 1367 |
let pool = test_db().await; |
| 1368 |
let feed_a = make_feed(&pool, "rss", "RSS Feed").await; |
| 1369 |
let feed_b = make_feed(&pool, "hn", "HN Feed").await; |
| 1370 |
make_item(&pool, &feed_a, "rss1").await; |
| 1371 |
make_item(&pool, &feed_b, "hn1").await; |
| 1372 |
let items = ItemsRepository::new(pool); |
| 1373 |
|
| 1374 |
let all = items |
| 1375 |
.list_search("Title", None, false, false, 10, 0) |
| 1376 |
.await |
| 1377 |
.unwrap(); |
| 1378 |
assert_eq!(all.len(), 2); |
| 1379 |
|
| 1380 |
let rss_only = items |
| 1381 |
.list_search("Title", Some("rss"), false, false, 10, 0) |
| 1382 |
.await |
| 1383 |
.unwrap(); |
| 1384 |
assert_eq!(rss_only.len(), 1); |
| 1385 |
assert_eq!(rss_only[0].busser_id.as_str(), "rss"); |
| 1386 |
} |
| 1387 |
|
| 1388 |
#[tokio::test] |
| 1389 |
async fn items_search_unread_only() { |
| 1390 |
let pool = test_db().await; |
| 1391 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1392 |
let item1 = make_item(&pool, &feed, "i1").await; |
| 1393 |
make_item(&pool, &feed, "i2").await; |
| 1394 |
let items = ItemsRepository::new(pool); |
| 1395 |
items.mark_read(item1.id, true).await.unwrap(); |
| 1396 |
let unread = items |
| 1397 |
.list_search("Title", None, true, false, 10, 0) |
| 1398 |
.await |
| 1399 |
.unwrap(); |
| 1400 |
assert_eq!(unread.len(), 1); |
| 1401 |
} |
| 1402 |
|
| 1403 |
#[tokio::test] |
| 1404 |
async fn items_search_starred_only() { |
| 1405 |
let pool = test_db().await; |
| 1406 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1407 |
let item1 = make_item(&pool, &feed, "i1").await; |
| 1408 |
make_item(&pool, &feed, "i2").await; |
| 1409 |
let items = ItemsRepository::new(pool); |
| 1410 |
items.mark_starred(item1.id, true).await.unwrap(); |
| 1411 |
let starred = items |
| 1412 |
.list_search("Title", None, false, true, 10, 0) |
| 1413 |
.await |
| 1414 |
.unwrap(); |
| 1415 |
assert_eq!(starred.len(), 1); |
| 1416 |
} |
| 1417 |
|
| 1418 |
|
| 1419 |
|
| 1420 |
#[tokio::test] |
| 1421 |
async fn tags_all_feed_tags_bulk() { |
| 1422 |
let pool = test_db().await; |
| 1423 |
let feed_a = make_feed(&pool, "rss", "A").await; |
| 1424 |
let feed_b = make_feed(&pool, "hn", "B").await; |
| 1425 |
let tags = TagsRepository::new(pool); |
| 1426 |
tags.set_tags(feed_a.id, &["tech".into(), "news".into()]) |
| 1427 |
.await |
| 1428 |
.unwrap(); |
| 1429 |
tags.set_tags(feed_b.id, &["fun".into()]).await.unwrap(); |
| 1430 |
let all = tags.all_feed_tags().await.unwrap(); |
| 1431 |
assert_eq!(all.len(), 3); |
| 1432 |
|
| 1433 |
assert!( |
| 1434 |
all.iter() |
| 1435 |
.any(|(id, tag)| *id == feed_a.id && tag == "news") |
| 1436 |
); |
| 1437 |
assert!( |
| 1438 |
all.iter() |
| 1439 |
.any(|(id, tag)| *id == feed_a.id && tag == "tech") |
| 1440 |
); |
| 1441 |
assert!(all.iter().any(|(id, tag)| *id == feed_b.id && tag == "fun")); |
| 1442 |
} |
| 1443 |
|
| 1444 |
|
| 1445 |
|
| 1446 |
#[tokio::test] |
| 1447 |
async fn state_list_returns_ordered() { |
| 1448 |
let pool = test_db().await; |
| 1449 |
let state = StateRepository::new(pool); |
| 1450 |
state.set("rss", "cursor", "abc").await.unwrap(); |
| 1451 |
state.set("rss", "auth_token", "xyz").await.unwrap(); |
| 1452 |
state.set("rss", "page", "2").await.unwrap(); |
| 1453 |
state.set("hn", "unrelated", "val").await.unwrap(); |
| 1454 |
let list = state.list("rss").await.unwrap(); |
| 1455 |
assert_eq!(list.len(), 3); |
| 1456 |
assert_eq!(list[0].key, "auth_token"); |
| 1457 |
assert_eq!(list[1].key, "cursor"); |
| 1458 |
assert_eq!(list[2].key, "page"); |
| 1459 |
} |
| 1460 |
|
| 1461 |
|
| 1462 |
|
| 1463 |
#[tokio::test] |
| 1464 |
async fn circuit_breaker_new_feed_not_broken() { |
| 1465 |
let pool = test_db().await; |
| 1466 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1467 |
assert!(!feed.circuit_broken); |
| 1468 |
} |
| 1469 |
|
| 1470 |
#[tokio::test] |
| 1471 |
async fn circuit_breaker_trips_at_threshold() { |
| 1472 |
let pool = test_db().await; |
| 1473 |
let feeds = FeedsRepository::new(pool.clone()); |
| 1474 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1475 |
|
| 1476 |
|
| 1477 |
for i in 0..(CIRCUIT_BREAKER_THRESHOLD - 1) { |
| 1478 |
let tripped = feeds |
| 1479 |
.record_fetch_failure(feed.id, &format!("error {i}")) |
| 1480 |
.await |
| 1481 |
.unwrap(); |
| 1482 |
assert!(!tripped, "should not trip at failure {}", i + 1); |
| 1483 |
} |
| 1484 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1485 |
assert!(!f.circuit_broken); |
| 1486 |
assert_eq!(f.consecutive_failures, CIRCUIT_BREAKER_THRESHOLD - 1); |
| 1487 |
|
| 1488 |
|
| 1489 |
let tripped = feeds |
| 1490 |
.record_fetch_failure(feed.id, "final error") |
| 1491 |
.await |
| 1492 |
.unwrap(); |
| 1493 |
assert!(tripped, "should trip at threshold"); |
| 1494 |
|
| 1495 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1496 |
assert!(f.circuit_broken); |
| 1497 |
assert_eq!(f.consecutive_failures, CIRCUIT_BREAKER_THRESHOLD); |
| 1498 |
} |
| 1499 |
|
| 1500 |
#[tokio::test] |
| 1501 |
async fn circuit_breaker_does_not_trip_again_once_broken() { |
| 1502 |
let pool = test_db().await; |
| 1503 |
let feeds = FeedsRepository::new(pool.clone()); |
| 1504 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1505 |
|
| 1506 |
|
| 1507 |
for _ in 0..CIRCUIT_BREAKER_THRESHOLD { |
| 1508 |
feeds.record_fetch_failure(feed.id, "err").await.unwrap(); |
| 1509 |
} |
| 1510 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1511 |
assert!(f.circuit_broken); |
| 1512 |
|
| 1513 |
|
| 1514 |
let tripped = feeds |
| 1515 |
.record_fetch_failure(feed.id, "extra error") |
| 1516 |
.await |
| 1517 |
.unwrap(); |
| 1518 |
assert!(!tripped, "should not re-trip"); |
| 1519 |
} |
| 1520 |
|
| 1521 |
#[tokio::test] |
| 1522 |
async fn circuit_breaker_reset_clears_state() { |
| 1523 |
let pool = test_db().await; |
| 1524 |
let feeds = FeedsRepository::new(pool.clone()); |
| 1525 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1526 |
|
| 1527 |
|
| 1528 |
for _ in 0..CIRCUIT_BREAKER_THRESHOLD { |
| 1529 |
feeds.record_fetch_failure(feed.id, "err").await.unwrap(); |
| 1530 |
} |
| 1531 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1532 |
assert!(f.circuit_broken); |
| 1533 |
assert_eq!(f.consecutive_failures, CIRCUIT_BREAKER_THRESHOLD); |
| 1534 |
assert!(f.last_error.is_some()); |
| 1535 |
|
| 1536 |
|
| 1537 |
feeds.reset_circuit_breaker(feed.id).await.unwrap(); |
| 1538 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1539 |
assert!(!f.circuit_broken); |
| 1540 |
assert_eq!(f.consecutive_failures, 0); |
| 1541 |
assert!(f.last_error.is_none()); |
| 1542 |
} |
| 1543 |
|
| 1544 |
#[tokio::test] |
| 1545 |
async fn circuit_breaker_success_resets_counter_but_not_broken_flag() { |
| 1546 |
let pool = test_db().await; |
| 1547 |
let feeds = FeedsRepository::new(pool.clone()); |
| 1548 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1549 |
|
| 1550 |
|
| 1551 |
for _ in 0..5 { |
| 1552 |
feeds.record_fetch_failure(feed.id, "err").await.unwrap(); |
| 1553 |
} |
| 1554 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1555 |
assert_eq!(f.consecutive_failures, 5); |
| 1556 |
assert!(!f.circuit_broken); |
| 1557 |
|
| 1558 |
|
| 1559 |
feeds.record_fetch_success(feed.id).await.unwrap(); |
| 1560 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1561 |
assert_eq!(f.consecutive_failures, 0); |
| 1562 |
assert!(!f.circuit_broken); |
| 1563 |
} |
| 1564 |
|
| 1565 |
#[tokio::test] |
| 1566 |
async fn circuit_broken_feed_excluded_from_list_enabled() { |
| 1567 |
let pool = test_db().await; |
| 1568 |
let feeds = FeedsRepository::new(pool.clone()); |
| 1569 |
let feed_ok = make_feed(&pool, "rss", "OK Feed").await; |
| 1570 |
let feed_broken = make_feed(&pool, "hn", "Broken Feed").await; |
| 1571 |
|
| 1572 |
|
| 1573 |
for _ in 0..CIRCUIT_BREAKER_THRESHOLD { |
| 1574 |
feeds |
| 1575 |
.record_fetch_failure(feed_broken.id, "err") |
| 1576 |
.await |
| 1577 |
.unwrap(); |
| 1578 |
} |
| 1579 |
|
| 1580 |
let enabled = feeds.list_enabled().await.unwrap(); |
| 1581 |
assert_eq!(enabled.len(), 1); |
| 1582 |
assert_eq!(enabled[0].id, feed_ok.id); |
| 1583 |
|
| 1584 |
|
| 1585 |
let all = feeds.list_all().await.unwrap(); |
| 1586 |
assert_eq!(all.len(), 2); |
| 1587 |
} |
| 1588 |
|
| 1589 |
#[tokio::test] |
| 1590 |
async fn circuit_breaker_set_and_clear() { |
| 1591 |
let pool = test_db().await; |
| 1592 |
let feeds = FeedsRepository::new(pool.clone()); |
| 1593 |
let feed = make_feed(&pool, "rss", "Feed").await; |
| 1594 |
assert!(!feed.circuit_broken); |
| 1595 |
|
| 1596 |
feeds.set_circuit_broken(feed.id, true).await.unwrap(); |
| 1597 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1598 |
assert!(f.circuit_broken); |
| 1599 |
|
| 1600 |
feeds.set_circuit_broken(feed.id, false).await.unwrap(); |
| 1601 |
let f = feeds.get(feed.id).await.unwrap().unwrap(); |
| 1602 |
assert!(!f.circuit_broken); |
| 1603 |
} |
| 1604 |
|