use anyhow::{Context, Result};
use sqlx::sqlite::{SqliteConnectOptions, SqlitePool, SqlitePoolOptions};
use sqlx::{ConnectOptions, FromRow, Row};
use std::str::FromStr;
use crate::config::Config;
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
pub enum RedeemError {
#[error("invite code not found")]
NotFound,
#[error("invite code expired")]
Expired,
#[error("invite code already redeemed")]
AlreadyRedeemed,
#[error("beta is at capacity")]
CapacityFull,
}
pub type Pool = SqlitePool;
#[derive(Debug, Clone, FromRow, PartialEq, Eq)]
pub struct Feed {
pub id: i64,
pub url: String,
pub title: Option<String>,
pub site_url: Option<String>,
pub etag: Option<String>,
pub last_modified: Option<String>,
pub last_polled: Option<String>,
pub next_poll: Option<String>,
#[sqlx(default)]
pub consecutive_errors: i64,
}
#[derive(Debug, Clone, FromRow, PartialEq, Eq)]
pub struct Entry {
pub id: i64,
pub feed_id: i64,
pub guid: String,
pub url: Option<String>,
pub title: Option<String>,
pub author: Option<String>,
pub published: Option<String>,
pub content_html: Option<String>,
pub fetched_at: String,
}
#[derive(Debug, Clone, FromRow, PartialEq, Eq)]
pub struct EntryListRow {
pub id: i64,
pub feed_id: i64,
pub guid: String,
pub url: Option<String>,
pub title: Option<String>,
pub published: Option<String>,
pub read: bool,
pub starred: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ListView {
Unread,
Starred,
All,
}
impl ListView {
fn predicate(self) -> &'static str {
match self {
ListView::Unread => "COALESCE(s.read, 0) = 0",
ListView::Starred => "COALESCE(s.starred, 0) = 1",
ListView::All => "1 = 1",
}
}
}
#[derive(Debug, Clone, FromRow, PartialEq, Eq)]
pub struct EntryState {
pub did: String,
pub entry_id: i64,
pub read: bool,
pub starred: bool,
pub updated_at: String,
}
#[derive(Debug, Clone, FromRow, PartialEq, Eq)]
pub struct ReadCursor {
pub did: String,
pub feed_url: String,
pub read_through: Option<String>,
pub read_ids: String,
pub unread_ids: String,
pub dirty: bool,
#[sqlx(default)]
pub pds_created: bool,
pub updated_at: String,
}
pub const ADOPTION_STAT_KEY: &str = "adoption.subscription";
#[derive(Debug, Clone, FromRow, PartialEq, Eq)]
pub struct NetworkStat {
pub key: String,
pub source: String,
pub value: i64,
pub truncated: bool,
pub observed_at: String,
}
#[derive(Debug, Clone, Default)]
pub struct NewFeed {
pub url: String,
pub title: Option<String>,
pub site_url: Option<String>,
pub etag: Option<String>,
pub last_modified: Option<String>,
pub last_polled: Option<String>,
pub next_poll: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct NewEntry {
pub guid: String,
pub url: Option<String>,
pub title: Option<String>,
pub author: Option<String>,
pub published: Option<String>,
pub content_html: Option<String>,
pub fetched_at: Option<String>,
pub keep_stored_content: bool,
}
const SCHEMA: &str = r#"
PRAGMA foreign_keys = ON;
CREATE TABLE IF NOT EXISTS feeds (
id INTEGER PRIMARY KEY AUTOINCREMENT,
url TEXT NOT NULL UNIQUE,
title TEXT,
site_url TEXT,
etag TEXT,
last_modified TEXT,
last_polled TEXT,
next_poll TEXT,
consecutive_errors INTEGER NOT NULL DEFAULT 0,
last_error_kind TEXT,
last_error TEXT,
-- What the poller does with this row; see `feed::FeedKind`. Written by the
-- Rust side at insert so SQL never re-derives it from the URL.
kind TEXT NOT NULL DEFAULT 'rss'
);
CREATE INDEX IF NOT EXISTS idx_feeds_next_poll ON feeds (next_poll);
-- NOTE: `idx_feeds_kind` is created in `apply_migrations`, AFTER `kind` is
-- ensured, for the same reason as the `intended_did` indexes below. 0.3.9 put
-- it here and crash-looped production on its first boot: on an existing volume
-- the CREATE TABLE above is a no-op, so the column does not exist yet.
CREATE TABLE IF NOT EXISTS entries (
id INTEGER PRIMARY KEY AUTOINCREMENT,
feed_id INTEGER NOT NULL REFERENCES feeds (id) ON DELETE CASCADE,
guid TEXT NOT NULL,
url TEXT,
title TEXT,
author TEXT,
published TEXT,
content_html TEXT,
fetched_at TEXT NOT NULL,
UNIQUE (feed_id, guid)
);
-- The list and prev/next queries order on `COALESCE(published, fetched_at)`
-- (#187). Measured on the real query shape (LEFT JOIN entry_state, EXISTS
-- sub_ref), this index serves them as well as it served bare `published`; a
-- `(feed_id, published, fetched_at)` replacement was tried and was ~3.8x
-- slower on the default prev/next query, which never chose it (review of #213).
CREATE INDEX IF NOT EXISTS idx_entries_feed_published ON entries (feed_id, published);
CREATE TABLE IF NOT EXISTS entry_state (
did TEXT NOT NULL,
entry_id INTEGER NOT NULL REFERENCES entries (id) ON DELETE CASCADE,
read INTEGER NOT NULL DEFAULT 0,
starred INTEGER NOT NULL DEFAULT 0,
updated_at TEXT NOT NULL,
PRIMARY KEY (did, entry_id)
);
CREATE INDEX IF NOT EXISTS idx_entry_state_did_read ON entry_state (did, read);
-- The FK child key. `entry_id` is the TRAILING column of the primary key, so
-- without this index it is not the leading column of anything and SQLite must
-- FULL SCAN entry_state for EVERY row deleted from `entries` to service
-- ON DELETE CASCADE.
--
-- That is not theoretical. Measured on 600k entry_state rows: 500 deletes took
-- 10.3s and 2,000 took 38.3s, against a busy_timeout of 5s — so any retention
-- sweep removing more than roughly 260 entries made every concurrent writer
-- (star, mark-read, OAuth session write) fail with SQLITE_BUSY. With this index
-- the same 32,850-row delete goes from ~10 minutes to 0.7s.
--
-- It also fixes the per-feed trim, whose starred-sparing subquery scans
-- entry_state on every poll of every feed and scales with TOTAL rows across all
-- users rather than with the feed being trimmed (2ms -> 21ms at 1M rows).
CREATE INDEX IF NOT EXISTS idx_entry_state_entry_id ON entry_state (entry_id);
-- Per-DID subscription projection. The shared `feeds`/`entries` cache is
-- deduped by URL and NOT owned by any single DID; `sub_ref` records which
-- feeds a given DID actually subscribes to (mirrored from the caller's PDS
-- subscription set on every resolve/sync). Every entry/feed READ and every
-- read/star MUTATION is scoped through this table so one user can never read
-- or mutate another user's cached articles. Rows are refreshed by
-- `replace_sub_refs`.
CREATE TABLE IF NOT EXISTS sub_ref (
did TEXT NOT NULL,
feed_id INTEGER NOT NULL REFERENCES feeds (id) ON DELETE CASCADE,
PRIMARY KEY (did, feed_id)
);
CREATE INDEX IF NOT EXISTS idx_sub_ref_feed ON sub_ref (feed_id);
CREATE TABLE IF NOT EXISTS read_cursor (
did TEXT NOT NULL,
feed_url TEXT NOT NULL,
read_through TEXT,
read_ids TEXT NOT NULL DEFAULT '[]',
unread_ids TEXT NOT NULL DEFAULT '[]',
dirty INTEGER NOT NULL DEFAULT 0,
pds_created INTEGER NOT NULL DEFAULT 0,
updated_at TEXT NOT NULL,
PRIMARY KEY (did, feed_url)
);
CREATE INDEX IF NOT EXISTS idx_read_cursor_dirty ON read_cursor (did, dirty);
-- The (did, feed_url) PRIMARY KEY can't serve a feed_url-only lookup (did is the
-- leading column). The retention path's orphan-cursor cleanup filters cursors by
-- feed_url alone, so give it an index.
CREATE INDEX IF NOT EXISTS idx_read_cursor_feed_url ON read_cursor (feed_url);
CREATE TABLE IF NOT EXISTS beta_access (
did TEXT PRIMARY KEY,
handle TEXT,
granted_by TEXT NOT NULL,
granted_at INTEGER NOT NULL,
invite_code_used TEXT
);
CREATE TABLE IF NOT EXISTS invite_codes (
code TEXT PRIMARY KEY,
creator_did TEXT NOT NULL,
status TEXT NOT NULL,
invitee_did TEXT,
-- The follower DID a bot-minted claim was minted FOR (recorded at mint time,
-- distinct from `invitee_did` which is stamped at redeem). This is the
-- server-side idempotency key: a second `POST /bot/claims` for a DID that
-- already holds an outstanding active code returns the SAME code instead of
-- minting a duplicate, so a bot-host state loss cannot re-mint per follower.
intended_did TEXT,
created_at INTEGER NOT NULL,
expires_at INTEGER NOT NULL,
redeemed_at INTEGER
);
CREATE INDEX IF NOT EXISTS idx_invite_codes_status ON invite_codes (status, expires_at);
-- NOTE: the `intended_did` indexes are created in `apply_migrations`, AFTER the
-- `intended_did` column is ensured. They MUST NOT live in this base SCHEMA batch:
-- on an existing pre-0.2.2 volume the `CREATE TABLE IF NOT EXISTS invite_codes`
-- above is a no-op (the table already exists without `intended_did`), so a
-- `CREATE INDEX ... (intended_did, ...)` here would fail with "no such column"
-- and crash-loop the boot before migrations ever run.
-- Network-observation counters (v0.2.8, design/NETWORK-SPEC.md §4.3). One row
-- per (metric, relay): the adoption probe records how many repos a given relay
-- has INDEXED as holding a collection. We store the COUNT, never the DID list —
-- persisting the DIDs would build a durable register of "accounts that use an
-- RSS reader" on our disk for a feature whose only output is an integer. This
-- table is a PROJECTION, not a source of truth: `DROP TABLE` it and the next
-- probe rebuilds it, and nothing in the reader path reads it. Bounded forever at
-- (metrics × relays) rows, so it never interacts with the DB-size watermark.
CREATE TABLE IF NOT EXISTS network_stat (
key TEXT NOT NULL, -- e.g. 'adoption.subscription'
source TEXT NOT NULL, -- the relay host the number came from
value INTEGER NOT NULL,
truncated INTEGER NOT NULL DEFAULT 0,
observed_at TEXT NOT NULL,
PRIMARY KEY (key, source)
);
-- Repo-operation timings, for comparing the two backends across a CUTOVER.
--
-- Persisted rather than held in memory because flipping the backend requires a
-- restart, and an in-memory table would lose the outgoing backend's numbers at
-- exactly the moment they became worth comparing against. These rows are the
-- only reason a "side by side" table can show two backends at once.
--
-- `repo_timing` is a bounded window of recent samples (pruned per backend+op);
-- `repo_timing_total` carries the all-time counts, which must survive that
-- pruning or a long-running backend would appear to have served fewer calls
-- than a fresh one.
CREATE TABLE IF NOT EXISTS repo_timing (
id INTEGER PRIMARY KEY AUTOINCREMENT,
backend TEXT NOT NULL,
op TEXT NOT NULL,
micros INTEGER NOT NULL,
ok INTEGER NOT NULL,
at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_repo_timing_key ON repo_timing(backend, op, id);
CREATE TABLE IF NOT EXISTS repo_timing_total (
backend TEXT NOT NULL,
op TEXT NOT NULL,
ok_count INTEGER NOT NULL DEFAULT 0,
err_count INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (backend, op)
);
"#;
fn now_rfc3339() -> String {
chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
}
pub async fn init(config: &Config) -> Result<Pool> {
let db_url = format!("sqlite://{}", config.db_path.display());
init_url(&db_url).await
}
const WAL_SIZE_LIMIT_BYTES: i64 = 64 * 1024 * 1024;
pub async fn init_url(db_url: &str) -> Result<Pool> {
let is_memory = db_url.contains(":memory:");
let mut opts = SqliteConnectOptions::from_str(db_url)
.with_context(|| format!("invalid sqlite url: {db_url}"))?
.create_if_missing(true)
.foreign_keys(true);
if !is_memory {
opts = opts.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal);
opts = opts.auto_vacuum(sqlx::sqlite::SqliteAutoVacuum::Incremental);
opts = opts.pragma("journal_size_limit", WAL_SIZE_LIMIT_BYTES.to_string());
}
opts = opts.busy_timeout(std::time::Duration::from_millis(5000));
opts = opts.log_statements(tracing::log::LevelFilter::Debug);
let pool = SqlitePoolOptions::new()
.min_connections(1)
.max_connections(if is_memory { 1 } else { 5 })
.connect_with(opts)
.await
.with_context(|| format!("failed to open sqlite pool: {db_url}"))?;
init_schema(&pool).await?;
Ok(pool)
}
pub async fn init_schema(pool: &SqlitePool) -> Result<()> {
sqlx::query(SCHEMA)
.execute(pool)
.await
.context("failed to create schema")?;
apply_migrations(pool).await?;
crate::oauth::store::init_schema(pool)
.await
.context("failed to create the OAuth schema")?;
Ok(())
}
async fn apply_migrations(pool: &SqlitePool) -> Result<()> {
ensure_column(
pool,
"PRAGMA table_info(feeds)",
"consecutive_errors",
"ALTER TABLE feeds ADD COLUMN consecutive_errors INTEGER NOT NULL DEFAULT 0",
)
.await?;
ensure_column(
pool,
"PRAGMA table_info(feeds)",
"last_error_kind",
"ALTER TABLE feeds ADD COLUMN last_error_kind TEXT",
)
.await?;
ensure_column(
pool,
"PRAGMA table_info(feeds)",
"last_error",
"ALTER TABLE feeds ADD COLUMN last_error TEXT",
)
.await?;
ensure_column(
pool,
"PRAGMA table_info(feeds)",
"kind",
"ALTER TABLE feeds ADD COLUMN kind TEXT NOT NULL DEFAULT 'rss'",
)
.await?;
sqlx::query("CREATE INDEX IF NOT EXISTS idx_feeds_kind ON feeds (kind)")
.execute(pool)
.await
.context("creating idx_feeds_kind")?;
let rows = sqlx::query("SELECT id, url, kind FROM feeds")
.fetch_all(pool)
.await
.context("reading feeds to re-derive kind")?;
let mut tx = pool.begin().await.context("begin kind re-derivation")?;
let (mut to_pollable, mut to_unpollable, mut unreadable) = (0u64, 0u64, 0u64);
for row in rows {
let (Ok(id), Ok(url), Ok(kind)) = (
row.try_get::<i64, _>("id"),
row.try_get::<String, _>("url"),
row.try_get::<String, _>("kind"),
) else {
unreadable += 1;
continue;
};
let want = crate::feed::FeedKind::of(&url);
if kind == want.as_str() {
continue;
}
sqlx::query("UPDATE feeds SET kind = ?1 WHERE id = ?2")
.bind(want.as_str())
.bind(id)
.execute(&mut *tx)
.await
.with_context(|| format!("re-deriving kind for feed {id}"))?;
if crate::feed::FeedKind::POLLABLE.contains(&want) {
to_pollable += 1;
} else {
sqlx::query(
"UPDATE feeds SET consecutive_errors = 0, last_error_kind = NULL, \
last_error = NULL, next_poll = NULL WHERE id = ?1",
)
.bind(id)
.execute(&mut *tx)
.await
.with_context(|| format!("clearing orphaned poll state for feed {id}"))?;
to_unpollable += 1;
}
}
tx.commit().await.context("commit kind re-derivation")?;
if to_pollable > 0 || to_unpollable > 0 {
tracing::info!(
to_pollable,
to_unpollable,
"feeds.kind re-derived from the URL"
);
}
if unreadable > 0 {
tracing::warn!(
unreadable,
"feeds rows are not readable as text; their kind was left alone"
);
}
sqlx::query(sqlx::AssertSqlSafe(format!(
"UPDATE feeds SET consecutive_errors = 0, last_error_kind = NULL, last_error = NULL \
WHERE kind NOT IN ({POLLABLE_KINDS_SQL}) AND last_polled IS NULL \
AND consecutive_errors > 0"
)))
.execute(pool)
.await
.context("clearing error counts on unpollable at:// feeds")?;
ensure_column(
pool,
"PRAGMA table_info(read_cursor)",
"pds_created",
"ALTER TABLE read_cursor ADD COLUMN pds_created INTEGER NOT NULL DEFAULT 0",
)
.await?;
ensure_column(
pool,
"PRAGMA table_info(invite_codes)",
"intended_did",
"ALTER TABLE invite_codes ADD COLUMN intended_did TEXT",
)
.await?;
sqlx::query(
"CREATE INDEX IF NOT EXISTS idx_invite_codes_intended \
ON invite_codes (intended_did, status)",
)
.execute(pool)
.await
.context("creating idx_invite_codes_intended")?;
sqlx::query(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_invite_codes_intended_active \
ON invite_codes (intended_did) \
WHERE intended_did IS NOT NULL AND status = 'active'",
)
.execute(pool)
.await
.context("creating idx_invite_codes_intended_active")?;
let ceiling = (chrono::Utc::now()
+ chrono::Duration::days(crate::feed::MAX_FUTURE_PUBLISHED_DAYS))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
sqlx::query("UPDATE entries SET published = NULL WHERE published > ?1")
.bind(&ceiling)
.execute(pool)
.await
.context("clearing stored future publication dates")?;
Ok(())
}
async fn ensure_column(
pool: &SqlitePool,
info_sql: &'static str,
column: &str,
alter_sql: &'static str,
) -> Result<()> {
let rows = sqlx::query(info_sql)
.fetch_all(pool)
.await
.with_context(|| format!("{info_sql} failed"))?;
let present = rows.iter().any(|r| r.get::<String, _>("name") == column);
if !present {
sqlx::query(alter_sql)
.execute(pool)
.await
.with_context(|| format!("adding column {column} via {alter_sql}"))?;
}
Ok(())
}
pub async fn upsert_feed(pool: &SqlitePool, feed: &NewFeed) -> Result<i64> {
let row = sqlx::query(
r#"
INSERT INTO feeds (url, title, site_url, etag, last_modified, last_polled, next_poll, kind)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
ON CONFLICT (url) DO UPDATE SET
title = COALESCE(excluded.title, feeds.title),
site_url = COALESCE(excluded.site_url, feeds.site_url),
etag = COALESCE(excluded.etag, feeds.etag),
last_modified = COALESCE(excluded.last_modified, feeds.last_modified),
last_polled = COALESCE(excluded.last_polled, feeds.last_polled),
next_poll = COALESCE(excluded.next_poll, feeds.next_poll),
-- Not COALESCE: `kind` is derived from the URL, and `excluded`
-- always carries the current answer. Preserving the stored value
-- would make a row's classification a function of when it was
-- first subscribed rather than of what it is.
kind = excluded.kind
RETURNING id
"#,
)
.bind(&feed.url)
.bind(&feed.title)
.bind(&feed.site_url)
.bind(&feed.etag)
.bind(&feed.last_modified)
.bind(&feed.last_polled)
.bind(&feed.next_poll)
.bind(crate::feed::FeedKind::of(&feed.url).as_str())
.fetch_one(pool)
.await
.with_context(|| format!("upsert_feed failed for {}", feed.url))?;
Ok(row.get::<i64, _>("id"))
}
pub async fn get_feed_by_url(pool: &SqlitePool, url: &str) -> Result<Option<Feed>> {
let feed = sqlx::query_as::<_, Feed>("SELECT * FROM feeds WHERE url = ?1")
.bind(url)
.fetch_optional(pool)
.await
.with_context(|| format!("get_feed_by_url failed for {url}"))?;
Ok(feed)
}
pub(crate) const POLLABLE_KINDS_SQL: &str = "'rss', 'publication'";
pub(crate) const AGED_KINDS_SQL: &str = "'rss'";
pub async fn unpollable_feeds(pool: &SqlitePool) -> Result<i64> {
sqlx::query_scalar(sqlx::AssertSqlSafe(format!(
"SELECT COUNT(*) FROM feeds WHERE kind NOT IN ({POLLABLE_KINDS_SQL})"
)))
.fetch_one(pool)
.await
.context("counting unpollable feeds")
}
#[cfg(test)]
pub(crate) async fn count_unpollable_feeds(pool: &SqlitePool) -> Result<i64> {
sqlx::query_scalar(sqlx::AssertSqlSafe(format!(
"SELECT COUNT(*) FROM feeds WHERE kind NOT IN ({POLLABLE_KINDS_SQL})"
)))
.fetch_one(pool)
.await
.context("counting unpollable feeds")
}
pub async fn due_feeds(pool: &SqlitePool, as_of: &str, limit: i64) -> Result<Vec<Feed>> {
let sql = format!(
r#"
SELECT * FROM feeds
WHERE (next_poll IS NULL OR next_poll <= ?1)
-- Only POLLABLE kinds are due; an `unsupported` row is skipped, not
-- failed. The why lives on `feed::FeedKind::POLLABLE`.
AND kind IN ({POLLABLE_KINDS_SQL})
ORDER BY next_poll IS NOT NULL, next_poll ASC
LIMIT ?2
"#
);
let feeds = sqlx::query_as::<_, Feed>(sqlx::AssertSqlSafe(sql))
.bind(as_of)
.bind(limit)
.fetch_all(pool)
.await
.context("due_feeds failed")?;
Ok(feeds)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FailingFeed {
pub url: String,
pub consecutive_errors: i64,
pub kind: Option<String>,
pub detail: Option<String>,
}
pub async fn failing_feeds(pool: &SqlitePool, limit: i64) -> Result<Vec<FailingFeed>> {
let sql = format!(
r#"
SELECT url, consecutive_errors, last_error_kind, last_error
FROM feeds
WHERE consecutive_errors > 0 AND kind IN ({POLLABLE_KINDS_SQL})
ORDER BY consecutive_errors DESC, url ASC
LIMIT ?1
"#
);
let rows: Vec<(String, i64, Option<String>, Option<String>)> =
sqlx::query_as(sqlx::AssertSqlSafe(sql))
.bind(limit)
.fetch_all(pool)
.await
.context("listing failing feeds")?;
Ok(rows
.into_iter()
.map(|(url, consecutive_errors, kind, detail)| FailingFeed {
url,
consecutive_errors,
kind,
detail,
})
.collect())
}
const MAX_ERROR_DETAIL_CHARS: usize = 300;
pub async fn bump_feed_errors(
pool: &SqlitePool,
url: &str,
kind: crate::feed::FailureKind,
detail: &str,
) -> Result<i64> {
let row = sqlx::query(
"UPDATE feeds SET consecutive_errors = consecutive_errors + 1, \
last_error_kind = ?2, last_error = ?3 \
WHERE url = ?1 RETURNING consecutive_errors",
)
.bind(url)
.bind(kind.as_str())
.bind(
detail
.chars()
.take(MAX_ERROR_DETAIL_CHARS)
.collect::<String>(),
)
.fetch_optional(pool)
.await
.with_context(|| format!("bump_feed_errors failed for {url}"))?;
Ok(row
.map(|r| r.get::<i64, _>("consecutive_errors"))
.unwrap_or(1))
}
pub async fn set_next_poll(pool: &SqlitePool, url: &str, delay: std::time::Duration) -> Result<()> {
let next = chrono::Utc::now()
+ chrono::Duration::from_std(delay).unwrap_or_else(|_| chrono::Duration::hours(1));
let next_poll = next.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let nf = NewFeed {
url: url.to_string(),
next_poll: Some(next_poll),
..Default::default()
};
upsert_feed(pool, &nf).await.map(|_| ())
}
pub async fn due_feeds_of_kind(
pool: &SqlitePool,
as_of: &str,
kind: crate::feed::FeedKind,
limit: i64,
) -> Result<Vec<Feed>> {
sqlx::query_as::<_, Feed>(
"SELECT * FROM feeds WHERE (next_poll IS NULL OR next_poll <= ?1) AND kind = ?2 \
ORDER BY next_poll IS NOT NULL, next_poll ASC LIMIT ?3",
)
.bind(as_of)
.bind(kind.as_str())
.bind(limit)
.fetch_all(pool)
.await
.context("due_feeds_of_kind failed")
}
pub async fn stagger_unscheduled(
pool: &SqlitePool,
kind: crate::feed::FeedKind,
spread: std::time::Duration,
) -> Result<u64> {
let ids: Vec<i64> = sqlx::query_scalar(
"SELECT id FROM feeds WHERE kind = ?1 AND next_poll IS NULL AND last_polled IS NULL \
ORDER BY id",
)
.bind(kind.as_str())
.fetch_all(pool)
.await
.context("listing unscheduled feeds to stagger")?;
if ids.is_empty() {
return Ok(0);
}
let now = chrono::Utc::now();
let spread = chrono::Duration::from_std(spread).unwrap_or_else(|_| chrono::Duration::hours(1));
let n = ids.len() as i32;
let mut tx = pool.begin().await.context("begin stagger")?;
for (i, id) in ids.iter().enumerate() {
let at = now + spread * i as i32 / n;
let next_poll = at.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
sqlx::query("UPDATE feeds SET next_poll = ?1 WHERE id = ?2 AND next_poll IS NULL")
.bind(next_poll)
.bind(id)
.execute(&mut *tx)
.await
.context("staggering a feed's first poll")?;
}
tx.commit().await.context("commit stagger")?;
Ok(ids.len() as u64)
}
pub async fn reset_feed_errors(pool: &SqlitePool, url: &str) -> Result<()> {
sqlx::query(
"UPDATE feeds SET consecutive_errors = 0, last_error_kind = NULL, last_error = NULL \
WHERE url = ?1",
)
.bind(url)
.execute(pool)
.await
.with_context(|| format!("reset_feed_errors failed for {url}"))?;
Ok(())
}
pub async fn feeds_for_did(pool: &SqlitePool, did: &str) -> Result<Vec<Feed>> {
let feeds = sqlx::query_as::<_, Feed>(
r#"
SELECT f.* FROM feeds f
JOIN sub_ref sr ON sr.feed_id = f.id AND sr.did = ?1
ORDER BY f.title IS NULL, f.title, f.url
"#,
)
.bind(did)
.fetch_all(pool)
.await
.with_context(|| format!("feeds_for_did failed for {did}"))?;
Ok(feeds)
}
pub async fn subscribed_feed_ids(pool: &SqlitePool, did: &str) -> Result<Vec<i64>> {
let ids: Vec<i64> = sqlx::query_scalar("SELECT feed_id FROM sub_ref WHERE did = ?1")
.bind(did)
.fetch_all(pool)
.await
.with_context(|| format!("subscribed_feed_ids failed for {did}"))?;
Ok(ids)
}
pub async fn count_subscriptions_for_did(pool: &SqlitePool, did: &str) -> Result<i64> {
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM sub_ref WHERE did = ?1")
.bind(did)
.fetch_one(pool)
.await
.with_context(|| format!("count_subscriptions_for_did failed for {did}"))?;
Ok(n)
}
pub async fn count_feeds(pool: &SqlitePool) -> Result<i64> {
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM feeds")
.fetch_one(pool)
.await
.context("count_feeds failed")?;
Ok(n)
}
pub async fn db_size_bytes(pool: &SqlitePool) -> Result<i64> {
let page_count: i64 = sqlx::query_scalar("PRAGMA page_count")
.fetch_one(pool)
.await
.context("PRAGMA page_count failed")?;
let freelist_count: i64 = sqlx::query_scalar("PRAGMA freelist_count")
.fetch_one(pool)
.await
.context("PRAGMA freelist_count failed")?;
let page_size: i64 = sqlx::query_scalar("PRAGMA page_size")
.fetch_one(pool)
.await
.context("PRAGMA page_size failed")?;
let used_pages = page_count.saturating_sub(freelist_count).max(0);
Ok(used_pages
.saturating_mul(page_size)
.saturating_add(wal_bytes(pool).await))
}
async fn wal_bytes(pool: &SqlitePool) -> i64 {
let Some(path) = main_db_path(pool).await else {
return 0;
};
std::fs::metadata(format!("{path}-wal"))
.map(|m| i64::try_from(m.len()).unwrap_or(i64::MAX))
.unwrap_or(0)
}
async fn main_db_path(pool: &SqlitePool) -> Option<String> {
sqlx::query_scalar("SELECT file FROM pragma_database_list WHERE name = 'main' AND file <> ''")
.fetch_optional(pool)
.await
.ok()
.flatten()
}
const RECLAIM_BATCH_PAGES: i64 = 2_000;
const RECLAIM_MAX_BATCHES: usize = 1_000;
pub async fn reclaim(pool: &SqlitePool) -> Result<()> {
match auto_vacuum_mode(pool).await? {
AutoVacuum::Incremental => {
let mut drained = true;
for batch in 0..RECLAIM_MAX_BATCHES {
let before: i64 = sqlx::query_scalar("PRAGMA freelist_count")
.fetch_one(pool)
.await
.context("PRAGMA freelist_count failed")?;
if before == 0 {
break;
}
sqlx::query(sqlx::AssertSqlSafe(format!(
"PRAGMA incremental_vacuum({RECLAIM_BATCH_PAGES})"
)))
.execute(pool)
.await
.context("PRAGMA incremental_vacuum failed")?;
let after: i64 = sqlx::query_scalar("PRAGMA freelist_count")
.fetch_one(pool)
.await
.context("PRAGMA freelist_count failed")?;
if after >= before {
tracing::warn!(
freelist_pages = after,
batches_run = batch + 1,
"reclaim stopped making progress with pages still on the \
freelist; the file will not shrink and the DB-size watermark \
may stay engaged until the next sweep"
);
drained = false;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
if batch + 1 == RECLAIM_MAX_BATCHES && after > 0 {
tracing::warn!(
batches_run = batch + 1,
freelist_pages = after,
"reclaim hit its batch backstop with pages still on the \
freelist; the rest waits for the next sweep"
);
drained = false;
}
}
if drained {
tracing::debug!("reclaim: freelist drained");
}
}
AutoVacuum::Full => {}
AutoVacuum::None => {
tracing::warn!(
"auto_vacuum=NONE: skipping reclaim. Freed pages stay allocated and \
the file will not shrink. Run `featherreader --migrate-auto-vacuum` \
once, while the volume has headroom, to move this database to \
INCREMENTAL mode."
);
}
}
match checkpoint_wal(pool).await {
Ok(true) => {}
Ok(false) => tracing::warn!(
"the WAL could not be truncated after reclaim (busy: a concurrent reader \
OR writer held it); the freed pages are gone but the file has not \
shrunk yet, and the DB-size watermark may stay engaged until the next \
sweep"
),
Err(err) => tracing::warn!(%err, "wal checkpoint after reclaim failed"),
}
Ok(())
}
async fn checkpoint_wal<'e, E>(conn: E) -> Result<bool>
where
E: sqlx::Executor<'e, Database = sqlx::Sqlite>,
{
let row: (i64, i64, i64) = sqlx::query_as("PRAGMA wal_checkpoint(TRUNCATE)")
.fetch_one(conn)
.await
.context("PRAGMA wal_checkpoint(TRUNCATE) failed")?;
Ok(row.0 == 0)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AutoVacuum {
None,
Full,
Incremental,
}
pub async fn auto_vacuum_mode(pool: &SqlitePool) -> Result<AutoVacuum> {
let mode: i64 = sqlx::query_scalar("PRAGMA auto_vacuum")
.fetch_one(pool)
.await
.context("PRAGMA auto_vacuum failed")?;
Ok(match mode {
1 => AutoVacuum::Full,
2 => AutoVacuum::Incremental,
_ => AutoVacuum::None,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum VacuumMigration {
NotNeeded(AutoVacuum),
RefusedNoHeadroom {
needed: u64,
available: u64,
file_bytes: Option<u64>,
},
Migrated {
bytes_before: i64,
bytes_after: i64,
file_before: Option<u64>,
file_after: Option<u64>,
},
}
pub async fn migrate_to_incremental_vacuum(
pool: &SqlitePool,
available_bytes: Option<u64>,
) -> Result<VacuumMigration> {
let mode = auto_vacuum_mode(pool).await?;
if mode != AutoVacuum::None {
return Ok(VacuumMigration::NotNeeded(mode));
}
let bytes_before = db_size_bytes(pool).await?;
let file_before = main_db_file_bytes(pool).await;
let temp_dir = main_db_path(pool).await.and_then(|p| {
std::path::Path::new(&p)
.parent()
.map(std::path::Path::to_path_buf)
});
let needed = (bytes_before.max(0) as u64).saturating_mul(2);
if let Some(available) = available_bytes {
if available < needed {
return Ok(VacuumMigration::RefusedNoHeadroom {
needed,
available,
file_bytes: file_before,
});
}
}
let mut conn = pool
.acquire()
.await
.context("acquiring a connection for the auto_vacuum migration")?;
sqlx::query("PRAGMA temp_store = FILE")
.execute(&mut *conn)
.await
.context("PRAGMA temp_store = FILE failed")?;
if let Some(dir) = temp_dir.clone() {
let quoted = dir.display().to_string().replace('\'', "''");
if let Err(err) = sqlx::query(sqlx::AssertSqlSafe(format!(
"PRAGMA temp_store_directory = '{quoted}'"
)))
.execute(&mut *conn)
.await
{
tracing::warn!(
%err, dir = %dir.display(),
"could not point SQLite's temp storage at the database volume; the \
headroom check may not cover where the VACUUM actually writes"
);
}
}
sqlx::query("PRAGMA auto_vacuum = INCREMENTAL")
.execute(&mut *conn)
.await
.context("PRAGMA auto_vacuum = INCREMENTAL failed")?;
sqlx::query("VACUUM")
.execute(&mut *conn)
.await
.context("VACUUM failed during the auto_vacuum migration")?;
match checkpoint_wal(&mut *conn).await {
Ok(true) => {}
Ok(false) => tracing::warn!(
"the WAL could not be truncated (a concurrent reader OR writer held it), \
so the reported size below includes it"
),
Err(err) => tracing::warn!(%err, "post-migration wal checkpoint failed"),
}
let after_raw: i64 = sqlx::query_scalar("PRAGMA auto_vacuum")
.fetch_one(&mut *conn)
.await
.context("PRAGMA auto_vacuum failed after the migration")?;
drop(conn);
let after = match after_raw {
1 => AutoVacuum::Full,
2 => AutoVacuum::Incremental,
_ => AutoVacuum::None,
};
anyhow::ensure!(
after == AutoVacuum::Incremental,
"the auto_vacuum migration ran but the database is still in {after:?} mode"
);
Ok(VacuumMigration::Migrated {
bytes_before,
bytes_after: db_size_bytes(pool).await?,
file_before,
file_after: main_db_file_bytes(pool).await,
})
}
async fn main_db_file_bytes(pool: &SqlitePool) -> Option<u64> {
let path = main_db_path(pool).await?;
std::fs::metadata(path).ok().map(|m| m.len())
}
pub async fn insert_entries(
pool: &SqlitePool,
feed_id: i64,
entries: &[NewEntry],
max_entries_per_feed: i64,
) -> Result<u64> {
let mut tx = pool.begin().await.context("begin insert_entries tx")?;
let mut count: u64 = 0;
for e in entries {
let fetched_at = e.fetched_at.clone().unwrap_or_else(now_rfc3339);
let res = sqlx::query(
r#"
INSERT INTO entries
(feed_id, guid, url, title, author, published, content_html, fetched_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
ON CONFLICT (feed_id, guid) DO UPDATE SET
url = excluded.url,
title = excluded.title,
author = excluded.author,
published = excluded.published,
content_html = CASE WHEN ?9 THEN entries.content_html
ELSE excluded.content_html END
"#,
)
.bind(feed_id)
.bind(&e.guid)
.bind(&e.url)
.bind(&e.title)
.bind(&e.author)
.bind(&e.published)
.bind(&e.content_html)
.bind(&fetched_at)
.bind(e.keep_stored_content)
.execute(&mut *tx)
.await
.with_context(|| format!("insert entry {} failed", e.guid))?;
count += res.rows_affected();
}
if max_entries_per_feed > 0 {
sqlx::query(
r#"
DELETE FROM entries
WHERE feed_id = ?1
AND id NOT IN (
SELECT id FROM entries
WHERE feed_id = ?1
ORDER BY COALESCE(published, fetched_at) DESC, id DESC
LIMIT ?2
)
-- Starred entries survive the per-feed trim, exactly as they
-- survive the retention sweep. This predicate was added to the
-- sweep and NOT here, which left the documented guarantee
-- ("starred entries are never evicted") false — and made this
-- path, which runs on every poll of every feed rather than daily,
-- the main producer of the very "starred but not cached" case the
-- saved-record rendering exists to paper over.
--
-- The sparing is BOUNDED and SCOPED, and both matter:
--
-- Bounded, because the first version spared every starred row
-- without limit, which did not weaken the cap so much as remove
-- it — measured at cap=5 with 50 starred rows, 55 survived, 11x
-- the cap. That is the same unbounded-sparing mistake the
-- retention hard ceiling was added to fix, reintroduced in the
-- other sweep. Worst case is now cap + cap.
--
-- Scoped, because `SELECT entry_id FROM entry_state WHERE
-- starred = 1` reads EVERY starred row on the instance, for every
-- poll of every feed — cost scaling with total users rather than
-- with the feed being trimmed.
AND id NOT IN (
SELECT e2.id FROM entries e2
WHERE e2.feed_id = ?1
AND EXISTS (
SELECT 1 FROM entry_state s
WHERE s.entry_id = e2.id AND s.starred = 1
)
ORDER BY COALESCE(e2.published, e2.fetched_at) DESC, e2.id DESC
LIMIT ?2
)
"#,
)
.bind(feed_id)
.bind(max_entries_per_feed)
.execute(&mut *tx)
.await
.with_context(|| format!("trimming feed {feed_id} to {max_entries_per_feed} entries"))?;
}
if max_entries_per_feed > 0 {
prune_orphan_cursor_ids_tx(&mut tx, Some(feed_id)).await?;
}
tx.commit().await.context("commit insert_entries tx")?;
Ok(count)
}
pub async fn mark_feed_due(
pool: &SqlitePool,
feed_url: &str,
not_polled_since: &str,
) -> Result<()> {
sqlx::query(
"UPDATE feeds SET next_poll = NULL \
WHERE url = ?1 AND (last_polled IS NULL OR last_polled < ?2)",
)
.bind(feed_url)
.bind(not_polled_since)
.execute(pool)
.await
.context("marking a feed due")?;
Ok(())
}
pub async fn prune_old_entries(
pool: &SqlitePool,
days: i64,
hard_days: i64,
publication_days: i64,
) -> Result<u64> {
let now = chrono::Utc::now();
let at = |d: i64, knob: &str| -> Option<String> {
let cutoff = chrono::Duration::try_days(d).and_then(|w| now.checked_sub_signed(w));
if cutoff.is_none() {
tracing::warn!(
days = d,
knob,
"retention window is too large to express as a date; treating it as \
disabled for this sweep rather than failing the sweeper"
);
}
cutoff.map(|t| t.to_rfc3339_opts(chrono::SecondsFormat::Secs, true))
};
let cutoff = (days > 0).then(|| at(days, "retention_days")).flatten();
let publication_cutoff = (publication_days > 0)
.then(|| at(publication_days, "publication_retention_days"))
.flatten();
let hard_cutoff = if hard_days > 0 && (days <= 0 || hard_days > days) {
at(hard_days, "retention_hard_days")
} else {
if hard_days > 0 {
tracing::warn!(
hard_days,
days,
"retention hard ceiling is not older than the retention window; \
ignoring it — set it above the window or to 0 to disable"
);
}
None
};
if cutoff.is_none() && hard_cutoff.is_none() && publication_cutoff.is_none() {
return Ok(0);
}
let hard_deleted = match &hard_cutoff {
Some(cutoff) => {
delete_in_batches(
pool,
&format!(
"SELECT id FROM entries WHERE COALESCE(published, fetched_at) < ?1 \
AND feed_id IN (SELECT id FROM feeds WHERE kind IN ({AGED_KINDS_SQL}))"
),
cutoff,
"hard ceiling",
)
.await?
}
None => 0,
};
let soft_deleted = match &cutoff {
Some(cutoff) => {
delete_in_batches(
pool,
&format!(
"SELECT e.id FROM entries e \
WHERE COALESCE(e.published, e.fetched_at) < ?1 \
AND e.feed_id IN \
(SELECT id FROM feeds WHERE kind IN ({AGED_KINDS_SQL})) \
AND NOT EXISTS ( \
SELECT 1 FROM entry_state s \
WHERE s.entry_id = e.id \
AND (s.starred = 1 OR s.read = 0) \
)"
),
cutoff,
"window",
)
.await?
}
None => 0,
};
let publication_deleted = match &publication_cutoff {
Some(cutoff) => {
delete_in_batches(
pool,
&format!(
"SELECT id FROM entries WHERE COALESCE(published, fetched_at) < ?1 \
AND feed_id IN (SELECT id FROM feeds WHERE kind NOT IN ({AGED_KINDS_SQL}))"
),
cutoff,
"archive ceiling",
)
.await?
}
None => 0,
};
let deleted = soft_deleted + hard_deleted + publication_deleted;
if deleted > 0 {
if let Err(err) = prune_orphan_cursor_ids(pool, None).await {
tracing::warn!(%err, "retention sweep: cursor id scrub failed after the deletes");
}
}
Ok(deleted)
}
const PRUNE_BATCH: i64 = 1_000;
const PRUNE_MAX_BATCHES: usize = 10_000;
const PRUNE_BATCH_HANDOFF: std::time::Duration = std::time::Duration::from_millis(10);
async fn delete_in_batches(
pool: &SqlitePool,
select_ids: &str,
cutoff: &str,
label: &str,
) -> Result<u64> {
let sql = format!("DELETE FROM entries WHERE id IN ({select_ids} LIMIT {PRUNE_BATCH})");
let mut total: u64 = 0;
for batch in 0..PRUNE_MAX_BATCHES {
let n = sqlx::query(sqlx::AssertSqlSafe(sql.clone()))
.bind(cutoff)
.execute(pool)
.await
.with_context(|| format!("prune_old_entries {label} (cutoff {cutoff})"))?
.rows_affected();
total += n;
if n == 0 {
return Ok(total);
}
tokio::time::sleep(PRUNE_BATCH_HANDOFF).await;
if batch + 1 == PRUNE_MAX_BATCHES && n == PRUNE_BATCH as u64 {
tracing::warn!(
label,
total,
"retention sweep hit its batch backstop; the rest waits for the next run"
);
}
}
Ok(total)
}
async fn prune_orphan_cursor_ids_tx(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
feed_id: Option<i64>,
) -> Result<u64> {
let feed_url = match feed_id {
Some(fid) => match feed_url_for_id_tx(tx, fid).await? {
Some(u) => Some(u),
None => return Ok(0), },
None => None,
};
let cursors: Vec<(String, String, String, String)> = match &feed_url {
Some(url) => sqlx::query(
"SELECT did, feed_url, read_ids, unread_ids FROM read_cursor WHERE feed_url = ?1",
)
.bind(url)
.fetch_all(&mut **tx)
.await
.context("prune_orphan_cursor_ids: load feed cursors")?,
None => sqlx::query("SELECT did, feed_url, read_ids, unread_ids FROM read_cursor")
.fetch_all(&mut **tx)
.await
.context("prune_orphan_cursor_ids: load all cursors")?,
}
.into_iter()
.map(|r| {
(
r.get::<String, _>("did"),
r.get::<String, _>("feed_url"),
r.get::<String, _>("read_ids"),
r.get::<String, _>("unread_ids"),
)
})
.collect();
if cursors.is_empty() {
return Ok(0);
}
let now = now_rfc3339();
let mut changed: u64 = 0;
for (did, curl, read_ids, unread_ids) in cursors {
let live: std::collections::HashSet<i64> = sqlx::query_scalar::<_, i64>(
"SELECT e.id FROM entries e JOIN feeds f ON f.id = e.feed_id WHERE f.url = ?1",
)
.bind(&curl)
.fetch_all(&mut **tx)
.await
.with_context(|| format!("prune_orphan_cursor_ids: live ids for {curl}"))?
.into_iter()
.collect();
let new_read = filter_id_set_to_live(&read_ids, &live);
let new_unread = filter_id_set_to_live(&unread_ids, &live);
if new_read == read_ids && new_unread == unread_ids {
continue; }
sqlx::query(
"UPDATE read_cursor SET read_ids = ?3, unread_ids = ?4, dirty = 1, updated_at = ?5 \
WHERE did = ?1 AND feed_url = ?2",
)
.bind(&did)
.bind(&curl)
.bind(&new_read)
.bind(&new_unread)
.bind(&now)
.execute(&mut **tx)
.await
.with_context(|| format!("prune_orphan_cursor_ids: rewrite cursor {did}/{curl}"))?;
changed += 1;
}
Ok(changed)
}
async fn prune_orphan_cursor_ids(pool: &SqlitePool, feed_id: Option<i64>) -> Result<u64> {
let feed_url = match feed_id {
Some(fid) => match sqlx::query_scalar::<_, String>("SELECT url FROM feeds WHERE id = ?1")
.bind(fid)
.fetch_optional(pool)
.await
.context("prune_orphan_cursor_ids: feed url")?
{
Some(u) => Some(u),
None => return Ok(0),
},
None => None,
};
let keys: Vec<(String, String)> = match &feed_url {
Some(url) => sqlx::query_as("SELECT did, feed_url FROM read_cursor WHERE feed_url = ?1")
.bind(url)
.fetch_all(pool)
.await
.context("prune_orphan_cursor_ids: load feed cursors")?,
None => sqlx::query_as("SELECT did, feed_url FROM read_cursor")
.fetch_all(pool)
.await
.context("prune_orphan_cursor_ids: load all cursors")?,
};
let mut changed: u64 = 0;
for (did, curl) in keys {
match scrub_one_cursor(pool, &did, &curl).await {
Ok(true) => changed += 1,
Ok(false) => {}
Err(err) => tracing::warn!(%err, %did, feed = %curl, "cursor id scrub failed"),
}
}
Ok(changed)
}
async fn scrub_one_cursor(pool: &SqlitePool, did: &str, feed_url: &str) -> Result<bool> {
let mut tx = pool.begin().await.context("begin scrub_one_cursor tx")?;
let (_, read_ids, unread_ids) = cursor_sets(&mut tx, did, feed_url).await?;
if is_empty_id_set(&read_ids) && is_empty_id_set(&unread_ids) {
return Ok(false);
}
let live: std::collections::HashSet<i64> = sqlx::query_scalar::<_, i64>(
"SELECT e.id FROM entries e JOIN feeds f ON f.id = e.feed_id WHERE f.url = ?1",
)
.bind(feed_url)
.fetch_all(&mut *tx)
.await
.with_context(|| format!("prune_orphan_cursor_ids: live ids for {feed_url}"))?
.into_iter()
.collect();
let new_read = filter_id_set_to_live(&read_ids, &live);
let new_unread = filter_id_set_to_live(&unread_ids, &live);
if new_read == read_ids && new_unread == unread_ids {
return Ok(false); }
sqlx::query(
"UPDATE read_cursor SET read_ids = ?3, unread_ids = ?4, dirty = 1, updated_at = ?5 \
WHERE did = ?1 AND feed_url = ?2",
)
.bind(did)
.bind(feed_url)
.bind(&new_read)
.bind(&new_unread)
.bind(now_rfc3339())
.execute(&mut *tx)
.await
.with_context(|| format!("prune_orphan_cursor_ids: rewrite cursor {did}/{feed_url}"))?;
tx.commit().await.context("commit scrub_one_cursor tx")?;
Ok(true)
}
fn is_empty_id_set(raw: &str) -> bool {
let t = raw.trim();
t.is_empty() || t == "[]"
}
fn filter_id_set_to_live(raw: &str, live: &std::collections::HashSet<i64>) -> String {
let ids: Vec<i64> = serde_json::from_str::<Vec<serde_json::Value>>(raw)
.ok()
.map(|vals| {
vals.into_iter()
.filter_map(|v| match v {
serde_json::Value::Number(n) => n.as_i64(),
serde_json::Value::String(s) => s.parse::<i64>().ok(),
_ => None,
})
.filter(|id| live.contains(id))
.collect()
})
.unwrap_or_default();
let as_strings: Vec<String> = ids.iter().map(|i| i.to_string()).collect();
serde_json::to_string(&as_strings).unwrap_or_else(|_| "[]".to_string())
}
pub async fn replace_sub_refs(pool: &SqlitePool, did: &str, feed_ids: &[i64]) -> Result<()> {
let mut tx = pool.begin().await.context("begin replace_sub_refs tx")?;
sqlx::query("DELETE FROM sub_ref WHERE did = ?1")
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("clear sub_ref for {did}"))?;
for &feed_id in feed_ids {
sqlx::query("INSERT OR IGNORE INTO sub_ref (did, feed_id) VALUES (?1, ?2)")
.bind(did)
.bind(feed_id)
.execute(&mut *tx)
.await
.with_context(|| format!("insert sub_ref {did}/{feed_id}"))?;
}
tx.commit().await.context("commit replace_sub_refs tx")?;
Ok(())
}
pub async fn did_subscribes_to_entry(pool: &SqlitePool, did: &str, entry_id: i64) -> Result<bool> {
let found: Option<i64> = sqlx::query_scalar(
r#"
SELECT 1
FROM entries e
JOIN sub_ref sr ON sr.feed_id = e.feed_id AND sr.did = ?1
WHERE e.id = ?2
"#,
)
.bind(did)
.bind(entry_id)
.fetch_optional(pool)
.await
.with_context(|| format!("did_subscribes_to_entry failed for {did}/{entry_id}"))?;
Ok(found.is_some())
}
fn list_entries_sql(view: ListView, feed_ids: Option<&[i64]>) -> (String, usize) {
list_query_sql(Projection::EntryList, view, feed_ids)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Projection {
EntryList,
Count,
Ids,
FeedCounts,
StarredUrls,
}
impl Projection {
const fn columns(self) -> &'static str {
match self {
Projection::EntryList => {
"e.id, e.feed_id, e.guid, e.url, e.title, e.published, \
COALESCE(s.read, 0) AS read, COALESCE(s.starred, 0) AS starred"
}
Projection::Count => "COUNT(*)",
Projection::Ids => "e.id",
Projection::FeedCounts => "e.feed_id, COUNT(*)",
Projection::StarredUrls => "e.url, e.guid",
}
}
}
fn list_query_sql(
projection: Projection,
view: ListView,
feed_ids: Option<&[i64]>,
) -> (String, usize) {
let cols = projection.columns();
let scoped = feed_ids.is_some();
let mut sql = format!(
"SELECT {cols} \
FROM entries e \
LEFT JOIN entry_state s ON s.entry_id = e.id AND s.did = ?1 \
WHERE {} \
AND EXISTS ( \
SELECT 1 FROM sub_ref sr \
WHERE sr.did = ?1 AND sr.feed_id = e.feed_id \
)",
view.predicate()
);
if scoped {
sql.push_str(" AND e.feed_id IN (SELECT value FROM json_each(?2))");
}
(sql, usize::from(scoped))
}
fn bind_list_scope<'q, O>(
q: sqlx::query::QueryAs<'q, sqlx::Sqlite, O, sqlx::sqlite::SqliteArguments>,
did: &'q str,
feed_ids: Option<&[i64]>,
) -> sqlx::query::QueryAs<'q, sqlx::Sqlite, O, sqlx::sqlite::SqliteArguments> {
let q = q.bind(did);
match feed_ids {
Some(ids) => q.bind(serde_json::to_string(ids).unwrap_or_else(|_| "[]".to_string())),
None => q,
}
}
pub async fn list_entries(
pool: &SqlitePool,
did: &str,
view: ListView,
feed_ids: Option<&[i64]>,
limit: i64,
offset: i64,
) -> Result<Vec<EntryListRow>> {
if feed_ids.is_some_and(<[i64]>::is_empty) || limit <= 0 {
return Ok(Vec::new());
}
let (mut sql, n) = list_entries_sql(view, feed_ids);
sql.push_str(&format!(
" ORDER BY COALESCE(e.published, e.fetched_at) DESC, e.id DESC LIMIT ?{} OFFSET ?{}",
n + 2,
n + 3
));
let q = sqlx::query_as::<_, EntryListRow>(sqlx::AssertSqlSafe(sql));
let rows = bind_list_scope(q, did, feed_ids)
.bind(limit)
.bind(offset.max(0))
.fetch_all(pool)
.await
.with_context(|| format!("list_entries({view:?}) failed for {did}"))?;
Ok(rows)
}
pub async fn count_entries_for_view(
pool: &SqlitePool,
did: &str,
view: ListView,
feed_ids: Option<&[i64]>,
) -> Result<i64> {
if feed_ids.is_some_and(<[i64]>::is_empty) {
return Ok(0);
}
let (sql, _) = list_query_sql(Projection::Count, view, feed_ids);
let q = sqlx::query_as::<_, (i64,)>(sqlx::AssertSqlSafe(sql));
let (n,) = bind_list_scope(q, did, feed_ids)
.fetch_one(pool)
.await
.with_context(|| format!("count_entries_for_view({view:?}) failed for {did}"))?;
Ok(n)
}
pub async fn list_entry_ids(
pool: &SqlitePool,
did: &str,
view: ListView,
feed_ids: Option<&[i64]>,
limit: i64,
) -> Result<Vec<i64>> {
if feed_ids.is_some_and(<[i64]>::is_empty) || limit <= 0 {
return Ok(Vec::new());
}
let (mut sql, n) = list_query_sql(Projection::Ids, view, feed_ids);
sql.push_str(&format!(
" ORDER BY COALESCE(e.published, e.fetched_at) DESC, e.id DESC LIMIT ?{}",
n + 2
));
let q = sqlx::query_as::<_, (i64,)>(sqlx::AssertSqlSafe(sql));
let rows = bind_list_scope(q, did, feed_ids)
.bind(limit)
.fetch_all(pool)
.await
.with_context(|| format!("list_entry_ids({view:?}) failed for {did}"))?;
Ok(rows.into_iter().map(|(id,)| id).collect())
}
pub async fn unread_counts_by_feed(
pool: &SqlitePool,
did: &str,
) -> Result<std::collections::HashMap<i64, i64>> {
let (sql, _) = list_query_sql(Projection::FeedCounts, ListView::Unread, None);
let rows =
sqlx::query_as::<_, (i64, i64)>(sqlx::AssertSqlSafe(format!("{sql} GROUP BY e.feed_id")))
.bind(did)
.fetch_all(pool)
.await
.with_context(|| format!("unread_counts_by_feed failed for {did}"))?;
Ok(rows.into_iter().collect())
}
pub enum StarredIdentities {
All(Vec<(Option<String>, String)>),
Truncated,
}
pub async fn starred_identities(
pool: &SqlitePool,
did: &str,
limit: i64,
) -> Result<StarredIdentities> {
let (mut sql, n) = list_query_sql(Projection::StarredUrls, ListView::Starred, None);
sql.push_str(&format!(" ORDER BY e.id LIMIT ?{}", n + 2));
let rows = sqlx::query_as::<_, (Option<String>, String)>(sqlx::AssertSqlSafe(sql))
.bind(did)
.bind(limit.saturating_add(1))
.fetch_all(pool)
.await
.with_context(|| format!("starred_identities failed for {did}"))?;
if rows.len() as i64 > limit {
return Ok(StarredIdentities::Truncated);
}
Ok(StarredIdentities::All(rows))
}
pub async fn mark_read(pool: &SqlitePool, did: &str, entry_id: i64, read: bool) -> Result<bool> {
let now = now_rfc3339();
let mut tx = pool.begin().await.context("begin mark_read tx")?;
let res = sqlx::query(
r#"
INSERT INTO entry_state (did, entry_id, read, starred, updated_at)
SELECT ?1, e.id, ?3, 0, ?4
FROM entries e
WHERE e.id = ?2
AND EXISTS (
SELECT 1 FROM sub_ref sr
WHERE sr.did = ?1 AND sr.feed_id = e.feed_id
)
ON CONFLICT (did, entry_id) DO UPDATE SET
read = excluded.read,
updated_at = excluded.updated_at
"#,
)
.bind(did)
.bind(entry_id)
.bind(read)
.bind(&now)
.execute(&mut *tx)
.await
.with_context(|| format!("mark_read failed for {did}/{entry_id}"))?;
if res.rows_affected() == 0 {
tx.rollback().await.ok();
return Ok(false);
}
project_entry_into_cursor(&mut tx, did, entry_id, read, &now).await?;
tx.commit().await.context("commit mark_read tx")?;
Ok(true)
}
pub async fn mark_starred(
pool: &SqlitePool,
did: &str,
entry_id: i64,
starred: bool,
) -> Result<bool> {
let res = sqlx::query(
r#"
INSERT INTO entry_state (did, entry_id, read, starred, updated_at)
SELECT ?1, e.id, 0, ?3, ?4
FROM entries e
WHERE e.id = ?2
AND EXISTS (
SELECT 1 FROM sub_ref sr
WHERE sr.did = ?1 AND sr.feed_id = e.feed_id
)
ON CONFLICT (did, entry_id) DO UPDATE SET
starred = excluded.starred,
updated_at = excluded.updated_at
"#,
)
.bind(did)
.bind(entry_id)
.bind(starred)
.bind(now_rfc3339())
.execute(pool)
.await
.with_context(|| format!("mark_starred failed for {did}/{entry_id}"))?;
Ok(res.rows_affected() > 0)
}
pub async fn compact_cursor(
pool: &SqlitePool,
did: &str,
feed_url: &str,
) -> Result<Option<String>> {
let mut tx = pool.begin().await.context("begin compact_cursor tx")?;
let (read_through, read_ids, unread_ids) = cursor_sets(&mut tx, did, feed_url).await?;
let oldest_unread: Option<String> = sqlx::query_scalar(
r#"
SELECT MIN(COALESCE(e.published, e.fetched_at))
FROM entries e
JOIN feeds f ON f.id = e.feed_id
LEFT JOIN entry_state s ON s.entry_id = e.id AND s.did = ?1
WHERE f.url = ?2 AND COALESCE(s.read, 0) = 0
"#,
)
.bind(did)
.bind(feed_url)
.fetch_one(&mut *tx)
.await
.with_context(|| format!("compact_cursor: oldest unread for {did}/{feed_url}"))?;
let watermark: Option<String> = match &oldest_unread {
Some(oldest) => sqlx::query_scalar(
r#"
SELECT MAX(COALESCE(e.published, e.fetched_at))
FROM entries e JOIN feeds f ON f.id = e.feed_id
WHERE f.url = ?1 AND COALESCE(e.published, e.fetched_at) < ?2
"#,
)
.bind(feed_url)
.bind(oldest)
.fetch_one(&mut *tx)
.await
.with_context(|| format!("compact_cursor: watermark for {did}/{feed_url}"))?,
None => sqlx::query_scalar(
r#"
SELECT MAX(COALESCE(e.published, e.fetched_at))
FROM entries e JOIN feeds f ON f.id = e.feed_id
WHERE f.url = ?1
"#,
)
.bind(feed_url)
.fetch_one(&mut *tx)
.await
.with_context(|| format!("compact_cursor: watermark for {did}/{feed_url}"))?,
};
let Some(watermark) = watermark else {
return Ok(None);
};
if read_through
.as_deref()
.is_some_and(|rt| rt >= &watermark[..])
{
return Ok(None);
}
let keep_above = ids_published_after(&mut tx, feed_url, &read_ids, &watermark).await?;
let keep_unread =
ids_published_at_or_before(&mut tx, feed_url, &unread_ids, &watermark).await?;
write_cursor_sets(
&mut tx,
did,
feed_url,
Some(&watermark),
&keep_above,
&keep_unread,
&now_rfc3339(),
)
.await?;
tx.commit().await.context("commit compact_cursor tx")?;
Ok(Some(watermark))
}
async fn ids_published_after(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
feed_url: &str,
ids: &str,
watermark: &str,
) -> Result<String> {
let live = ids_matching_watermark(tx, feed_url, watermark, true).await?;
Ok(filter_id_set_to_live(ids, &live))
}
async fn ids_published_at_or_before(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
feed_url: &str,
ids: &str,
watermark: &str,
) -> Result<String> {
let live = ids_matching_watermark(tx, feed_url, watermark, false).await?;
Ok(filter_id_set_to_live(ids, &live))
}
async fn ids_matching_watermark(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
feed_url: &str,
watermark: &str,
after: bool,
) -> Result<std::collections::HashSet<i64>> {
let sql = if after {
"SELECT e.id FROM entries e JOIN feeds f ON f.id = e.feed_id \
WHERE f.url = ?1 AND COALESCE(e.published, e.fetched_at) > ?2"
} else {
"SELECT e.id FROM entries e JOIN feeds f ON f.id = e.feed_id \
WHERE f.url = ?1 AND COALESCE(e.published, e.fetched_at) <= ?2"
};
Ok(sqlx::query_scalar::<_, i64>(sql)
.bind(feed_url)
.bind(watermark)
.fetch_all(&mut **tx)
.await
.context("compact_cursor: ids on one side of the watermark")?
.into_iter()
.collect())
}
pub async fn clear_star_by_identity(
pool: &SqlitePool,
did: &str,
url: Option<&str>,
guid: Option<&str>,
) -> Result<u64> {
if url.is_none_or(str::is_empty) && guid.is_none_or(str::is_empty) {
return Ok(0);
}
let res = sqlx::query(
r#"
UPDATE entry_state
SET starred = 0, updated_at = ?4
WHERE did = ?1
AND starred = 1
AND entry_id IN (
SELECT id FROM entries
WHERE (?2 IS NOT NULL AND url = ?2)
OR (?3 IS NOT NULL AND guid = ?3)
)
"#,
)
.bind(did)
.bind(url.filter(|u| !u.is_empty()))
.bind(guid.filter(|g| !g.is_empty()))
.bind(now_rfc3339())
.execute(pool)
.await
.with_context(|| format!("clear_star_by_identity failed for {did}"))?;
Ok(res.rows_affected())
}
pub async fn mark_feed_read(pool: &SqlitePool, did: &str, feed_id: i64, read: bool) -> Result<u64> {
let now = now_rfc3339();
let mut tx = pool.begin().await.context("begin mark_feed_read tx")?;
let res = sqlx::query(
r#"
INSERT INTO entry_state (did, entry_id, read, starred, updated_at)
SELECT ?1, e.id, ?2, 0, ?3 FROM entries e
WHERE e.feed_id = ?4
AND EXISTS (
SELECT 1 FROM sub_ref sr
WHERE sr.did = ?1 AND sr.feed_id = e.feed_id
)
ON CONFLICT (did, entry_id) DO UPDATE SET
read = excluded.read,
updated_at = excluded.updated_at
"#,
)
.bind(did)
.bind(read)
.bind(&now)
.bind(feed_id)
.execute(&mut *tx)
.await
.with_context(|| format!("mark_feed_read failed for {did}/feed {feed_id}"))?;
if res.rows_affected() > 0 {
project_feed_into_cursor(&mut tx, did, feed_id, read, &now).await?;
}
tx.commit().await.context("commit mark_feed_read tx")?;
Ok(res.rows_affected())
}
fn json_id_set_toggle(raw: &str, id: i64, present: bool) -> String {
let mut ids: Vec<i64> = serde_json::from_str::<Vec<serde_json::Value>>(raw)
.ok()
.map(|vals| {
vals.into_iter()
.filter_map(|v| match v {
serde_json::Value::Number(n) => n.as_i64(),
serde_json::Value::String(s) => s.parse::<i64>().ok(),
_ => None,
})
.collect()
})
.unwrap_or_default();
if present {
if !ids.contains(&id) {
ids.push(id);
}
} else {
ids.retain(|&x| x != id);
}
let as_strings: Vec<String> = ids.iter().map(|i| i.to_string()).collect();
serde_json::to_string(&as_strings).unwrap_or_else(|_| "[]".to_string())
}
async fn feed_url_for_id_tx(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
feed_id: i64,
) -> Result<Option<String>> {
let url: Option<String> = sqlx::query_scalar("SELECT url FROM feeds WHERE id = ?1")
.bind(feed_id)
.fetch_optional(&mut **tx)
.await
.with_context(|| format!("feed_url_for_id_tx failed for feed {feed_id}"))?;
Ok(url)
}
async fn cursor_sets(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
did: &str,
feed_url: &str,
) -> Result<(Option<String>, String, String)> {
let row = sqlx::query(
"SELECT read_through, read_ids, unread_ids FROM read_cursor \
WHERE did = ?1 AND feed_url = ?2",
)
.bind(did)
.bind(feed_url)
.fetch_optional(&mut **tx)
.await
.with_context(|| format!("cursor_sets failed for {did}/{feed_url}"))?;
Ok(match row {
Some(r) => (
r.get::<Option<String>, _>("read_through"),
r.get::<String, _>("read_ids"),
r.get::<String, _>("unread_ids"),
),
None => (None, "[]".to_string(), "[]".to_string()),
})
}
async fn write_cursor_sets(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
did: &str,
feed_url: &str,
read_through: Option<&str>,
read_ids: &str,
unread_ids: &str,
now: &str,
) -> Result<()> {
sqlx::query(
r#"
INSERT INTO read_cursor
(did, feed_url, read_through, read_ids, unread_ids, dirty, updated_at)
VALUES (?1, ?2, ?3, ?4, ?5, 1, ?6)
ON CONFLICT (did, feed_url) DO UPDATE SET
read_through = excluded.read_through,
read_ids = excluded.read_ids,
unread_ids = excluded.unread_ids,
dirty = 1,
updated_at = excluded.updated_at
"#,
)
.bind(did)
.bind(feed_url)
.bind(read_through)
.bind(read_ids)
.bind(unread_ids)
.bind(now)
.execute(&mut **tx)
.await
.with_context(|| format!("write_cursor_sets failed for {did}/{feed_url}"))?;
Ok(())
}
async fn project_entry_into_cursor(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
did: &str,
entry_id: i64,
read: bool,
now: &str,
) -> Result<()> {
let feed_id: Option<i64> = sqlx::query_scalar("SELECT feed_id FROM entries WHERE id = ?1")
.bind(entry_id)
.fetch_optional(&mut **tx)
.await
.with_context(|| format!("project_entry_into_cursor: feed_id for entry {entry_id}"))?;
let feed_id = match feed_id {
Some(f) => f,
None => return Ok(()), };
let feed_url = match feed_url_for_id_tx(tx, feed_id).await? {
Some(u) => u,
None => return Ok(()),
};
let (read_through, read_ids, unread_ids) = cursor_sets(tx, did, &feed_url).await?;
let read_ids = json_id_set_toggle(&read_ids, entry_id, read);
let unread_ids = json_id_set_toggle(&unread_ids, entry_id, !read);
write_cursor_sets(
tx,
did,
&feed_url,
read_through.as_deref(),
&read_ids,
&unread_ids,
now,
)
.await
}
async fn project_feed_into_cursor(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
did: &str,
feed_id: i64,
read: bool,
now: &str,
) -> Result<()> {
let feed_url = match feed_url_for_id_tx(tx, feed_id).await? {
Some(u) => u,
None => return Ok(()),
};
let ids: Vec<i64> = sqlx::query_scalar(
r#"
SELECT e.id FROM entries e
WHERE e.feed_id = ?2
AND EXISTS (
SELECT 1 FROM sub_ref sr
WHERE sr.did = ?1 AND sr.feed_id = e.feed_id
)
"#,
)
.bind(did)
.bind(feed_id)
.fetch_all(&mut **tx)
.await
.with_context(|| format!("project_feed_into_cursor: entry ids for {did}/feed {feed_id}"))?;
let (read_through, mut read_ids, mut unread_ids) = cursor_sets(tx, did, &feed_url).await?;
for id in ids {
read_ids = json_id_set_toggle(&read_ids, id, read);
unread_ids = json_id_set_toggle(&unread_ids, id, !read);
}
write_cursor_sets(
tx,
did,
&feed_url,
read_through.as_deref(),
&read_ids,
&unread_ids,
now,
)
.await
}
#[cfg(test)]
mod test_helpers {
use super::*;
const FIXTURE_MAX: i64 = 10_000;
pub(crate) async fn entries_for_feed(
pool: &SqlitePool,
did: &str,
feed_id: i64,
) -> Result<Vec<EntryListRow>> {
list_entries(pool, did, ListView::All, Some(&[feed_id]), FIXTURE_MAX, 0).await
}
pub(crate) async fn get_unread_for_did(
pool: &SqlitePool,
did: &str,
) -> Result<Vec<EntryListRow>> {
list_entries(pool, did, ListView::Unread, None, FIXTURE_MAX, 0).await
}
pub(crate) async fn get_starred_for_did(
pool: &SqlitePool,
did: &str,
) -> Result<Vec<EntryListRow>> {
list_entries(pool, did, ListView::Starred, None, FIXTURE_MAX, 0).await
}
}
#[cfg(test)]
pub(crate) use test_helpers::{entries_for_feed, get_starred_for_did, get_unread_for_did};
pub async fn upsert_cursor(pool: &SqlitePool, cursor: &ReadCursor) -> Result<()> {
sqlx::query(
r#"
INSERT INTO read_cursor
(did, feed_url, read_through, read_ids, unread_ids, dirty, updated_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
ON CONFLICT (did, feed_url) DO UPDATE SET
read_through = excluded.read_through,
read_ids = excluded.read_ids,
unread_ids = excluded.unread_ids,
dirty = excluded.dirty,
updated_at = excluded.updated_at
"#,
)
.bind(&cursor.did)
.bind(&cursor.feed_url)
.bind(&cursor.read_through)
.bind(&cursor.read_ids)
.bind(&cursor.unread_ids)
.bind(cursor.dirty)
.bind(&cursor.updated_at)
.execute(pool)
.await
.with_context(|| {
format!(
"upsert_cursor failed for {}/{}",
cursor.did, cursor.feed_url
)
})?;
Ok(())
}
pub async fn get_cursor(
pool: &SqlitePool,
did: &str,
feed_url: &str,
) -> Result<Option<ReadCursor>> {
let cursor = sqlx::query_as::<_, ReadCursor>(
"SELECT * FROM read_cursor WHERE did = ?1 AND feed_url = ?2",
)
.bind(did)
.bind(feed_url)
.fetch_optional(pool)
.await
.context("get_cursor failed")?;
Ok(cursor)
}
pub async fn parked_readstate_dids(pool: &SqlitePool) -> Result<i64> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(DISTINCT rc.did)
FROM read_cursor rc
WHERE rc.dirty = 1
AND NOT EXISTS (SELECT 1 FROM oauth_session s WHERE s.sub = rc.did)
"#,
)
.fetch_one(pool)
.await
.context("counting parked read-state DIDs")?;
Ok(row.0)
}
pub async fn dirty_cursors(pool: &SqlitePool, did: &str) -> Result<Vec<ReadCursor>> {
let cursors =
sqlx::query_as::<_, ReadCursor>("SELECT * FROM read_cursor WHERE did = ?1 AND dirty = 1")
.bind(did)
.fetch_all(pool)
.await
.with_context(|| format!("dirty_cursors failed for {did}"))?;
Ok(cursors)
}
pub async fn record_network_stat(pool: &SqlitePool, stat: &NetworkStat) -> Result<()> {
sqlx::query(
r#"
INSERT INTO network_stat (key, source, value, truncated, observed_at)
VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT (key, source) DO UPDATE SET
value = excluded.value,
truncated = excluded.truncated,
observed_at = excluded.observed_at
WHERE NOT (excluded.truncated = 1 AND excluded.value <= network_stat.value)
"#,
)
.bind(&stat.key)
.bind(&stat.source)
.bind(stat.value)
.bind(stat.truncated)
.bind(&stat.observed_at)
.execute(pool)
.await
.with_context(|| {
format!(
"record_network_stat failed for {}/{}",
stat.key, stat.source
)
})?;
Ok(())
}
pub async fn latest_network_stat(pool: &SqlitePool, key: &str) -> Result<Option<NetworkStat>> {
let stat = sqlx::query_as::<_, NetworkStat>(
"SELECT key, source, value, truncated, observed_at FROM network_stat \
WHERE key = ?1 ORDER BY value DESC, observed_at DESC LIMIT 1",
)
.bind(key)
.fetch_optional(pool)
.await
.with_context(|| format!("latest_network_stat failed for {key}"))?;
Ok(stat)
}
pub async fn mark_cursor_pds_created(pool: &SqlitePool, did: &str, feed_url: &str) -> Result<()> {
sqlx::query("UPDATE read_cursor SET pds_created = 1 WHERE did = ?1 AND feed_url = ?2")
.bind(did)
.bind(feed_url)
.execute(pool)
.await
.with_context(|| format!("mark_cursor_pds_created failed for {did}/{feed_url}"))?;
Ok(())
}
pub async fn set_cursor_pds_created(
pool: &SqlitePool,
did: &str,
feed_url: &str,
created: bool,
) -> Result<()> {
sqlx::query("UPDATE read_cursor SET pds_created = ?3 WHERE did = ?1 AND feed_url = ?2")
.bind(did)
.bind(feed_url)
.bind(created)
.execute(pool)
.await
.with_context(|| format!("set_cursor_pds_created failed for {did}/{feed_url}"))?;
Ok(())
}
pub async fn clear_cursor_dirty(
pool: &SqlitePool,
did: &str,
feed_url: &str,
flushed_updated_at: &str,
) -> Result<()> {
sqlx::query(
"UPDATE read_cursor SET dirty = 0 \
WHERE did = ?1 AND feed_url = ?2 AND updated_at = ?3",
)
.bind(did)
.bind(feed_url)
.bind(flushed_updated_at)
.execute(pool)
.await
.context("clear_cursor_dirty failed")?;
Ok(())
}
pub(crate) fn now_unix() -> i64 {
chrono::Utc::now().timestamp()
}
const CODE_ALPHABET: &[u8] = b"ABCDEFGHJKLMNPQRSTUVWXYZ23456789";
const CODE_PREFIX: &str = "FEATHER-";
const CODE_BODY_LEN: usize = 8;
pub fn generate_invite_code() -> Result<String> {
let n = CODE_ALPHABET.len() as u16; let limit = 256 / n * n; let mut out = String::with_capacity(CODE_PREFIX.len() + CODE_BODY_LEN);
out.push_str(CODE_PREFIX);
let mut got = 0;
let mut buf = [0u8; 1];
while got < CODE_BODY_LEN {
getrandom::fill(&mut buf).context("getrandom failed while minting invite code")?;
let b = buf[0] as u16;
if b < limit {
out.push(CODE_ALPHABET[(b % n) as usize] as char);
got += 1;
}
}
Ok(out)
}
pub async fn has_beta_access(pool: &SqlitePool, did: &str) -> Result<bool> {
let row = sqlx::query("SELECT 1 FROM beta_access WHERE did = ?1")
.bind(did)
.fetch_optional(pool)
.await
.with_context(|| format!("has_beta_access failed for {did}"))?;
Ok(row.is_some())
}
pub async fn count_beta_access(pool: &SqlitePool) -> Result<i64> {
let row = sqlx::query("SELECT COUNT(*) AS n FROM beta_access")
.fetch_one(pool)
.await
.context("count_beta_access failed")?;
Ok(row.get::<i64, _>("n"))
}
pub async fn count_active_codes(pool: &SqlitePool) -> Result<i64> {
let now = now_unix();
let row = sqlx::query(
"SELECT COUNT(*) AS n FROM invite_codes WHERE status = 'active' AND expires_at >= ?1",
)
.bind(now)
.fetch_one(pool)
.await
.context("count_active_codes failed")?;
Ok(row.get::<i64, _>("n"))
}
pub async fn grant_access(
pool: &SqlitePool,
did: &str,
handle: Option<&str>,
granted_by: &str,
invite_code_used: Option<&str>,
) -> Result<()> {
sqlx::query(
r#"
INSERT INTO beta_access (did, handle, granted_by, granted_at, invite_code_used)
VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT (did) DO UPDATE SET
handle = COALESCE(excluded.handle, beta_access.handle),
granted_by = excluded.granted_by,
invite_code_used = COALESCE(excluded.invite_code_used, beta_access.invite_code_used)
"#,
)
.bind(did)
.bind(handle)
.bind(granted_by)
.bind(now_unix())
.bind(invite_code_used)
.execute(pool)
.await
.with_context(|| format!("grant_access failed for {did}"))?;
Ok(())
}
pub async fn mint_code(pool: &SqlitePool, creator_did: &str, ttl_secs: i64) -> Result<String> {
mint_code_inner(pool, creator_did, ttl_secs, None).await
}
pub async fn mint_code_for_did(
pool: &SqlitePool,
creator_did: &str,
ttl_secs: i64,
intended_did: &str,
) -> Result<String> {
mint_code_inner(pool, creator_did, ttl_secs, Some(intended_did)).await
}
async fn mint_code_inner(
pool: &SqlitePool,
creator_did: &str,
ttl_secs: i64,
intended_did: Option<&str>,
) -> Result<String> {
let code = generate_invite_code()?;
let now = now_unix();
let expires_at = now.saturating_add(ttl_secs.max(0));
sqlx::query(
r#"
INSERT INTO invite_codes
(code, creator_did, status, invitee_did, intended_did, created_at, expires_at, redeemed_at)
VALUES (?1, ?2, 'active', NULL, ?3, ?4, ?5, NULL)
"#,
)
.bind(&code)
.bind(creator_did)
.bind(intended_did)
.bind(now)
.bind(expires_at)
.execute(pool)
.await
.with_context(|| format!("mint_code failed for creator {creator_did}"))?;
Ok(code)
}
pub fn is_intended_active_conflict(err: &anyhow::Error) -> bool {
for cause in err.chain() {
if let Some(sqlx::Error::Database(db)) = cause.downcast_ref::<sqlx::Error>() {
let msg = db.message();
let is_unique = db.code().as_deref() == Some("2067")
|| db.code().as_deref() == Some("19")
|| msg.contains("UNIQUE constraint failed");
if is_unique && msg.contains("invite_codes.intended_did") {
return true;
}
}
}
false
}
pub async fn find_active_code_for_did(
pool: &SqlitePool,
intended_did: &str,
) -> Result<Option<String>> {
let now = now_unix();
let row = sqlx::query(
"SELECT code FROM invite_codes
WHERE intended_did = ?1 AND status = 'active' AND expires_at >= ?2
ORDER BY expires_at ASC
LIMIT 1",
)
.bind(intended_did)
.bind(now)
.fetch_optional(pool)
.await
.with_context(|| format!("find_active_code_for_did failed for {intended_did}"))?;
Ok(row.map(|r| r.get::<String, _>("code")))
}
pub async fn redeem_code(
pool: &SqlitePool,
code: &str,
did: &str,
handle: Option<&str>,
cap: i64,
) -> Result<std::result::Result<(), RedeemError>> {
let now = now_unix();
let mut tx = pool.begin().await.context("begin redeem_code tx")?;
sqlx::query("UPDATE invite_codes SET status = status WHERE code = ?1")
.bind(code)
.execute(&mut *tx)
.await
.context("redeem_code: acquire write lock")?;
let row =
sqlx::query("SELECT status, expires_at, intended_did FROM invite_codes WHERE code = ?1")
.bind(code)
.fetch_optional(&mut *tx)
.await
.context("redeem_code: lookup")?;
let row = match row {
Some(r) => r,
None => return Ok(Err(RedeemError::NotFound)),
};
let status: String = row.get("status");
let expires_at: i64 = row.get("expires_at");
let intended_did: Option<String> = row.get("intended_did");
if let Some(bound) = intended_did.as_deref() {
if bound != did {
return Ok(Err(RedeemError::NotFound));
}
}
if status == "expired" || now > expires_at {
return Ok(Err(RedeemError::Expired));
}
if status != "active" {
return Ok(Err(RedeemError::AlreadyRedeemed));
}
let count: i64 = sqlx::query("SELECT COUNT(*) AS n FROM beta_access")
.fetch_one(&mut *tx)
.await
.context("redeem_code: count")?
.get("n");
if count >= cap {
return Ok(Err(RedeemError::CapacityFull));
}
let flipped = sqlx::query(
r#"
UPDATE invite_codes
SET status = 'redeemed', invitee_did = ?2, redeemed_at = ?3
WHERE code = ?1 AND status = 'active'
"#,
)
.bind(code)
.bind(did)
.bind(now)
.execute(&mut *tx)
.await
.context("redeem_code: flip")?;
if flipped.rows_affected() == 0 {
return Ok(Err(RedeemError::AlreadyRedeemed));
}
sqlx::query(
r#"
INSERT INTO beta_access (did, handle, granted_by, granted_at, invite_code_used)
VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT (did) DO UPDATE SET
handle = COALESCE(excluded.handle, beta_access.handle),
invite_code_used = excluded.invite_code_used
"#,
)
.bind(did)
.bind(handle)
.bind(
sqlx::query("SELECT creator_did FROM invite_codes WHERE code = ?1")
.bind(code)
.fetch_one(&mut *tx)
.await
.context("redeem_code: creator lookup")?
.get::<String, _>("creator_did"),
)
.bind(now)
.bind(code)
.execute(&mut *tx)
.await
.context("redeem_code: grant")?;
tx.commit().await.context("commit redeem_code tx")?;
Ok(Ok(()))
}
pub async fn expire_old_codes(pool: &SqlitePool) -> Result<u64> {
let now = now_unix();
let res = sqlx::query(
"UPDATE invite_codes SET status = 'expired' WHERE status = 'active' AND expires_at < ?1",
)
.bind(now)
.execute(pool)
.await
.context("expire_old_codes failed")?;
Ok(res.rows_affected())
}
pub async fn ensure_seed(pool: &SqlitePool, dids: &[String]) -> Result<u64> {
let mut tx = pool.begin().await.context("begin ensure_seed tx")?;
let now = now_unix();
let mut created = 0u64;
for did in dids {
let res = sqlx::query(
r#"
INSERT INTO beta_access (did, handle, granted_by, granted_at, invite_code_used)
VALUES (?1, NULL, 'admin', ?2, NULL)
ON CONFLICT (did) DO NOTHING
"#,
)
.bind(did)
.bind(now)
.execute(&mut *tx)
.await
.with_context(|| format!("ensure_seed insert failed for {did}"))?;
created += res.rows_affected();
}
tx.commit().await.context("commit ensure_seed tx")?;
Ok(created)
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct PurgeCounts {
pub entry_state: u64,
pub read_cursor: u64,
pub sub_ref: u64,
pub beta_access: u64,
pub invite_codes: u64,
pub invitee_scrubbed: u64,
pub granted_by_scrubbed: u64,
}
impl PurgeCounts {
pub fn total(&self) -> u64 {
self.entry_state + self.read_cursor + self.sub_ref + self.beta_access + self.invite_codes
}
}
pub const REDACTED_DID: &str = "__redacted__";
pub async fn purge_did_data(pool: &SqlitePool, did: &str) -> Result<PurgeCounts> {
let mut tx = pool.begin().await.context("begin purge_did_data tx")?;
let entry_state = sqlx::query("DELETE FROM entry_state WHERE did = ?1")
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("purge entry_state for {did}"))?
.rows_affected();
let read_cursor = sqlx::query("DELETE FROM read_cursor WHERE did = ?1")
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("purge read_cursor for {did}"))?
.rows_affected();
let sub_ref = sqlx::query("DELETE FROM sub_ref WHERE did = ?1")
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("purge sub_ref for {did}"))?
.rows_affected();
let beta_access = sqlx::query("DELETE FROM beta_access WHERE did = ?1")
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("purge beta_access for {did}"))?
.rows_affected();
let invite_codes = sqlx::query("DELETE FROM invite_codes WHERE creator_did = ?1")
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("purge invite_codes for {did}"))?
.rows_affected();
let invitee_scrubbed =
sqlx::query("UPDATE invite_codes SET invitee_did = NULL WHERE invitee_did = ?1")
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("scrub invitee_did for {did}"))?
.rows_affected();
sqlx::query(
"UPDATE invite_codes \
SET intended_did = NULL, \
status = CASE WHEN status = 'active' THEN 'expired' ELSE status END \
WHERE intended_did = ?1",
)
.bind(did)
.execute(&mut *tx)
.await
.with_context(|| format!("scrub intended_did for {did}"))?;
let granted_by_scrubbed =
sqlx::query("UPDATE beta_access SET granted_by = ?2 WHERE granted_by = ?1")
.bind(did)
.bind(REDACTED_DID)
.execute(&mut *tx)
.await
.with_context(|| format!("scrub granted_by for {did}"))?
.rows_affected();
tx.commit().await.context("commit purge_did_data tx")?;
Ok(PurgeCounts {
entry_state,
read_cursor,
sub_ref,
beta_access,
invite_codes,
invitee_scrubbed,
granted_by_scrubbed,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PollHealth {
pub feeds_tracked: i64,
pub polled_last_hour: i64,
pub overdue: i64,
pub last_poll_secs_ago: Option<i64>,
pub oldest_poll_secs_ago: Option<i64>,
pub never_polled: i64,
pub in_backoff: i64,
pub badly_broken: i64,
pub failure_kinds: Vec<(String, i64)>,
}
const BADLY_BROKEN_ERRORS: i64 = 6;
pub async fn poll_health(pool: &SqlitePool, now: &str, hour_ago: &str) -> Result<PollHealth> {
let aggregate = format!(
r#"
SELECT
COUNT(*),
COALESCE(SUM(CASE WHEN last_polled IS NOT NULL AND last_polled >= ?2 THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE WHEN next_poll IS NULL OR next_poll <= ?1 THEN 1 ELSE 0 END), 0),
MAX(last_polled),
-- NULL-AWARE. `MIN` skips NULLs, so an instance where most feeds
-- had NEVER been polled reported the freshest of the few that had —
-- the figure read healthiest in the most degraded state, which is
-- the opposite of what a health page is for. A never-polled feed IS
-- the worst staleness, so it wins outright.
CASE WHEN SUM(CASE WHEN last_polled IS NULL THEN 1 ELSE 0 END) > 0
THEN NULL ELSE MIN(last_polled) END,
SUM(CASE WHEN last_polled IS NULL THEN 1 ELSE 0 END),
COALESCE(SUM(CASE WHEN consecutive_errors > 0 THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE WHEN consecutive_errors >= ?3 THEN 1 ELSE 0 END), 0)
FROM feeds
WHERE kind IN ({POLLABLE_KINDS_SQL})
"#
);
#[allow(clippy::type_complexity)]
let row: (i64, i64, i64, Option<String>, Option<String>, i64, i64, i64) =
sqlx::query_as(sqlx::AssertSqlSafe(aggregate))
.bind(now)
.bind(hour_ago)
.bind(BADLY_BROKEN_ERRORS)
.fetch_one(pool)
.await
.context("computing poll health")?;
let histogram = format!(
r#"
-- **`failure_kind`, not `kind`.** Aliasing this `kind` collided with
-- the `feeds.kind` column added for the poller: SQLite resolved
-- `GROUP BY kind` to the table column, so every failing feed collapsed
-- into ONE bucket labelled from an arbitrary row — a public page
-- reporting "10 fetch" for ten unrelated causes. Caught by
-- `an_unrecognised_failure_kind_folds_into_unknown`.
SELECT COALESCE(last_error_kind, 'unknown') AS failure_kind, COUNT(*) AS n
FROM feeds
WHERE consecutive_errors > 0 AND kind IN ({POLLABLE_KINDS_SQL})
GROUP BY failure_kind
ORDER BY n DESC, failure_kind ASC
"#
);
let kinds: Vec<(String, i64)> = sqlx::query_as(sqlx::AssertSqlSafe(histogram))
.fetch_all(pool)
.await
.context("computing the failure-cause histogram")?;
let mut folded: std::collections::BTreeMap<String, i64> = std::collections::BTreeMap::new();
for (kind, n) in kinds {
let key = if kind == "unknown" || crate::feed::FailureKind::parse(&kind).is_some() {
kind
} else {
"unknown".to_string()
};
*folded.entry(key).or_insert(0) += n;
}
let mut kinds: Vec<(String, i64)> = folded.into_iter().collect();
kinds.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
Ok(PollHealth {
feeds_tracked: row.0,
polled_last_hour: row.1,
overdue: row.2,
last_poll_secs_ago: secs_between(row.3.as_deref(), now),
oldest_poll_secs_ago: secs_between(row.4.as_deref(), now),
never_polled: row.5,
in_backoff: row.6,
badly_broken: row.7,
failure_kinds: kinds,
})
}
fn secs_between(then: Option<&str>, now: &str) -> Option<i64> {
let then = chrono::DateTime::parse_from_rfc3339(then?).ok()?;
let now = chrono::DateTime::parse_from_rfc3339(now).ok()?;
Some((now - then).num_seconds().max(0))
}
#[cfg(test)]
mod tests {
use super::*;
const PUB_A: &str = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3laa";
#[tokio::test]
async fn a_due_publication_is_handed_to_the_poller() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
upsert_feed(
&pool,
&NewFeed {
url: PUB_A.into(),
..Default::default()
},
)
.await?;
let due = due_feeds(&pool, "2999-01-01T00:00:00Z", 50).await?;
assert!(
due.iter().any(|f| f.url == PUB_A),
"a publication row is not handed to the poller"
);
Ok(())
}
#[tokio::test]
async fn admitting_publications_staggers_their_first_poll() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
for i in 0..19 {
let url =
format!("at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3l{i:02}");
upsert_feed(
&pool,
&NewFeed {
url,
..Default::default()
},
)
.await?;
}
upsert_feed(
&pool,
&NewFeed {
url: "https://rss.example/feed.xml".into(),
..Default::default()
},
)
.await?;
let n = stagger_unscheduled(
&pool,
crate::feed::FeedKind::Publication,
std::time::Duration::from_secs(3600),
)
.await?;
assert_eq!(n, 19, "not every unscheduled publication was scheduled");
let slots: Vec<String> = sqlx::query_scalar(
"SELECT next_poll FROM feeds WHERE kind = 'publication' ORDER BY next_poll",
)
.fetch_all(&pool)
.await?;
let distinct: std::collections::BTreeSet<_> = slots.iter().collect();
assert_eq!(distinct.len(), 19, "publications share slots: {slots:?}");
let rss: Option<String> = sqlx::query_scalar(
"SELECT next_poll FROM feeds WHERE url = 'https://rss.example/feed.xml'",
)
.fetch_one(&pool)
.await?;
assert_eq!(rss, None, "an RSS row was rescheduled");
let again = stagger_unscheduled(
&pool,
crate::feed::FeedKind::Publication,
std::time::Duration::from_secs(3600),
)
.await?;
assert_eq!(
again, 0,
"a second boot re-staggered rows that already had a slot"
);
Ok(())
}
#[tokio::test]
async fn admitted_rows_do_not_outrank_an_overdue_rss_feed() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
for i in 0..5 {
let url =
format!("at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3l{i:02}");
upsert_feed(
&pool,
&NewFeed {
url,
..Default::default()
},
)
.await?;
}
upsert_feed(
&pool,
&NewFeed {
url: "https://overdue.example/feed.xml".into(),
next_poll: Some("2000-01-01T00:00:00Z".into()),
..Default::default()
},
)
.await?;
stagger_unscheduled(
&pool,
crate::feed::FeedKind::Publication,
std::time::Duration::from_secs(3600),
)
.await?;
let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let first = due_feeds(&pool, &now, 1).await?;
assert_eq!(first[0].url, "https://overdue.example/feed.xml");
Ok(())
}
#[tokio::test]
async fn a_repoll_refreshes_published_but_never_fetched_at() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/f.xml".to_string(),
..Default::default()
},
)
.await?;
let seen = |at: &str| {
vec![NewEntry {
guid: "g".to_string(),
published: Some(at.to_string()),
fetched_at: Some(at.to_string()),
..Default::default()
}]
};
insert_entries(&pool, feed_id, &seen("2026-01-01T00:00:00Z"), 0).await?;
insert_entries(&pool, feed_id, &seen("2026-09-20T00:00:00Z"), 0).await?;
let (published, fetched_at): (Option<String>, String) =
sqlx::query_as("SELECT published, fetched_at FROM entries WHERE guid = 'g'")
.fetch_one(&pool)
.await?;
assert_eq!(
published.as_deref(),
Some("2026-09-20T00:00:00Z"),
"the second poll's date did not replace the first"
);
assert_eq!(
fetched_at, "2026-01-01T00:00:00Z",
"fetched_at moved, so an undated entry would never age either"
);
Ok(())
}
#[tokio::test]
async fn validators_survive_a_partial_upsert() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let url = "https://example.com/feed.xml";
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
etag: Some("\"abc123\"".to_string()),
last_modified: Some("Wed, 01 Jan 2026 00:00:00 GMT".to_string()),
..Default::default()
},
)
.await?;
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
next_poll: Some("2026-07-12T00:00:00Z".to_string()),
..Default::default()
},
)
.await?;
let feed = get_feed_by_url(&pool, url).await?.expect("feed");
assert_eq!(
feed.etag.as_deref(),
Some("\"abc123\""),
"a partial upsert erased the ETag, disabling conditional GET"
);
assert_eq!(
feed.last_modified.as_deref(),
Some("Wed, 01 Jan 2026 00:00:00 GMT"),
"a partial upsert erased Last-Modified"
);
assert_eq!(feed.next_poll.as_deref(), Some("2026-07-12T00:00:00Z"));
Ok(())
}
#[tokio::test]
async fn a_ceiling_inside_the_window_is_ignored_not_applied() -> Result<()> {
for hard in [0_i64, 1, 7, 14] {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/f.xml".to_string(),
..Default::default()
},
)
.await?;
let old = (chrono::Utc::now() - chrono::Duration::days(30))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "starred-30d".to_string(),
published: Some(old.clone()),
..Default::default()
},
NewEntry {
guid: "unread-30d".to_string(),
published: Some(old.clone()),
..Default::default()
},
],
0,
)
.await?;
sqlx::query(
"INSERT INTO entry_state (did, entry_id, read, starred, updated_at)
SELECT 'did:plc:x', id, 1, 1, '2026-01-01T00:00:00Z'
FROM entries WHERE guid = 'starred-30d'",
)
.execute(&pool)
.await?;
sqlx::query(
"INSERT INTO entry_state (did, entry_id, read, starred, updated_at)
SELECT 'did:plc:x', id, 0, 0, '2026-01-01T00:00:00Z'
FROM entries WHERE guid = 'unread-30d'",
)
.execute(&pool)
.await?;
prune_old_entries(&pool, 14, hard, 0).await?;
let left: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await?;
assert_eq!(
left, 2,
"hard_days={hard} destroyed starred/unread rows at the soft window"
);
}
Ok(())
}
#[tokio::test]
async fn a_disabled_window_does_not_disable_the_ceiling() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/f.xml".to_string(),
..Default::default()
},
)
.await?;
let age = |d: i64| {
(chrono::Utc::now() - chrono::Duration::days(d))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
};
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "starred-400d".to_string(),
published: Some(age(400)),
..Default::default()
},
NewEntry {
guid: "starred-30d".to_string(),
published: Some(age(30)),
..Default::default()
},
],
0,
)
.await?;
sqlx::query(
"INSERT INTO entry_state (did, entry_id, read, starred, updated_at)
SELECT 'did:plc:x', id, 1, 1, '2026-01-01T00:00:00Z' FROM entries",
)
.execute(&pool)
.await?;
let deleted = prune_old_entries(&pool, 0, 180, 0).await?;
assert_eq!(
deleted, 1,
"retention_days=0 skipped the hard ceiling, leaving the cache unbounded"
);
let left: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
.fetch_all(&pool)
.await?;
assert_eq!(
left,
vec!["starred-30d".to_string()],
"the ceiling removed the wrong rows with the window disabled"
);
Ok(())
}
#[tokio::test]
async fn both_knobs_off_deletes_nothing() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/f.xml".to_string(),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: "ancient".to_string(),
published: Some(
(chrono::Utc::now() - chrono::Duration::days(9999))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
),
..Default::default()
}],
0,
)
.await?;
assert_eq!(prune_old_entries(&pool, 0, 0, 0).await?, 0);
let left: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await?;
assert_eq!(left, 1);
Ok(())
}
#[tokio::test]
async fn per_feed_trim_stays_bounded_when_everything_is_starred() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/f.xml".to_string(),
..Default::default()
},
)
.await?;
let entries: Vec<NewEntry> = (0..100)
.map(|i| NewEntry {
guid: format!("g-{i}"),
published: Some(format!("2026-01-{:02}T00:00:00Z", (i % 28) + 1)),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
sqlx::query(
"INSERT INTO entry_state (did, entry_id, read, starred, updated_at)
SELECT 'did:plc:x', id, 0, 1, '2026-01-01T00:00:00Z'
FROM entries LIMIT 50",
)
.execute(&pool)
.await?;
insert_entries(&pool, feed_id, &[], 5).await?;
let left: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await?;
assert!(
left <= 10,
"per-feed trim kept {left} rows for a cap of 5; sparing removed the bound"
);
Ok(())
}
#[tokio::test]
async fn init_insert_readback() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/feed.xml".to_string(),
title: Some("Example".to_string()),
site_url: Some("https://example.com".to_string()),
next_poll: Some("2026-07-12T00:00:00Z".to_string()),
..Default::default()
},
)
.await?;
assert!(feed_id > 0);
let feed = get_feed_by_url(&pool, "https://example.com/feed.xml")
.await?
.expect("feed should exist");
assert_eq!(feed.id, feed_id);
assert_eq!(feed.title.as_deref(), Some("Example"));
assert_eq!(feed.site_url.as_deref(), Some("https://example.com"));
let feed_id2 = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/feed.xml".to_string(),
title: Some("Example (renamed)".to_string()),
..Default::default()
},
)
.await?;
assert_eq!(feed_id, feed_id2, "same URL must reuse the same row");
let n = insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "guid-1".to_string(),
url: Some("https://example.com/a".to_string()),
title: Some("First".to_string()),
published: Some("2026-07-10T08:00:00Z".to_string()),
content_html: Some("<p>hello</p>".to_string()),
..Default::default()
},
NewEntry {
guid: "guid-2".to_string(),
url: Some("https://example.com/b".to_string()),
title: Some("Second".to_string()),
published: Some("2026-07-11T08:00:00Z".to_string()),
..Default::default()
},
],
0, )
.await?;
assert_eq!(n, 2);
let did = "did:plc:abc123";
replace_sub_refs(&pool, did, &[feed_id]).await?;
let entries = entries_for_feed(&pool, did, feed_id).await?;
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].guid, "guid-2");
assert_eq!(entries[1].guid, "guid-1");
let body: Option<String> =
sqlx::query_scalar("SELECT content_html FROM entries WHERE guid = 'guid-1'")
.fetch_one(&pool)
.await?;
assert_eq!(body.as_deref(), Some("<p>hello</p>"));
let n2 = insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: "guid-1".to_string(),
title: Some("First (edited)".to_string()),
..Default::default()
}],
0,
)
.await?;
assert_eq!(n2, 1);
assert_eq!(entries_for_feed(&pool, did, feed_id).await?.len(), 2);
let e1 = entries.iter().find(|e| e.guid == "guid-1").unwrap().id;
assert_eq!(get_unread_for_did(&pool, did).await?.len(), 2);
mark_read(&pool, did, e1, true).await?;
let unread = get_unread_for_did(&pool, did).await?;
assert_eq!(unread.len(), 1);
assert_eq!(unread[0].guid, "guid-2");
mark_starred(&pool, did, e1, true).await?;
let starred = get_starred_for_did(&pool, did).await?;
assert_eq!(starred.len(), 1);
assert_eq!(starred[0].id, e1);
mark_feed_read(&pool, did, feed_id, true).await?;
assert_eq!(get_unread_for_did(&pool, did).await?.len(), 0);
let cursor = ReadCursor {
did: did.to_string(),
feed_url: "https://example.com/feed.xml".to_string(),
read_through: Some("2026-07-11T08:00:00Z".to_string()),
read_ids: "[]".to_string(),
unread_ids: "[]".to_string(),
dirty: true,
pds_created: false,
updated_at: now_rfc3339(),
};
upsert_cursor(&pool, &cursor).await?;
let fetched = get_cursor(&pool, did, "https://example.com/feed.xml")
.await?
.expect("cursor should exist");
assert_eq!(
fetched.read_through.as_deref(),
Some("2026-07-11T08:00:00Z")
);
assert!(fetched.dirty);
let dirty = dirty_cursors(&pool, did).await?;
assert_eq!(dirty.len(), 1);
let flushed_at = dirty[0].updated_at.clone();
clear_cursor_dirty(&pool, did, "https://example.com/feed.xml", &flushed_at).await?;
assert_eq!(dirty_cursors(&pool, did).await?.len(), 0);
Ok(())
}
async fn seed_big_entries(pool: &SqlitePool, did: &str, count: usize) -> Result<i64> {
let feed_id = upsert_feed(
pool,
&NewFeed {
url: "https://example.com/big.xml".to_string(),
..Default::default()
},
)
.await?;
let body = "x".repeat(20_000);
let entries: Vec<NewEntry> = (0..count)
.map(|i| NewEntry {
guid: format!("guid-{i:04}"),
url: Some(format!("https://example.com/a/{i}")),
title: Some(format!("Article {i}")),
published: Some(format!("2026-01-{:02}T00:00:00Z", (i % 28) + 1)),
content_html: Some(body.clone()),
..Default::default()
})
.collect();
insert_entries(pool, feed_id, &entries, 0).await?;
replace_sub_refs(pool, did, &[feed_id]).await?;
Ok(feed_id)
}
#[tokio::test]
async fn a_stored_future_date_is_cleared_at_startup() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://clock.example/f.xml".into(),
..Default::default()
},
)
.await?;
let tomorrow = (chrono::Utc::now() + chrono::Duration::days(1))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
for (guid, published) in [
("bogus", "2999-01-01T00:00:00Z"),
("soon", tomorrow.as_str()),
] {
sqlx::query(
"INSERT INTO entries (feed_id, guid, published, fetched_at) \
VALUES (?1, ?2, ?3, '2026-07-11T00:00:00Z')",
)
.bind(feed_id)
.bind(guid)
.bind(published)
.execute(&pool)
.await?;
}
apply_migrations(&pool).await?;
let dated: Vec<(String, Option<String>)> =
sqlx::query_as("SELECT guid, published FROM entries ORDER BY guid")
.fetch_all(&pool)
.await?;
assert_eq!(
dated[0],
("bogus".to_string(), None),
"a 2999 date survived startup"
);
assert_eq!(
dated[1].1.as_deref(),
Some(tomorrow.as_str()),
"a near-future date was cleared"
);
Ok(())
}
#[tokio::test]
async fn an_undated_entry_leads_the_reading_list_as_it_leads_the_cap() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:undated";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://undated.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "dated-old".to_string(),
title: Some("Old".to_string()),
published: Some("2024-01-01T00:00:00Z".to_string()),
..Default::default()
},
NewEntry {
guid: "dated-new".to_string(),
title: Some("Newer".to_string()),
published: Some("2025-01-01T00:00:00Z".to_string()),
..Default::default()
},
NewEntry {
guid: "undated".to_string(),
title: Some("Undated".to_string()),
..Default::default()
},
],
0,
)
.await?;
replace_sub_refs(&pool, did, &[feed_id]).await?;
let rows = list_entries(&pool, did, ListView::All, None, 100, 0).await?;
let order: Vec<&str> = rows.iter().map(|r| r.guid.as_str()).collect();
assert_eq!(
order,
vec!["undated", "dated-new", "dated-old"],
"the list disagrees with the cap about an undated entry's date",
);
let ids = list_entry_ids(&pool, did, ListView::All, None, 100).await?;
let by_guid: std::collections::HashMap<i64, &str> =
rows.iter().map(|r| (r.id, r.guid.as_str())).collect();
let id_order: Vec<&str> = ids.iter().filter_map(|i| by_guid.get(i).copied()).collect();
assert_eq!(
id_order,
vec!["undated", "dated-new", "dated-old"],
"the id projection orders differently from the list it projects",
);
Ok(())
}
#[tokio::test]
async fn list_entries_is_bounded_and_pages_without_overlap() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:pager";
seed_big_entries(&pool, did, 250).await?;
let page1 = list_entries(&pool, did, ListView::All, None, 100, 0).await?;
assert_eq!(page1.len(), 100, "limit was not applied");
let page2 = list_entries(&pool, did, ListView::All, None, 100, 100).await?;
let page3 = list_entries(&pool, did, ListView::All, None, 100, 200).await?;
assert_eq!(page3.len(), 50, "the last page should be the remainder");
let walked: Vec<i64> = page1
.iter()
.chain(&page2)
.chain(&page3)
.map(|e| e.id)
.collect();
let unique: std::collections::HashSet<i64> = walked.iter().copied().collect();
assert_eq!(unique.len(), 250, "paging repeated or skipped rows");
let whole = list_entries(&pool, did, ListView::All, None, 1_000, 0).await?;
assert_eq!(
walked,
whole.iter().map(|e| e.id).collect::<Vec<_>>(),
"paging changed the ordering"
);
let mut expected: Vec<(i64, i64)> = whole
.iter()
.map(|e| {
let day: i64 = e.published.as_deref().unwrap()[8..10].parse().unwrap();
(day, e.id)
})
.collect();
expected.sort_by(|a, b| b.cmp(a));
assert_eq!(
walked,
expected.iter().map(|(_, id)| *id).collect::<Vec<_>>(),
"ties are not broken by newest id"
);
assert_eq!(
count_entries_for_view(&pool, did, ListView::All, None).await?,
250,
"the unpaged count must survive paging"
);
Ok(())
}
#[tokio::test]
async fn the_list_projection_does_not_name_the_body_column() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:projection";
seed_big_entries(&pool, did, 3).await?;
let (sql, _) = list_entries_sql(ListView::All, None);
assert!(
!sql.contains("content_html") && !sql.contains("e.*"),
"the list query reads the article body: {sql}"
);
let rows = list_entries(&pool, did, ListView::All, None, 10, 0).await?;
assert_eq!(rows.len(), 3);
let widest = rows
.iter()
.map(|r| {
r.guid.len()
+ r.url.as_deref().map_or(0, str::len)
+ r.title.as_deref().map_or(0, str::len)
})
.max()
.unwrap_or(0);
assert!(
widest < 1_000,
"a list row carries {widest} bytes of text; the 20,000-byte body leaked in"
);
Ok(())
}
#[tokio::test]
async fn a_large_scope_is_one_bind_and_still_filters() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:widescope";
let mut all_ids = Vec::new();
for i in 0..300 {
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: format!("https://wide{i}.example/f.xml"),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: format!("w-{i}"),
..Default::default()
}],
0,
)
.await?;
all_ids.push(feed_id);
}
replace_sub_refs(&pool, did, &all_ids).await?;
let scope: Vec<i64> = all_ids.iter().copied().take(200).collect();
let rows = list_entries(&pool, did, ListView::All, Some(&scope), 1_000, 0).await?;
assert_eq!(rows.len(), 200, "the scope filter did not narrow correctly");
let in_scope: std::collections::HashSet<i64> = scope.iter().copied().collect();
assert!(
rows.iter().all(|r| in_scope.contains(&r.feed_id)),
"a feed outside the scope came back"
);
assert_eq!(
count_entries_for_view(&pool, did, ListView::All, Some(&scope)).await?,
200
);
let (sql, n) = list_query_sql(Projection::Ids, ListView::All, Some(&scope));
assert_eq!(n, 1, "the scope must contribute exactly one placeholder");
assert!(
sql.contains("json_each(?2)") && !sql.contains("?3"),
"the scope is still expanded into per-id placeholders: {sql}"
);
Ok(())
}
#[tokio::test]
async fn a_feed_scope_narrows_the_query_not_the_page() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:scope";
let wanted = seed_big_entries(&pool, did, 10).await?;
let other = upsert_feed(
&pool,
&NewFeed {
url: "https://other.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
let noise: Vec<NewEntry> = (0..40)
.map(|i| NewEntry {
guid: format!("noise-{i}"),
published: Some("2027-01-01T00:00:00Z".to_string()),
..Default::default()
})
.collect();
insert_entries(&pool, other, &noise, 0).await?;
replace_sub_refs(&pool, did, &[wanted, other]).await?;
let scoped = list_entries(&pool, did, ListView::All, Some(&[wanted]), 10, 0).await?;
assert_eq!(
scoped.len(),
10,
"the scoped page came back short — the filter ran after the LIMIT"
);
assert!(scoped.iter().all(|e| e.feed_id == wanted));
assert!(list_entries(&pool, did, ListView::All, Some(&[]), 10, 0)
.await?
.is_empty());
assert_eq!(
count_entries_for_view(&pool, did, ListView::All, Some(&[])).await?,
0
);
Ok(())
}
#[tokio::test]
async fn list_rows_carry_their_own_read_and_star_bits() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:bits";
seed_big_entries(&pool, did, 3).await?;
let ids: Vec<i64> = list_entries(&pool, did, ListView::All, None, 10, 0)
.await?
.iter()
.map(|e| e.id)
.collect();
mark_read(&pool, did, ids[0], true).await?;
mark_starred(&pool, did, ids[1], true).await?;
let all = list_entries(&pool, did, ListView::All, None, 10, 0).await?;
let by_id = |id: i64| all.iter().find(|e| e.id == id).expect("row present");
assert!(by_id(ids[0]).read && !by_id(ids[0]).starred);
assert!(!by_id(ids[1]).read && by_id(ids[1]).starred);
assert!(!by_id(ids[2]).read && !by_id(ids[2]).starred);
let unread = list_entries(&pool, did, ListView::Unread, None, 10, 0).await?;
assert_eq!(unread.len(), 2);
assert!(unread.iter().all(|e| !e.read));
let starred = list_entries(&pool, did, ListView::Starred, None, 10, 0).await?;
assert_eq!(starred.len(), 1);
assert_eq!(starred[0].id, ids[1]);
Ok(())
}
#[tokio::test]
async fn unread_counts_are_per_feed_and_exclude_read_rows() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:counts";
let a = seed_big_entries(&pool, did, 5).await?;
let b = upsert_feed(
&pool,
&NewFeed {
url: "https://b.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
b,
&[
NewEntry {
guid: "b-1".to_string(),
..Default::default()
},
NewEntry {
guid: "b-2".to_string(),
..Default::default()
},
],
0,
)
.await?;
replace_sub_refs(&pool, did, &[a, b]).await?;
let first_a = list_entries(&pool, did, ListView::All, Some(&[a]), 1, 0).await?[0].id;
mark_read(&pool, did, first_a, true).await?;
let counts = unread_counts_by_feed(&pool, did).await?;
assert_eq!(counts.get(&a).copied(), Some(4));
assert_eq!(counts.get(&b).copied(), Some(2));
replace_sub_refs(&pool, did, &[b]).await?;
let counts = unread_counts_by_feed(&pool, did).await?;
assert_eq!(counts.get(&a), None);
assert_eq!(counts.get(&b).copied(), Some(2));
Ok(())
}
#[tokio::test]
async fn compaction_folds_read_ids_into_the_water_mark() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:compact";
let feed_url = "https://compact.example/f.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
..Default::default()
},
)
.await?;
let entries: Vec<NewEntry> = (0..40)
.map(|i| NewEntry {
guid: format!("c-{i:03}"),
published: Some(format!("2026-01-{:02}T00:00:00Z", i + 1)),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
replace_sub_refs(&pool, did, &[feed_id]).await?;
let all = list_entries(&pool, did, ListView::All, None, 100, 0).await?;
let mut oldest_first = all.clone();
oldest_first.reverse();
for row in oldest_first.iter().take(30) {
mark_read(&pool, did, row.id, true).await?;
}
let before = get_cursor(&pool, did, feed_url).await?.expect("cursor");
assert!(before.read_through.is_none(), "read_through starts unset");
let before_ids: Vec<String> = serde_json::from_str(&before.read_ids)?;
assert_eq!(before_ids.len(), 30, "every read is its own exception");
let watermark = compact_cursor(&pool, did, feed_url)
.await?
.expect("the water-mark must advance");
let after = get_cursor(&pool, did, feed_url).await?.expect("cursor");
assert_eq!(after.read_through.as_deref(), Some(watermark.as_str()));
let after_ids: Vec<String> = serde_json::from_str(&after.read_ids)?;
assert!(
after_ids.is_empty(),
"a contiguous read prefix must fold entirely into the water-mark, left {after_ids:?}"
);
assert_eq!(watermark, "2026-01-30T00:00:00Z");
assert!(after.dirty, "a rewritten cursor must be re-flushed");
Ok(())
}
#[tokio::test]
async fn compaction_stops_below_the_oldest_unread_entry() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:gap";
let feed_url = "https://gap.example/f.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
..Default::default()
},
)
.await?;
let entries: Vec<NewEntry> = (0..10)
.map(|i| NewEntry {
guid: format!("g-{i:02}"),
published: Some(format!("2026-02-{:02}T00:00:00Z", i + 1)),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
replace_sub_refs(&pool, did, &[feed_id]).await?;
let mut oldest_first = list_entries(&pool, did, ListView::All, None, 100, 0).await?;
oldest_first.reverse();
for (i, row) in oldest_first.iter().enumerate() {
if i != 2 {
mark_read(&pool, did, row.id, true).await?;
}
}
let watermark = compact_cursor(&pool, did, feed_url)
.await?
.expect("advances");
assert_eq!(
watermark, "2026-02-02T00:00:00Z",
"the water-mark jumped the unread hole"
);
let after = get_cursor(&pool, did, feed_url).await?.expect("cursor");
let kept: Vec<String> = serde_json::from_str(&after.read_ids)?;
assert_eq!(
kept.len(),
7,
"the 7 reads ABOVE the hole must stay as explicit exceptions"
);
let unread: Vec<String> = serde_json::from_str(&after.unread_ids)?;
assert!(
unread.is_empty(),
"redundant unread exceptions survived: {unread:?}"
);
assert_eq!(
compact_cursor(&pool, did, feed_url).await?,
None,
"a second compaction moved a water-mark that was already correct"
);
Ok(())
}
#[tokio::test]
async fn compaction_never_invents_a_water_mark() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:none";
let feed_url = "https://none.example/f.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
..Default::default()
},
)
.await?;
replace_sub_refs(&pool, did, &[feed_id]).await?;
assert_eq!(compact_cursor(&pool, did, feed_url).await?, None);
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "n-1".to_string(),
published: Some("2026-03-01T00:00:00Z".to_string()),
..Default::default()
},
NewEntry {
guid: "n-2".to_string(),
published: Some("2026-03-02T00:00:00Z".to_string()),
..Default::default()
},
],
0,
)
.await?;
assert_eq!(
compact_cursor(&pool, did, feed_url).await?,
None,
"a water-mark appeared with nothing read — that asserts the backlog is read"
);
Ok(())
}
#[tokio::test]
async fn a_star_can_be_cleared_after_unsubscribing_from_its_feed() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:unsub";
let feed_id = seed_big_entries(&pool, did, 3).await?;
let rows = list_entries(&pool, did, ListView::All, None, 10, 0).await?;
let target = rows[0].clone();
mark_starred(&pool, did, target.id, true).await?;
assert_eq!(get_starred_for_did(&pool, did).await?.len(), 1);
replace_sub_refs(&pool, did, &[]).await?;
assert!(
get_starred_for_did(&pool, did).await?.is_empty(),
"fixture precondition: the star must be invisible to the scoped read"
);
assert!(
matches!(
starred_identities(&pool, did, 1_000).await?,
StarredIdentities::All(ref v) if v.is_empty()
),
"fixture precondition: the identity lookup must miss it too"
);
let still_starred: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM entry_state WHERE did = ?1 AND starred = 1")
.bind(did)
.fetch_one(&pool)
.await?;
assert_eq!(
still_starred, 1,
"the star is still there, just unreachable"
);
let cleared =
clear_star_by_identity(&pool, did, target.url.as_deref(), Some(&target.guid)).await?;
assert_eq!(cleared, 1, "the star survived the unsave");
let after: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM entry_state WHERE did = ?1 AND starred = 1")
.bind(did)
.fetch_one(&pool)
.await?;
assert_eq!(after, 0);
replace_sub_refs(&pool, did, &[feed_id]).await?;
assert!(
get_starred_for_did(&pool, did).await?.is_empty(),
"the star came back after resubscribing — the desync is still there"
);
Ok(())
}
#[tokio::test]
async fn clearing_a_star_touches_only_that_did_and_that_article() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let mine = "did:plc:mine";
let theirs = "did:plc:theirs";
let feed_id = seed_big_entries(&pool, mine, 3).await?;
replace_sub_refs(&pool, theirs, &[feed_id]).await?;
let rows = list_entries(&pool, mine, ListView::All, None, 10, 0).await?;
for r in &rows {
mark_starred(&pool, mine, r.id, true).await?;
mark_starred(&pool, theirs, r.id, true).await?;
}
let target = &rows[1];
assert_eq!(
clear_star_by_identity(&pool, mine, target.url.as_deref(), Some(&target.guid)).await?,
1
);
let count = |did: &'static str| {
let pool = pool.clone();
async move {
sqlx::query_scalar::<_, i64>(
"SELECT COUNT(*) FROM entry_state WHERE did = ?1 AND starred = 1",
)
.bind(did)
.fetch_one(&pool)
.await
.unwrap()
}
};
assert_eq!(count(mine).await, 2, "it cleared more than the one article");
assert_eq!(count(theirs).await, 3, "it cleared another DID's stars");
assert_eq!(
clear_star_by_identity(&pool, mine, Some("https://nope.example/x"), Some("nope"))
.await?,
0
);
assert_eq!(count(mine).await, 2);
let before: String = sqlx::query_scalar(
"SELECT updated_at FROM entry_state WHERE did = ?1 AND entry_id = ?2",
)
.bind(mine)
.bind(target.id)
.fetch_one(&pool)
.await?;
assert_eq!(
clear_star_by_identity(&pool, mine, target.url.as_deref(), Some(&target.guid)).await?,
0,
"a second clear reported rows it did not change"
);
let after: String = sqlx::query_scalar(
"SELECT updated_at FROM entry_state WHERE did = ?1 AND entry_id = ?2",
)
.bind(mine)
.bind(target.id)
.fetch_one(&pool)
.await?;
assert_eq!(before, after, "a no-op clear rewrote updated_at");
assert_eq!(clear_star_by_identity(&pool, mine, None, None).await?, 0);
assert_eq!(
clear_star_by_identity(&pool, mine, Some(""), Some("")).await?,
0
);
assert_eq!(count(mine).await, 2);
Ok(())
}
#[tokio::test]
async fn starred_identities_span_the_whole_set() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:ident";
seed_big_entries(&pool, did, 150).await?;
for row in list_entries(&pool, did, ListView::All, None, 1_000, 0).await? {
mark_starred(&pool, did, row.id, true).await?;
}
let identities = match starred_identities(&pool, did, 20_000).await? {
StarredIdentities::All(v) => v,
StarredIdentities::Truncated => panic!("150 rows must not read as truncated"),
};
assert_eq!(
identities.len(),
150,
"the identity set was truncated to a page"
);
assert!(identities
.iter()
.all(|(url, guid)| url.is_some() && !guid.is_empty()));
assert!(
matches!(
starred_identities(&pool, did, 10).await?,
StarredIdentities::Truncated
),
"a truncated identity set reported itself as complete"
);
assert!(
matches!(
starred_identities(&pool, did, 150).await?,
StarredIdentities::All(ref v) if v.len() == 150
),
"a set exactly at the cap was misreported as truncated"
);
Ok(())
}
#[tokio::test]
async fn entry_ids_are_ordered_and_capped() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:ids";
seed_big_entries(&pool, did, 60).await?;
let capped = list_entry_ids(&pool, did, ListView::All, None, 25).await?;
assert_eq!(capped.len(), 25);
let rows = list_entries(&pool, did, ListView::All, None, 25, 0).await?;
assert_eq!(
capped,
rows.iter().map(|e| e.id).collect::<Vec<_>>(),
"the id list and the row list disagree on ordering"
);
Ok(())
}
#[tokio::test]
async fn mark_read_dirties_the_feed_cursor() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_url = "https://example.com/feed.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
title: Some("Example".to_string()),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "g1".to_string(),
published: Some("2026-07-10T00:00:00Z".to_string()),
..Default::default()
},
NewEntry {
guid: "g2".to_string(),
published: Some("2026-07-11T00:00:00Z".to_string()),
..Default::default()
},
],
0,
)
.await?;
let did = "did:plc:reader";
replace_sub_refs(&pool, did, &[feed_id]).await?;
assert!(get_cursor(&pool, did, feed_url).await?.is_none());
assert_eq!(dirty_cursors(&pool, did).await?.len(), 0);
let e1 = entries_for_feed(&pool, did, feed_id).await?[0].id;
assert!(mark_read(&pool, did, e1, true).await?);
let cursor = get_cursor(&pool, did, feed_url)
.await?
.expect("mark_read must create the feed's read_cursor");
assert!(cursor.dirty, "cursor must be dirty after mark_read");
assert!(
cursor.read_ids.contains(&e1.to_string()),
"the read entry id must be in read_ids: {}",
cursor.read_ids
);
let dirty = dirty_cursors(&pool, did).await?;
assert_eq!(dirty.len(), 1, "flusher must see the newly dirty cursor");
assert_eq!(dirty[0].feed_url, feed_url);
assert!(mark_read(&pool, did, e1, false).await?);
let cursor = get_cursor(&pool, did, feed_url).await?.unwrap();
assert!(cursor.dirty);
assert!(
cursor.unread_ids.contains(&e1.to_string()),
"unread id must be in unread_ids: {}",
cursor.unread_ids
);
assert!(
!cursor.read_ids.contains(&e1.to_string()),
"id must have left read_ids: {}",
cursor.read_ids
);
assert!(mark_feed_read(&pool, did, feed_id, true).await? > 0);
let cursor = get_cursor(&pool, did, feed_url).await?.unwrap();
assert!(cursor.dirty);
assert_eq!(dirty_cursors(&pool, did).await?.len(), 1);
let outsider = "did:plc:outsider";
assert!(!mark_read(&pool, outsider, e1, true).await?);
assert_eq!(dirty_cursors(&pool, outsider).await?.len(), 0);
let snap = dirty_cursors(&pool, did).await?[0].clone();
clear_cursor_dirty(&pool, did, feed_url, "1999-01-01T00:00:00Z").await?;
assert_eq!(
dirty_cursors(&pool, did).await?.len(),
1,
"stale-snapshot clear must be a no-op"
);
clear_cursor_dirty(&pool, did, feed_url, &snap.updated_at).await?;
assert_eq!(dirty_cursors(&pool, did).await?.len(), 0);
Ok(())
}
#[test]
fn json_id_set_toggle_is_set_like() {
let s = json_id_set_toggle("[]", 5, true);
assert_eq!(s, r#"["5"]"#);
assert_eq!(json_id_set_toggle(&s, 5, true), r#"["5"]"#); let s = json_id_set_toggle(&s, 7, true);
assert_eq!(s, r#"["5","7"]"#);
let s = json_id_set_toggle(&s, 5, false);
assert_eq!(s, r#"["7"]"#);
assert_eq!(json_id_set_toggle("[1,2]", 3, true), r#"["1","2","3"]"#);
assert_eq!(json_id_set_toggle("garbage", 1, true), r#"["1"]"#);
}
#[tokio::test]
async fn per_did_isolation_scopes_reads_and_mutations() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_a = upsert_feed(
&pool,
&NewFeed {
url: "https://a.example/feed.xml".to_string(),
title: Some("A".to_string()),
..Default::default()
},
)
.await?;
let feed_b = upsert_feed(
&pool,
&NewFeed {
url: "https://b.example/feed.xml".to_string(),
title: Some("B".to_string()),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_a,
&[NewEntry {
guid: "a-1".to_string(),
url: Some("https://a.example/1".to_string()),
title: Some("A one".to_string()),
published: Some("2026-07-10T00:00:00Z".to_string()),
content_html: Some("<p>secret A body</p>".to_string()),
..Default::default()
}],
0,
)
.await?;
insert_entries(
&pool,
feed_b,
&[NewEntry {
guid: "b-1".to_string(),
url: Some("https://b.example/1".to_string()),
title: Some("B one".to_string()),
published: Some("2026-07-11T00:00:00Z".to_string()),
content_html: Some("<p>secret B body</p>".to_string()),
..Default::default()
}],
0,
)
.await?;
let did_a = "did:plc:aaaa";
let did_b = "did:plc:bbbb";
replace_sub_refs(&pool, did_a, &[feed_a]).await?;
replace_sub_refs(&pool, did_b, &[feed_b]).await?;
let b_entry_id = entries_for_feed(&pool, did_b, feed_b).await?[0].id;
assert_eq!(entries_for_feed(&pool, did_a, feed_a).await?.len(), 1);
assert!(
entries_for_feed(&pool, did_a, feed_b).await?.is_empty(),
"A must not read entries of a feed it does not subscribe to"
);
let unread_a = get_unread_for_did(&pool, did_a).await?;
assert_eq!(unread_a.len(), 1);
assert_eq!(unread_a[0].guid, "a-1");
let unread_b = get_unread_for_did(&pool, did_b).await?;
assert_eq!(unread_b.len(), 1);
assert_eq!(unread_b[0].guid, "b-1");
assert!(did_subscribes_to_entry(&pool, did_b, b_entry_id).await?);
assert!(
!did_subscribes_to_entry(&pool, did_a, b_entry_id).await?,
"A does not subscribe to B's feed"
);
assert!(
!mark_read(&pool, did_a, b_entry_id, true).await?,
"non-subscriber mark_read must be a no-op (→ 404), never a mutation"
);
assert_eq!(get_unread_for_did(&pool, did_b).await?.len(), 1);
assert!(mark_read(&pool, did_b, b_entry_id, true).await?);
assert_eq!(get_unread_for_did(&pool, did_b).await?.len(), 0);
assert!(
!mark_starred(&pool, did_a, b_entry_id, true).await?,
"non-subscriber mark_starred must be a no-op (→ 404)"
);
assert!(
get_starred_for_did(&pool, did_a).await?.is_empty(),
"A's starred list stays empty after the rejected attempt"
);
assert!(mark_starred(&pool, did_b, b_entry_id, true).await?);
assert_eq!(get_starred_for_did(&pool, did_b).await?.len(), 1);
assert!(get_starred_for_did(&pool, did_a).await?.is_empty());
let a_feeds = feeds_for_did(&pool, did_a).await?;
assert_eq!(a_feeds.len(), 1);
assert_eq!(a_feeds[0].id, feed_a);
let b_feeds = feeds_for_did(&pool, did_b).await?;
assert_eq!(b_feeds.len(), 1);
assert_eq!(b_feeds[0].id, feed_b);
replace_sub_refs(&pool, did_b, &[]).await?;
assert!(get_unread_for_did(&pool, did_b).await?.is_empty());
assert!(get_starred_for_did(&pool, did_b).await?.is_empty());
assert!(entries_for_feed(&pool, did_b, feed_b).await?.is_empty());
assert!(feeds_for_did(&pool, did_b).await?.is_empty());
Ok(())
}
#[tokio::test]
async fn pds_outage_fallback_fails_closed_not_open() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did_a = "did:plc:aaaa";
let feed_a = upsert_feed(
&pool,
&NewFeed {
url: "https://a.example/feed.xml".to_string(),
title: Some("A".to_string()),
..Default::default()
},
)
.await?;
let feed_orphan = upsert_feed(
&pool,
&NewFeed {
url: "https://orphan.example/feed.xml".to_string(),
title: Some("Orphan".to_string()),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_a,
&[NewEntry {
guid: "a-1".to_string(),
url: Some("https://a.example/1".to_string()),
title: Some("A one".to_string()),
published: Some("2026-07-10T00:00:00Z".to_string()),
content_html: Some("<p>A body</p>".to_string()),
..Default::default()
}],
0,
)
.await?;
insert_entries(
&pool,
feed_orphan,
&[NewEntry {
guid: "orphan-1".to_string(),
url: Some("https://orphan.example/1".to_string()),
title: Some("Orphan one".to_string()),
published: Some("2026-07-11T00:00:00Z".to_string()),
content_html: Some("<p>secret orphan body</p>".to_string()),
..Default::default()
}],
0,
)
.await?;
replace_sub_refs(&pool, did_a, &[feed_a]).await?;
replace_sub_refs(&pool, "did:plc:seed", &[feed_orphan]).await?;
let orphan_entry_id = entries_for_feed(&pool, "did:plc:seed", feed_orphan).await?[0].id;
replace_sub_refs(&pool, "did:plc:seed", &[]).await?;
let fallback = feeds_for_did(&pool, did_a).await?;
let fallback_ids: Vec<i64> = fallback.iter().map(|f| f.id).collect();
assert_eq!(
fallback_ids,
vec![feed_a],
"outage fallback must serve ONLY A's own last-known sub_ref, \
never widen to the orphan cached feed"
);
assert!(
!fallback_ids.contains(&feed_orphan),
"FAIL-OPEN regression: outage fallback leaked an unsubscribed \
cached feed into A's surface"
);
assert!(
!did_subscribes_to_entry(&pool, did_a, orphan_entry_id).await?,
"A must not be authorized for an orphan feed's entry during an outage"
);
assert!(
entries_for_feed(&pool, did_a, feed_orphan)
.await?
.is_empty(),
"entries_for_feed must not expose the orphan feed to A during an outage"
);
let unread_guids: Vec<String> = get_unread_for_did(&pool, did_a)
.await?
.into_iter()
.map(|e| e.guid)
.collect();
assert!(
!unread_guids.iter().any(|g| g == "orphan-1"),
"orphan entry leaked into A's unread list during an outage"
);
assert!(
get_starred_for_did(&pool, did_a).await?.is_empty(),
"A has no starred entries; the orphan must not appear"
);
assert!(
!mark_read(&pool, did_a, orphan_entry_id, true).await?,
"A must not mark an orphan feed's entry read during an outage"
);
assert!(
!mark_starred(&pool, did_a, orphan_entry_id, true).await?,
"A must not star an orphan feed's entry during an outage"
);
assert_eq!(
mark_feed_read(&pool, did_a, feed_orphan, true).await?,
0,
"A must not mark-all-read the orphan feed during an outage"
);
let es_count: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM entry_state WHERE did = ?1 AND entry_id = ?2")
.bind(did_a)
.bind(orphan_entry_id)
.fetch_one(&pool)
.await?;
assert_eq!(es_count, 0, "no cross-tenant mutation during the outage");
Ok(())
}
#[test]
fn code_gen_shape_and_alphabet() {
for _ in 0..200 {
let code = generate_invite_code().unwrap();
assert!(code.starts_with("FEATHER-"), "bad prefix: {code}");
let body = &code["FEATHER-".len()..];
assert_eq!(body.len(), CODE_BODY_LEN, "bad body length: {code}");
for c in body.chars() {
assert!(
CODE_ALPHABET.contains(&(c as u8)),
"char {c:?} not in alphabet ({code})"
);
assert!(
!matches!(c, 'I' | 'O' | '0' | '1'),
"ambiguous char {c:?} leaked into {code}"
);
}
}
assert_ne!(
generate_invite_code().unwrap(),
generate_invite_code().unwrap()
);
}
#[tokio::test]
async fn busy_timeout_is_applied() -> Result<()> {
let dir = std::env::temp_dir().join(format!("fr-busy-{}", std::process::id()));
std::fs::create_dir_all(&dir).ok();
let path = dir.join("busy.db");
let url = format!("sqlite://{}", path.display());
let pool = init_url(&url).await?;
let row = sqlx::query("PRAGMA busy_timeout").fetch_one(&pool).await?;
let timeout: i64 = row.get(0);
assert_eq!(timeout, 5000, "busy_timeout should be 5000 ms");
pool.close().await;
std::fs::remove_dir_all(&dir).ok();
Ok(())
}
#[tokio::test]
async fn redeem_valid_grants_seat() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let code = mint_code(&pool, "did:plc:creator", 3600).await?;
assert!(!has_beta_access(&pool, "did:plc:new").await?);
let out = redeem_code(&pool, &code, "did:plc:new", Some("new.bsky"), 100).await?;
assert_eq!(out, Ok(()));
assert!(has_beta_access(&pool, "did:plc:new").await?);
assert_eq!(count_beta_access(&pool).await?, 1);
let again = redeem_code(&pool, &code, "did:plc:other", None, 100).await?;
assert_eq!(again, Err(RedeemError::AlreadyRedeemed));
Ok(())
}
#[tokio::test]
async fn redeem_not_found() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let out = redeem_code(&pool, "FEATHER-NOPENOPE", "did:plc:x", None, 100).await?;
assert_eq!(out, Err(RedeemError::NotFound));
Ok(())
}
async fn insert_expired_code(pool: &SqlitePool, code: &str, creator: &str) -> Result<()> {
let now = now_unix();
sqlx::query(
r#"INSERT INTO invite_codes
(code, creator_did, status, invitee_did, created_at, expires_at, redeemed_at)
VALUES (?1, ?2, 'active', NULL, ?3, ?4, NULL)"#,
)
.bind(code)
.bind(creator)
.bind(now - 100)
.bind(now - 10) .execute(pool)
.await?;
Ok(())
}
#[tokio::test]
async fn redeem_expired() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
insert_expired_code(&pool, "FEATHER-EXPIRED0", "did:plc:creator").await?;
let out = redeem_code(&pool, "FEATHER-EXPIRED0", "did:plc:new", None, 100).await?;
assert_eq!(out, Err(RedeemError::Expired));
assert_eq!(count_beta_access(&pool).await?, 0);
Ok(())
}
#[tokio::test]
async fn redeem_capacity_full() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
ensure_seed(&pool, &["did:plc:admin".to_string()]).await?;
assert_eq!(count_beta_access(&pool).await?, 1);
let code = mint_code(&pool, "did:plc:admin", 3600).await?;
let out = redeem_code(&pool, &code, "did:plc:new", None, 1).await?;
assert_eq!(out, Err(RedeemError::CapacityFull));
assert!(!has_beta_access(&pool, "did:plc:new").await?);
let ok = redeem_code(&pool, &code, "did:plc:new", None, 2).await?;
assert_eq!(ok, Ok(()));
Ok(())
}
#[tokio::test]
async fn count_active_codes_excludes_expired_and_redeemed() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
assert_eq!(count_active_codes(&pool).await?, 0);
let a = mint_code(&pool, "did:plc:bot", 3600).await?;
let _b = mint_code(&pool, "did:plc:bot", 3600).await?;
assert_eq!(count_active_codes(&pool).await?, 2);
insert_expired_code(&pool, "FEATHER-EXPIRED0", "did:plc:bot").await?;
assert_eq!(count_active_codes(&pool).await?, 2);
let out = redeem_code(&pool, &a, "did:plc:new", None, 100).await?;
assert_eq!(out, Ok(()));
assert_eq!(count_active_codes(&pool).await?, 1);
Ok(())
}
#[tokio::test]
async fn expire_and_seed() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
insert_expired_code(&pool, "FEATHER-EXPIRED1", "did:plc:creator").await?;
let live = mint_code(&pool, "did:plc:creator", 3600).await?;
let n = expire_old_codes(&pool).await?;
assert_eq!(n, 1, "exactly the past-expiry code should flip");
assert_eq!(
redeem_code(&pool, &live, "did:plc:new", None, 100).await?,
Ok(())
);
let created = ensure_seed(
&pool,
&["did:plc:seed1".to_string(), "did:plc:seed2".to_string()],
)
.await?;
assert_eq!(created, 2);
let created2 = ensure_seed(&pool, &["did:plc:seed1".to_string()]).await?;
assert_eq!(created2, 0, "re-seeding an existing DID is a no-op");
assert!(has_beta_access(&pool, "did:plc:seed1").await?);
Ok(())
}
#[tokio::test]
async fn the_expiry_sweep_spares_redeemed_codes() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let code = mint_code(&pool, "did:plc:creator", 3600).await?;
assert!(redeem_code(&pool, &code, "did:plc:new", None, 100)
.await?
.is_ok());
sqlx::query("UPDATE invite_codes SET expires_at = ?1 WHERE code = ?2")
.bind(now_unix() - 10)
.bind(&code)
.execute(&pool)
.await?;
insert_expired_code(&pool, "FEATHER-EXPIRED2", "did:plc:creator").await?;
let n = expire_old_codes(&pool).await?;
assert_eq!(n, 1, "the sweep counted the redeemed code");
let status: String = sqlx::query_scalar("SELECT status FROM invite_codes WHERE code = ?1")
.bind(&code)
.fetch_one(&pool)
.await?;
assert_eq!(status, "redeemed", "the sweep rewrote a redemption");
Ok(())
}
#[tokio::test]
async fn count_helpers_track_feeds_and_subs() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
assert_eq!(count_feeds(&pool).await?, 0);
let mut ids = Vec::new();
for i in 0..3 {
let id = upsert_feed(
&pool,
&NewFeed {
url: format!("https://f{i}.example/feed.xml"),
..Default::default()
},
)
.await?;
ids.push(id);
}
assert_eq!(count_feeds(&pool).await?, 3);
let did = "did:plc:capcheck";
assert_eq!(count_subscriptions_for_did(&pool, did).await?, 0);
replace_sub_refs(&pool, did, &ids).await?;
assert_eq!(count_subscriptions_for_did(&pool, did).await?, 3);
Ok(())
}
#[tokio::test]
async fn insert_entries_trims_over_cap_keeping_newest() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://firehose.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let batch: Vec<NewEntry> = (0..5)
.map(|i| NewEntry {
guid: format!("g-{i}"),
title: Some(format!("E{i}")),
published: Some(format!("2026-07-0{}T00:00:00Z", i + 1)),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &batch, 2).await?;
let did = "did:plc:trim";
replace_sub_refs(&pool, did, &[feed_id]).await?;
let kept = entries_for_feed(&pool, did, feed_id).await?;
assert_eq!(
kept.len(),
2,
"over-cap feed trimmed to the newest 2 entries"
);
assert_eq!(kept[0].guid, "g-4");
assert_eq!(kept[1].guid, "g-3");
Ok(())
}
#[tokio::test]
async fn insert_entries_trims_keeps_fresh_undated_over_stale_dated() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://undated.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let batch = vec![
NewEntry {
guid: "old-dated-1".to_string(),
title: Some("Old A".to_string()),
published: Some("2026-07-01T00:00:00Z".to_string()),
fetched_at: Some("2026-07-01T00:00:00Z".to_string()),
..Default::default()
},
NewEntry {
guid: "old-dated-2".to_string(),
title: Some("Old B".to_string()),
published: Some("2026-07-02T00:00:00Z".to_string()),
fetched_at: Some("2026-07-02T00:00:00Z".to_string()),
..Default::default()
},
NewEntry {
guid: "fresh-undated".to_string(),
title: Some("Fresh undated".to_string()),
published: None,
fetched_at: Some("2026-07-11T00:00:00Z".to_string()),
..Default::default()
},
];
insert_entries(&pool, feed_id, &batch, 2).await?;
let did = "did:plc:undated";
replace_sub_refs(&pool, did, &[feed_id]).await?;
let kept = entries_for_feed(&pool, did, feed_id).await?;
assert_eq!(kept.len(), 2, "over-cap feed trimmed to 2 entries");
let guids: Vec<&str> = kept.iter().map(|e| e.guid.as_str()).collect();
assert!(
guids.contains(&"fresh-undated"),
"the freshly-fetched undated entry must survive the trim, kept: {guids:?}"
);
assert!(
guids.contains(&"old-dated-2"),
"the newer dated entry survives; the OLDEST dated entry is the one evicted, kept: {guids:?}"
);
assert!(
!guids.contains(&"old-dated-1"),
"the oldest dated entry is the one that should be evicted, kept: {guids:?}"
);
Ok(())
}
#[tokio::test]
async fn db_size_is_positive_and_grows() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let before = db_size_bytes(&pool).await?;
assert!(before > 0, "a schema-initialised DB has a non-zero size");
seed_big_entries(&pool, "did:plc:growth", 400).await?;
let after = db_size_bytes(&pool).await?;
assert!(
after > before,
"the database grew by {} bytes after 400 seeded entries; the size is \
not tracking the data, so the disk watermark can never trip",
after.saturating_sub(before),
);
Ok(())
}
#[tokio::test]
async fn purge_did_data_removes_only_the_callers_rows() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://example.com/feed.xml".to_string(),
title: Some("Example".to_string()),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: "g-1".to_string(),
url: Some("https://example.com/a".to_string()),
title: Some("First".to_string()),
published: Some("2026-07-10T08:00:00Z".to_string()),
..Default::default()
}],
0,
)
.await?;
let entry_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'g-1'")
.fetch_one(&pool)
.await?;
let victim = "did:plc:victim";
let bystander = "did:plc:bystander";
for did in [victim, bystander] {
replace_sub_refs(&pool, did, &[feed_id]).await?;
assert!(mark_read(&pool, did, entry_id, true).await?);
assert!(mark_starred(&pool, did, entry_id, true).await?);
upsert_cursor(
&pool,
&ReadCursor {
did: did.to_string(),
feed_url: "https://example.com/feed.xml".to_string(),
read_through: Some("2026-07-10T08:00:00Z".to_string()),
read_ids: "[]".to_string(),
unread_ids: "[]".to_string(),
dirty: false,
pds_created: false,
updated_at: now_rfc3339(),
},
)
.await?;
grant_access(&pool, did, Some("h.example"), "admin", None).await?;
mint_code(&pool, did, 3600).await?;
}
let counts = purge_did_data(&pool, victim).await?;
assert_eq!(
counts.entry_state, 1,
"one entry_state row (read+star merge)"
);
assert_eq!(counts.read_cursor, 1);
assert_eq!(counts.sub_ref, 1);
assert_eq!(counts.beta_access, 1);
assert_eq!(counts.invite_codes, 1);
assert_eq!(counts.total(), 5);
let es: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entry_state WHERE did = ?1")
.bind(victim)
.fetch_one(&pool)
.await?;
assert_eq!(es, 0, "victim still had entry_state rows");
let rc: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM read_cursor WHERE did = ?1")
.bind(victim)
.fetch_one(&pool)
.await?;
assert_eq!(rc, 0, "victim still had read_cursor rows");
let sr: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM sub_ref WHERE did = ?1")
.bind(victim)
.fetch_one(&pool)
.await?;
assert_eq!(sr, 0, "victim still had sub_ref rows");
assert!(
!has_beta_access(&pool, victim).await?,
"victim still had a beta seat"
);
let victim_codes: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM invite_codes WHERE creator_did = ?1")
.bind(victim)
.fetch_one(&pool)
.await?;
assert_eq!(victim_codes, 0);
assert!(has_beta_access(&pool, bystander).await?);
let bystander_subs = count_subscriptions_for_did(&pool, bystander).await?;
assert_eq!(bystander_subs, 1, "bystander's sub_ref survived");
let bystander_codes: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM invite_codes WHERE creator_did = ?1")
.bind(bystander)
.fetch_one(&pool)
.await?;
assert_eq!(bystander_codes, 1);
assert_eq!(count_feeds(&pool).await?, 1);
let again = purge_did_data(&pool, victim).await?;
assert_eq!(again.total(), 0);
Ok(())
}
#[tokio::test]
async fn purge_did_data_scrubs_cross_did_back_references() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let inviter = "did:plc:inviter";
let leaver = "did:plc:leaver";
let friend = "did:plc:friend";
let inviter_code = mint_code(&pool, inviter, 3600).await?;
grant_access(&pool, inviter, None, "admin", None).await?;
assert_eq!(
redeem_code(&pool, &inviter_code, leaver, Some("leaver.bsky"), 100).await?,
Ok(())
);
let leaver_code = mint_code(&pool, leaver, 3600).await?;
assert_eq!(
redeem_code(&pool, &leaver_code, friend, Some("friend.bsky"), 100).await?,
Ok(())
);
let invitee_before: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM invite_codes WHERE invitee_did = ?1")
.bind(leaver)
.fetch_one(&pool)
.await?;
assert_eq!(
invitee_before, 1,
"leaver should be an invitee before purge"
);
let granted_before: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM beta_access WHERE granted_by = ?1")
.bind(leaver)
.fetch_one(&pool)
.await?;
assert_eq!(granted_before, 1, "leaver should be a granter before purge");
let counts = purge_did_data(&pool, leaver).await?;
assert_eq!(
counts.invitee_scrubbed, 1,
"the redeemed code's invitee_did"
);
assert_eq!(counts.granted_by_scrubbed, 1, "the seat leaver granted");
let invitee_after: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM invite_codes WHERE invitee_did = ?1")
.bind(leaver)
.fetch_one(&pool)
.await?;
assert_eq!(invitee_after, 0, "leaver survived in invitee_did");
let granted_after: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM beta_access WHERE granted_by = ?1")
.bind(leaver)
.fetch_one(&pool)
.await?;
assert_eq!(granted_after, 0, "leaver survived in granted_by");
assert!(
has_beta_access(&pool, friend).await?,
"friend's seat must survive the leaver's scrub"
);
let friend_granted_by: String =
sqlx::query_scalar("SELECT granted_by FROM beta_access WHERE did = ?1")
.bind(friend)
.fetch_one(&pool)
.await?;
assert_eq!(friend_granted_by, REDACTED_DID);
let inviter_code_rows: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM invite_codes WHERE creator_did = ?1")
.bind(inviter)
.fetch_one(&pool)
.await?;
assert_eq!(inviter_code_rows, 1, "inviter's code row must survive");
Ok(())
}
#[tokio::test]
async fn the_migration_clears_error_counts_on_unpollable_at_uri_rows() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
for url in [
"https://real.example/feed.xml",
"at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/3lab",
] {
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
sqlx::query(
"UPDATE feeds SET consecutive_errors = 35, last_error_kind = 'fetch', \
last_error = 'unsupported scheme' WHERE url = ?1",
)
.bind(url)
.execute(&pool)
.await?;
}
apply_migrations(&pool).await?;
let (at_errors, at_kind, at_detail): (i64, Option<String>, Option<String>) =
sqlx::query_as(sqlx::AssertSqlSafe(format!(
"SELECT consecutive_errors, last_error_kind, last_error FROM feeds \
WHERE kind = '{}'",
crate::feed::FeedKind::Unsupported.as_str()
)))
.fetch_one(&pool)
.await?;
assert_eq!(at_errors, 0, "an unpollable row kept its failure count");
assert_eq!(at_kind, None, "an unpollable row kept its failure kind");
assert_eq!(at_detail, None, "an unpollable row kept its failure detail");
let http_errors: i64 = sqlx::query_scalar(
"SELECT consecutive_errors FROM feeds WHERE url = 'https://real.example/feed.xml'",
)
.fetch_one(&pool)
.await?;
assert_eq!(http_errors, 35, "a real feed's history was discarded");
Ok(())
}
#[tokio::test]
async fn an_at_uri_feed_is_never_due_for_polling() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
for url in [
"https://example.com/feed.xml",
"at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/3lab",
"at://alice.example.com/site.standard.publication/3lab",
] {
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
}
let due = due_feeds(&pool, "2026-09-20T00:00:00Z", 50).await?;
let urls: Vec<&str> = due.iter().map(|f| f.url.as_str()).collect();
assert_eq!(
urls,
["https://example.com/feed.xml"],
"an at:// feed was handed to the poller"
);
Ok(())
}
#[tokio::test]
async fn feed_error_count_bumps_and_resets() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let url = "https://broken.example/feed.xml";
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
let feed = get_feed_by_url(&pool, url).await?.expect("feed exists");
assert_eq!(feed.consecutive_errors, 0);
let mut last = std::time::Duration::ZERO;
for expected in 1..=3 {
let count = bump_feed_errors(
&pool,
url,
crate::feed::FailureKind::Fetch,
"connection refused",
)
.await?;
assert_eq!(count, expected, "bump returns the new count");
let backoff = crate::feed::backoff_for(count as u32);
assert!(
backoff >= last,
"backoff must not shrink as errors accumulate"
);
last = backoff;
}
assert!(crate::feed::backoff_for(2) > crate::feed::backoff_for(1));
assert_eq!(
get_feed_by_url(&pool, url)
.await?
.unwrap()
.consecutive_errors,
3
);
reset_feed_errors(&pool, url).await?;
assert_eq!(
get_feed_by_url(&pool, url)
.await?
.unwrap()
.consecutive_errors,
0
);
Ok(())
}
#[tokio::test]
async fn a_successful_poll_clears_the_recorded_failure_reason() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let url = "https://recovers.example/feed.xml";
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
bump_feed_errors(&pool, url, crate::feed::FailureKind::Fetch, "SENTINEL_WHY").await?;
let failing: (Option<String>, Option<String>) =
sqlx::query_as("SELECT last_error_kind, last_error FROM feeds WHERE url = ?1")
.bind(url)
.fetch_one(&pool)
.await?;
assert_eq!(
failing.0.as_deref(),
Some("fetch"),
"the kind was not stored"
);
assert_eq!(
failing.1.as_deref(),
Some("SENTINEL_WHY"),
"the detail was not stored"
);
reset_feed_errors(&pool, url).await?;
let recovered: (Option<String>, Option<String>) =
sqlx::query_as("SELECT last_error_kind, last_error FROM feeds WHERE url = ?1")
.bind(url)
.fetch_one(&pool)
.await?;
assert_eq!(
recovered.0, None,
"a healthy feed still names a failure kind"
);
assert_eq!(
recovered.1, None,
"a healthy feed still carries error detail"
);
Ok(())
}
#[tokio::test]
async fn an_unrecognised_failure_kind_folds_into_unknown() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
for (url, kind) in [
("https://a.example/f.xml", Some("fetch")),
("https://b.example/f.xml", Some("quota")), ("https://c.example/f.xml", None), ] {
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
sqlx::query(
"UPDATE feeds SET consecutive_errors = 1, last_error_kind = ?2 WHERE url = ?1",
)
.bind(url)
.bind(kind)
.execute(&pool)
.await?;
}
let now = chrono::Utc::now();
let health = poll_health(
&pool,
&now.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
&(now - chrono::Duration::hours(1)).to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
)
.await?;
let mut kinds = health.failure_kinds.clone();
kinds.sort();
assert_eq!(
kinds,
vec![("fetch".to_string(), 1), ("unknown".to_string(), 2)],
"an unrecognised kind reached the public histogram as its own bucket: {:?}",
health.failure_kinds
);
Ok(())
}
#[tokio::test]
async fn the_last_error_columns_migrate_onto_a_table_that_predates_them() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
sqlx::query("DROP TABLE feeds").execute(&pool).await?;
sqlx::query(
"CREATE TABLE feeds (
id INTEGER PRIMARY KEY AUTOINCREMENT,
url TEXT NOT NULL UNIQUE,
title TEXT,
site_url TEXT,
etag TEXT,
last_modified TEXT,
last_polled TEXT,
next_poll TEXT,
consecutive_errors INTEGER NOT NULL DEFAULT 0
)",
)
.execute(&pool)
.await?;
sqlx::query("INSERT INTO feeds (url, consecutive_errors) VALUES (?1, 7)")
.bind("https://legacy.example/feed.xml")
.execute(&pool)
.await?;
apply_migrations(&pool).await?;
let cols: Vec<String> = sqlx::query("PRAGMA table_info(feeds)")
.fetch_all(&pool)
.await?
.iter()
.map(|r| r.get::<String, _>("name"))
.collect();
assert!(cols.iter().any(|c| c == "last_error_kind"), "{cols:?}");
assert!(cols.iter().any(|c| c == "last_error"), "{cols:?}");
let row: (i64, Option<String>, Option<String>) = sqlx::query_as(
"SELECT consecutive_errors, last_error_kind, last_error FROM feeds WHERE url = ?1",
)
.bind("https://legacy.example/feed.xml")
.fetch_one(&pool)
.await?;
assert_eq!(row.0, 7, "the migration disturbed an existing error count");
assert_eq!(row.1, None, "a legacy row was given a cause it never had");
assert_eq!(row.2, None);
apply_migrations(&pool).await?;
Ok(())
}
#[tokio::test]
async fn the_stored_error_detail_is_truncated() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let url = "https://verbose.example/feed.xml";
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
bump_feed_errors(
&pool,
url,
crate::feed::FailureKind::Body,
&"x".repeat(10_000),
)
.await?;
let stored: (Option<String>,) =
sqlx::query_as("SELECT last_error FROM feeds WHERE url = ?1")
.bind(url)
.fetch_one(&pool)
.await?;
assert_eq!(stored.0.unwrap().chars().count(), MAX_ERROR_DETAIL_CHARS);
Ok(())
}
#[tokio::test]
async fn a_new_database_is_created_in_incremental_vacuum_mode() -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-autovac-{}.db", std::process::id()));
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
let pool = init_url(&format!("sqlite://{}", path.display())).await?;
assert_eq!(
auto_vacuum_mode(&pool).await?,
AutoVacuum::Incremental,
"a fresh database is still in the mode where reclaim needs a full VACUUM"
);
let limit: i64 = sqlx::query_scalar("PRAGMA journal_size_limit")
.fetch_one(&pool)
.await?;
assert_eq!(
limit, WAL_SIZE_LIMIT_BYTES,
"journal_size_limit not applied"
);
assert_eq!(
migrate_to_incremental_vacuum(&pool, None).await?,
VacuumMigration::NotNeeded(AutoVacuum::Incremental)
);
pool.close().await;
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
Ok(())
}
#[tokio::test]
async fn the_vacuum_migration_refuses_without_headroom() -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-autovac-none-{}.db", std::process::id()));
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
let url = format!("sqlite://{}", path.display());
let opts = SqliteConnectOptions::from_str(&url)?
.create_if_missing(true)
.foreign_keys(true)
.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
.auto_vacuum(sqlx::sqlite::SqliteAutoVacuum::None);
let pool = SqlitePoolOptions::new()
.min_connections(1)
.max_connections(1)
.connect_with(opts)
.await?;
init_schema(&pool).await?;
assert_eq!(auto_vacuum_mode(&pool).await?, AutoVacuum::None);
let refused = migrate_to_incremental_vacuum(&pool, Some(0)).await?;
assert!(
matches!(refused, VacuumMigration::RefusedNoHeadroom { .. }),
"expected a refusal, got {refused:?}"
);
assert_eq!(
auto_vacuum_mode(&pool).await?,
AutoVacuum::None,
"a refused migration must not have changed the mode"
);
let done = migrate_to_incremental_vacuum(&pool, Some(u64::MAX)).await?;
let VacuumMigration::Migrated {
bytes_after,
file_after,
..
} = done
else {
panic!("expected a migration, got {done:?}");
};
assert_eq!(auto_vacuum_mode(&pool).await?, AutoVacuum::Incremental);
let file_after = file_after.expect("an on-disk database has a file size") as i64;
assert!(
bytes_after <= file_after * 2,
"bytes_after ({bytes_after}) is inflated by an untruncated WAL against a \
{file_after}-byte file"
);
pool.close().await;
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
Ok(())
}
#[tokio::test]
#[ignore]
async fn r6_measure_retention_sweep() -> Result<()> {
const FEEDS: i64 = 500;
const PER_FEED: i64 = 2_000; const PINNED: i64 = 50_000;
async fn run(
label: &str,
extra_indexes: &[&str],
pinned: i64,
old_list_form: bool,
) -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-r6-{}-{label}.db", std::process::id()));
struct Fixture(std::path::PathBuf);
impl Fixture {
fn wipe(&self) {
for p in [
self.0.display().to_string(),
format!("{}-wal", self.0.display()),
format!("{}-shm", self.0.display()),
] {
std::fs::remove_file(&p).ok();
}
}
}
impl Drop for Fixture {
fn drop(&mut self) {
self.wipe();
}
}
let fixture = Fixture(path.clone());
fixture.wipe();
let pool = init_url(&format!("sqlite://{}", path.display())).await?;
sqlx::query(
"WITH RECURSIVE n(value) AS ( \
SELECT 1 UNION ALL SELECT value + 1 FROM n WHERE value < ?1 \
) \
INSERT INTO feeds (url) \
SELECT 'https://f' || value || '.example/x.xml' FROM n",
)
.bind(FEEDS)
.execute(&pool)
.await
.context("seeding feeds")?;
sqlx::query(
"WITH RECURSIVE n(value) AS ( \
SELECT 1 UNION ALL SELECT value + 1 FROM n WHERE value < ?1 \
) \
INSERT INTO entries (feed_id, guid, title, published, fetched_at) \
SELECT f.id, \
'g' || f.id || '-' || s.value, \
'Entry ' || s.value, \
CASE WHEN s.value % 2 = 0 THEN '2020-01-01T00:00:00Z' \
ELSE '2099-01-01T00:00:00Z' END, \
'2026-01-01T00:00:00Z' \
FROM feeds f, n s",
)
.bind(PER_FEED)
.execute(&pool)
.await?;
sqlx::query(
"INSERT INTO entry_state (did, entry_id, read, starred, updated_at) \
SELECT 'did:plc:reader', id, \
CASE WHEN id % 10 <> 0 THEN 1 \
WHEN id % 20 = 0 THEN 1 ELSE 0 END, \
CASE WHEN id % 10 <> 0 THEN 0 \
WHEN id % 20 = 0 THEN 1 ELSE 0 END, \
'2026-01-01T00:00:00Z' \
FROM entries LIMIT ?1",
)
.bind(pinned)
.execute(&pool)
.await?;
let pages_before: i64 = sqlx::query_scalar("PRAGMA page_count")
.fetch_one(&pool)
.await?;
let page_size: i64 = sqlx::query_scalar("PRAGMA page_size")
.fetch_one(&pool)
.await?;
for idx in extra_indexes {
sqlx::query(sqlx::AssertSqlSafe((*idx).to_string()))
.execute(&pool)
.await
.with_context(|| format!("creating {idx}"))?;
}
let pages_after: i64 = sqlx::query_scalar("PRAGMA page_count")
.fetch_one(&pool)
.await?;
let index_bytes = (pages_after - pages_before) * page_size;
let t_ins = std::time::Instant::now();
sqlx::query(
"WITH RECURSIVE n(value) AS ( \
SELECT 1 UNION ALL SELECT value + 1 FROM n WHERE value < 10000 \
) \
INSERT INTO entries (feed_id, guid, published, fetched_at) \
SELECT 1, 'ins-' || value, '2099-06-01T00:00:00Z', '2026-01-01T00:00:00Z' \
FROM n",
)
.execute(&pool)
.await?;
let insert_10k = t_ins.elapsed();
sqlx::query("ANALYZE").execute(&pool).await?;
let planned = if old_list_form {
"EXPLAIN QUERY PLAN SELECT id FROM entries \
WHERE COALESCE(published, fetched_at) < '2026-06-01T00:00:00Z' \
AND id NOT IN (SELECT entry_id FROM entry_state \
WHERE starred = 1 OR read = 0) \
LIMIT 1000"
} else {
"EXPLAIN QUERY PLAN SELECT e.id FROM entries e \
WHERE COALESCE(e.published, e.fetched_at) < '2026-06-01T00:00:00Z' \
AND NOT EXISTS (SELECT 1 FROM entry_state s \
WHERE s.entry_id = e.id \
AND (s.starred = 1 OR s.read = 0)) \
LIMIT 1000"
};
let plan: Vec<String> = sqlx::query(sqlx::AssertSqlSafe(planned))
.fetch_all(&pool)
.await?
.into_iter()
.map(|r| r.get::<String, _>("detail"))
.collect();
let cutoff = (chrono::Utc::now() - chrono::Duration::days(30))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let t = std::time::Instant::now();
let deleted = if old_list_form {
let mut n = 0u64;
loop {
let got = sqlx::query(
"DELETE FROM entries WHERE id IN ( \
SELECT id FROM entries \
WHERE COALESCE(published, fetched_at) < ?1 \
AND id NOT IN ( \
SELECT entry_id FROM entry_state \
WHERE starred = 1 OR read = 0 \
) \
LIMIT 1000)",
)
.bind(&cutoff)
.execute(&pool)
.await?
.rows_affected();
n += got;
if got == 0 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
n
} else {
prune_old_entries(&pool, 30, 3650, 0).await?
};
let elapsed = t.elapsed();
let matching: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM entry_state WHERE starred = 1 OR read = 0",
)
.fetch_one(&pool)
.await?;
println!("\n=== {label} (entry_state = {pinned}, pinned = {matching}) ===");
println!(
" index cost: {:.1} MiB on disk, 10k inserts in {insert_10k:?}",
index_bytes as f64 / 1024.0 / 1024.0
);
for l in &plan {
println!(" plan: {l}");
}
println!(
" deleted {deleted} in {elapsed:?} ({:?}/batch)",
elapsed / (deleted as u32 / PRUNE_BATCH as u32).max(1)
);
pool.close().await;
drop(fixture);
Ok(())
}
const AGE_IDX: &str =
"CREATE INDEX idx_entries_age ON entries(COALESCE(published, fetched_at))";
const PINNED_IDX: &str = "CREATE INDEX idx_es_pinned ON entry_state(entry_id) \
WHERE starred = 1 OR read = 0";
for pinned in [PINNED, 600_000] {
run("as shipped (NOT EXISTS)", &[], pinned, false).await?;
run("old NOT IN list form", &[], pinned, true).await?;
run("old NOT IN + pinned index", &[PINNED_IDX], pinned, true).await?;
run("as shipped + pinned index", &[PINNED_IDX], pinned, false).await?;
run("as shipped + age index", &[AGE_IDX], pinned, false).await?;
}
Ok(())
}
#[tokio::test]
async fn the_vacuum_migration_never_waits_on_its_own_pool() -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-nodeadlock-{}.db", std::process::id()));
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
let url = format!("sqlite://{}", path.display());
let opts = SqliteConnectOptions::from_str(&url)?
.create_if_missing(true)
.foreign_keys(true)
.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
.auto_vacuum(sqlx::sqlite::SqliteAutoVacuum::None);
let pool = SqlitePoolOptions::new()
.min_connections(1)
.max_connections(1)
.connect_with(opts)
.await?;
init_schema(&pool).await?;
let t0 = std::time::Instant::now();
let outcome = migrate_to_incremental_vacuum(&pool, Some(u64::MAX)).await?;
let elapsed = t0.elapsed();
assert!(
matches!(outcome, VacuumMigration::Migrated { .. }),
"expected a migration, got {outcome:?}"
);
assert!(
elapsed < std::time::Duration::from_secs(5),
"the migration took {elapsed:?} on an empty database — it is waiting on \
its own pool while holding a connection"
);
pool.close().await;
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
Ok(())
}
#[tokio::test]
async fn reclaim_does_not_full_vacuum_in_none_mode() -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-noneclaim-{}.db", std::process::id()));
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
let url = format!("sqlite://{}", path.display());
let opts = SqliteConnectOptions::from_str(&url)?
.create_if_missing(true)
.foreign_keys(true)
.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
.auto_vacuum(sqlx::sqlite::SqliteAutoVacuum::None);
let pool = SqlitePoolOptions::new()
.min_connections(1)
.max_connections(1)
.connect_with(opts)
.await?;
init_schema(&pool).await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://none.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
let entries: Vec<NewEntry> = (0..1500)
.map(|i| NewEntry {
guid: format!("n-{i}"),
content_html: Some("x".repeat(800)),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
sqlx::query("PRAGMA wal_checkpoint(TRUNCATE)")
.execute(&pool)
.await?;
let used_full = db_size_bytes(&pool).await?;
sqlx::query("DELETE FROM entries").execute(&pool).await?;
let pages_before: i64 = sqlx::query_scalar("PRAGMA page_count")
.fetch_one(&pool)
.await?;
reclaim(&pool).await?;
let pages_after: i64 = sqlx::query_scalar("PRAGMA page_count")
.fetch_one(&pool)
.await?;
assert_eq!(
pages_after, pages_before,
"reclaim shrank the file in NONE mode, so it ran the full VACUUM this \
branch exists to avoid"
);
let used_after = db_size_bytes(&pool).await?;
assert!(
used_after < used_full,
"used size did not fall after the delete ({used_after} !< {used_full}); \
without a VACUUM the DB-size watermark would latch the poller off"
);
pool.close().await;
for p in [
path.display().to_string(),
format!("{}-wal", path.display()),
format!("{}-shm", path.display()),
] {
std::fs::remove_file(&p).ok();
}
Ok(())
}
#[tokio::test]
async fn db_size_drops_after_prune_and_reclaim() -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-reclaim-{}.db", std::process::id()));
let url = format!("sqlite://{}", path.display());
let pool = init_url(&url).await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://bulk.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let entries: Vec<NewEntry> = (0..2000)
.map(|i| NewEntry {
guid: format!("guid-{i}"),
title: Some(format!("Entry number {i} with some padding text")),
content_html: Some("<p>".to_string() + &"x".repeat(400) + "</p>"),
published: Some("2026-01-01T00:00:00Z".to_string()),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
let full = db_size_bytes(&pool).await?;
assert!(full > 0);
sqlx::query("DELETE FROM entries WHERE feed_id = ?1")
.bind(feed_id)
.execute(&pool)
.await?;
let after_delete = db_size_bytes(&pool).await?;
assert!(
after_delete < full,
"used size must drop once rows are deleted (freed pages excluded): \
{after_delete} !< {full}"
);
reclaim(&pool).await?;
let after_reclaim = db_size_bytes(&pool).await?;
assert!(
after_reclaim <= after_delete,
"reclaim must not grow used size: {after_reclaim} !<= {after_delete}"
);
assert!(
after_reclaim < full,
"after prune+reclaim the DB is smaller than when full: \
{after_reclaim} !< {full}"
);
drop(pool);
let _ = std::fs::remove_file(&path);
let _ = std::fs::remove_file(format!("{}-wal", path.display()));
let _ = std::fs::remove_file(format!("{}-shm", path.display()));
Ok(())
}
#[tokio::test]
async fn cursor_pds_created_defaults_false_and_flips() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:f4";
let feed_url = "https://example.com/feed.xml";
upsert_cursor(
&pool,
&ReadCursor {
did: did.to_string(),
feed_url: feed_url.to_string(),
read_through: None,
read_ids: r#"["1"]"#.to_string(),
unread_ids: "[]".to_string(),
dirty: true,
pds_created: false,
updated_at: now_rfc3339(),
},
)
.await?;
let c = get_cursor(&pool, did, feed_url).await?.unwrap();
assert!(!c.pds_created, "first flush must emit a create, not update");
for (d, f) in [
(did, "https://other.example/feed.xml"),
("did:plc:other", feed_url),
] {
upsert_cursor(
&pool,
&ReadCursor {
did: d.to_string(),
feed_url: f.to_string(),
read_through: None,
read_ids: "[]".to_string(),
unread_ids: "[]".to_string(),
dirty: false,
pds_created: false,
updated_at: now_rfc3339(),
},
)
.await?;
}
mark_cursor_pds_created(&pool, did, feed_url).await?;
let c = get_cursor(&pool, did, feed_url).await?.unwrap();
assert!(c.pds_created);
for (d, f) in [
(did, "https://other.example/feed.xml"),
("did:plc:other", feed_url),
] {
let bystander = get_cursor(&pool, d, f).await?.unwrap();
assert!(
!bystander.pds_created,
"marking ({did}, {feed_url}) also flagged ({d}, {f})"
);
}
Ok(())
}
async fn count_entries(pool: &SqlitePool) -> Result<i64> {
Ok(sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM entries")
.fetch_one(pool)
.await?)
}
#[tokio::test]
async fn the_window_and_the_ceiling_spare_a_publication_but_not_an_rss_entry() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let rss = upsert_feed(
&pool,
&NewFeed {
url: "https://aged.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let publication = upsert_feed(
&pool,
&NewFeed {
url: "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab"
.to_string(),
..Default::default()
},
)
.await?;
let kinds: Vec<String> = sqlx::query_scalar("SELECT kind FROM feeds ORDER BY id")
.fetch_all(&pool)
.await?;
assert_eq!(kinds, vec!["rss".to_string(), "publication".to_string()]);
let ancient = (chrono::Utc::now() - chrono::Duration::days(365))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
for feed_id in [rss, publication] {
insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: format!("ancient-{feed_id}"),
published: Some(ancient.clone()),
fetched_at: Some(ancient.clone()),
..Default::default()
}],
0,
)
.await?;
}
replace_sub_refs(&pool, "did:plc:reader", &[rss, publication]).await?;
for id in sqlx::query_scalar::<_, i64>("SELECT id FROM entries ORDER BY id")
.fetch_all(&pool)
.await?
{
mark_read(&pool, "did:plc:reader", id, true).await?;
}
assert_eq!(count_entries(&pool).await?, 2);
let deleted = prune_old_entries(&pool, 14, 180, 3_650).await?;
assert_eq!(deleted, 1, "exactly one of the two should have gone");
let surviving: Vec<i64> = sqlx::query_scalar("SELECT feed_id FROM entries")
.fetch_all(&pool)
.await?;
assert_eq!(
surviving,
vec![publication],
"the publication's year-old document was swept — under the 14-day \
window that is every document a real publication has, so the feed a \
reader subscribed to would be permanently empty",
);
Ok(())
}
#[tokio::test]
async fn the_archive_ceiling_reaps_a_publication_entry_past_it() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let rss = upsert_feed(
&pool,
&NewFeed {
url: "https://not-swept.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let publication = upsert_feed(
&pool,
&NewFeed {
url: "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab"
.to_string(),
..Default::default()
},
)
.await?;
let ancient = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let recent = now_rfc3339();
insert_entries(
&pool,
publication,
&[
NewEntry {
guid: "past-the-ceiling".into(),
published: Some(ancient.clone()),
fetched_at: Some(ancient),
..Default::default()
},
NewEntry {
guid: "inside-the-ceiling".into(),
published: Some(recent.clone()),
fetched_at: Some(recent),
..Default::default()
},
],
0,
)
.await?;
let long_ago = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
insert_entries(
&pool,
rss,
&[NewEntry {
guid: "rss-past-the-archive-ceiling".into(),
published: Some(long_ago.clone()),
fetched_at: Some(long_ago),
..Default::default()
}],
0,
)
.await?;
replace_sub_refs(&pool, "did:plc:reader", &[rss, publication]).await?;
for id in sqlx::query_scalar::<_, i64>("SELECT id FROM entries ORDER BY id")
.fetch_all(&pool)
.await?
{
mark_starred(&pool, "did:plc:reader", id, true).await?;
}
let deleted = prune_old_entries(&pool, 0, 0, 365).await?;
assert_eq!(
deleted, 1,
"the entry past the archive ceiling was not reaped"
);
let mut guids: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries")
.fetch_all(&pool)
.await?;
guids.sort();
assert_eq!(
guids,
vec![
"inside-the-ceiling".to_string(),
"rss-past-the-archive-ceiling".to_string(),
],
"the archive ceiling must reap the publication's over-age entry and \
ONLY that — an RSS entry on an instance with both RSS knobs at zero \
is one the operator chose to keep",
);
assert_eq!(
prune_old_entries(&pool, 0, 0, 0).await?,
0,
"publication_retention_days = 0 still deleted something",
);
Ok(())
}
#[tokio::test]
async fn an_unrepresentable_retention_window_disables_the_pass_it_belongs_to() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://absurd.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let ancient = (chrono::Utc::now() - chrono::Duration::days(1_000))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: "ancient".into(),
published: Some(ancient.clone()),
fetched_at: Some(ancient),
..Default::default()
}],
0,
)
.await?;
let absurd = u32::MAX as i64;
assert_eq!(
prune_old_entries(&pool, absurd, 0, 0).await?,
0,
"an absurd rolling window deleted something",
);
assert_eq!(
prune_old_entries(&pool, 0, absurd, 0).await?,
0,
"an absurd hard ceiling deleted something",
);
assert_eq!(
prune_old_entries(&pool, 0, 0, absurd).await?,
0,
"an absurd archive ceiling deleted something",
);
assert_eq!(
count_entries(&pool).await?,
1,
"the entry went away under a window that cannot even be expressed",
);
assert_eq!(
prune_old_entries(&pool, 30, 0, 0).await?,
1,
"a 30-day window did not delete a 1000-day-old entry",
);
Ok(())
}
#[test]
fn the_sql_aged_kind_list_matches_the_rust_one() {
let expected = crate::feed::FeedKind::AGED
.iter()
.map(|k| format!("'{}'", k.as_str()))
.collect::<Vec<_>>()
.join(", ");
assert_eq!(AGED_KINDS_SQL, expected);
}
#[tokio::test]
async fn prune_old_entries_deletes_only_old_and_cascades_entry_state() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://ret.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let recent = now_rfc3339();
let ancient = (chrono::Utc::now() - chrono::Duration::days(365))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "fresh".into(),
published: Some(recent.clone()),
fetched_at: Some(recent.clone()),
..Default::default()
},
NewEntry {
guid: "ancient".into(),
published: Some(ancient.clone()),
fetched_at: Some(ancient.clone()),
..Default::default()
},
NewEntry {
guid: "undated-fresh".into(),
published: None,
fetched_at: Some(recent.clone()),
..Default::default()
},
],
0,
)
.await?;
assert_eq!(count_entries(&pool).await?, 3);
replace_sub_refs(&pool, "did:plc:reader", &[feed_id]).await?;
let ancient_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'ancient'")
.fetch_one(&pool)
.await?;
let wrote = mark_read(&pool, "did:plc:reader", ancient_id, true).await?;
assert!(wrote, "mark_read must write with a sub_ref in place");
let state_before: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM entry_state WHERE entry_id = ?1")
.bind(ancient_id)
.fetch_one(&pool)
.await?;
assert_eq!(state_before, 1);
let deleted = prune_old_entries(&pool, 90, 3650, 0).await?;
assert_eq!(deleted, 1, "only the year-old entry should be pruned");
assert_eq!(
count_entries(&pool).await?,
2,
"fresh + undated-fresh survive"
);
let surviving: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
.fetch_all(&pool)
.await?;
assert_eq!(surviving, vec!["fresh", "undated-fresh"]);
let state_after: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM entry_state WHERE entry_id = ?1")
.bind(ancient_id)
.fetch_one(&pool)
.await?;
assert_eq!(state_after, 0, "entry_state must cascade on entry delete");
assert_eq!(prune_old_entries(&pool, 0, 3650, 0).await?, 0);
assert_eq!(count_entries(&pool).await?, 2);
Ok(())
}
#[tokio::test]
async fn prune_removes_orphan_ids_from_read_cursor() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:reader";
let feed_url = "https://orphan.example/feed.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
..Default::default()
},
)
.await?;
replace_sub_refs(&pool, did, &[feed_id]).await?;
let ancient = (chrono::Utc::now() - chrono::Duration::days(365))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let recent = now_rfc3339();
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "old".into(),
published: Some(ancient.clone()),
fetched_at: Some(ancient.clone()),
..Default::default()
},
NewEntry {
guid: "new".into(),
published: Some(recent.clone()),
fetched_at: Some(recent.clone()),
..Default::default()
},
],
0,
)
.await?;
let old_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'old'")
.fetch_one(&pool)
.await?;
let new_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'new'")
.fetch_one(&pool)
.await?;
mark_read(&pool, did, old_id, true).await?;
mark_read(&pool, did, new_id, true).await?;
let before = get_cursor(&pool, did, feed_url).await?.unwrap();
let ids_before: Vec<String> = serde_json::from_str(&before.read_ids)?;
assert!(ids_before.contains(&old_id.to_string()));
assert!(ids_before.contains(&new_id.to_string()));
let deleted = prune_old_entries(&pool, 90, 3650, 0).await?;
assert_eq!(deleted, 1);
let after = get_cursor(&pool, did, feed_url).await?.unwrap();
let ids_after: Vec<String> = serde_json::from_str(&after.read_ids)?;
assert_eq!(
ids_after,
vec![new_id.to_string()],
"orphaned (deleted) entry id must be removed; live id kept"
);
assert!(
after.dirty,
"cursor must be marked dirty after orphan scrub"
);
Ok(())
}
#[tokio::test]
async fn insert_entries_trim_scrubs_orphan_cursor_ids() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:reader";
let feed_url = "https://trim.example/feed.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
..Default::default()
},
)
.await?;
replace_sub_refs(&pool, did, &[feed_id]).await?;
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "a".into(),
published: Some("2026-01-01T00:00:00Z".into()),
..Default::default()
},
NewEntry {
guid: "b".into(),
published: Some("2026-01-02T00:00:00Z".into()),
..Default::default()
},
],
2,
)
.await?;
let a_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'a'")
.fetch_one(&pool)
.await?;
mark_read(&pool, did, a_id, true).await?;
insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: "c".into(),
published: Some("2026-01-03T00:00:00Z".into()),
..Default::default()
}],
1,
)
.await?;
let a_still: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries WHERE guid = 'a'")
.fetch_one(&pool)
.await?;
assert_eq!(a_still, 0, "oldest entry trimmed by the per-feed cap");
let cursor = get_cursor(&pool, did, feed_url).await?.unwrap();
let ids: Vec<String> = serde_json::from_str(&cursor.read_ids)?;
assert!(
!ids.contains(&a_id.to_string()),
"trimmed entry id must be scrubbed from the cursor"
);
Ok(())
}
#[tokio::test]
async fn a_sweep_larger_than_one_batch_still_drains() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://bulk.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
let old = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let count = (PRUNE_BATCH * 2 + 137) as usize;
let entries: Vec<NewEntry> = (0..count)
.map(|i| NewEntry {
guid: format!("bulk-{i}"),
published: Some(old.clone()),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
assert_eq!(count_entries(&pool).await? as usize, count);
let deleted = prune_old_entries(&pool, 30, 180, 0).await?;
assert_eq!(deleted as usize, count, "the sweep left rows behind");
assert_eq!(count_entries(&pool).await?, 0);
Ok(())
}
enum SweepOutcome {
Completed(u64),
Contended,
}
fn is_sqlite_busy(err: &anyhow::Error) -> bool {
err.chain().any(|e| {
e.downcast_ref::<sqlx::Error>().is_some_and(|e| match e {
sqlx::Error::Database(db) => db
.code()
.and_then(|c| c.parse::<i32>().ok())
.is_some_and(is_busy_code),
_ => false,
})
})
}
fn is_busy_code(code: i32) -> bool {
code & 0xFF == 5
}
#[test]
fn busy_codes_cover_the_wal_variants() {
for code in [
5, 261, 517, 773, ] {
assert!(
is_busy_code(code),
"{code} is in the BUSY family but would be treated as a hard failure"
);
}
for code in [
0, 1, 6, 262, 11, ] {
assert!(
!is_busy_code(code),
"{code} is not contention, but would be swallowed as though it were"
);
}
}
async fn sweep_tolerating_busy(
pool: &SqlitePool,
select_ids: &str,
cutoff: &str,
label: &str,
) -> Result<SweepOutcome> {
match delete_in_batches(pool, select_ids, cutoff, label).await {
Ok(n) => Ok(SweepOutcome::Completed(n)),
Err(err) if is_sqlite_busy(&err) => Ok(SweepOutcome::Contended),
Err(err) => Err(err),
}
}
#[tokio::test]
async fn a_busy_sweep_is_reported_as_contended_not_as_a_failure() -> Result<()> {
struct TempDb(std::path::PathBuf);
impl Drop for TempDb {
fn drop(&mut self) {
for suffix in ["", "-wal", "-shm"] {
std::fs::remove_file(format!("{}{suffix}", self.0.display())).ok();
}
}
}
let path = std::env::temp_dir().join(format!("fr-busysweep-{}.db", std::process::id()));
drop(TempDb(path.clone()));
let _tmp = TempDb(path.clone());
let url = format!("sqlite://{}", path.display());
let pool = init_url(&url).await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://busy.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
let old = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let entries: Vec<NewEntry> = (0..4)
.map(|i| NewEntry {
guid: format!("busy-{i}"),
url: Some(format!("https://busy.example/{i}")),
title: Some(format!("e{i}")),
published: Some(old.clone()),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 1_000).await?;
let sweep_pool = SqlitePoolOptions::new()
.max_connections(1)
.connect_with(
url.parse::<sqlx::sqlite::SqliteConnectOptions>()?
.busy_timeout(std::time::Duration::from_millis(2)),
)
.await?;
let mut blocker = pool.acquire().await?;
sqlx::query("BEGIN IMMEDIATE")
.execute(&mut *blocker)
.await?;
let cutoff = (chrono::Utc::now() - chrono::Duration::days(180))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let outcome = sweep_tolerating_busy(
&sweep_pool,
"SELECT id FROM entries WHERE COALESCE(published, fetched_at) < ?1",
&cutoff,
"busy-sweep-test",
)
.await;
sqlx::query("ROLLBACK").execute(&mut *blocker).await.ok();
match outcome {
Ok(SweepOutcome::Contended) => Ok(()),
Ok(SweepOutcome::Completed(n)) => panic!(
"the sweep completed ({n} rows) while the write lock was held — \
the fixture is not actually contending, so this test proves nothing"
),
Err(err) => panic!(
"a contended sweep was reported as a failure rather than as \
inconclusive; production logs this and retries on the next \
tick (scheduler.rs): {err:#}"
),
}
}
#[tokio::test]
async fn a_non_busy_sweep_error_still_fails() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let outcome = sweep_tolerating_busy(
&pool,
"SELECT id FROM no_such_table WHERE created < ?1",
"2026-01-01T00:00:00Z",
"bad-query-test",
)
.await;
match outcome {
Err(err) => {
assert!(
!is_sqlite_busy(&err),
"fixture drifted: this must be a non-BUSY error, got {err:#}"
);
Ok(())
}
Ok(SweepOutcome::Contended) => panic!(
"a malformed sweep was reported as lock contention — the \
tolerance is a blanket error swallow, not a narrowing"
),
Ok(SweepOutcome::Completed(n)) => {
panic!("a sweep over a missing table reported {n} rows deleted")
}
}
}
#[tokio::test]
async fn a_sweep_with_no_contention_completes() -> Result<()> {
struct TempDb(std::path::PathBuf);
impl Drop for TempDb {
fn drop(&mut self) {
for suffix in ["", "-wal", "-shm"] {
std::fs::remove_file(format!("{}{suffix}", self.0.display())).ok();
}
}
}
let path = std::env::temp_dir().join(format!("fr-calmsweep-{}.db", std::process::id()));
drop(TempDb(path.clone()));
let _tmp = TempDb(path.clone());
let pool = init_url(&format!("sqlite://{}", path.display())).await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://calm.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
let old = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let entries: Vec<NewEntry> = (0..3)
.map(|i| NewEntry {
guid: format!("calm-{i}"),
published: Some(old.clone()),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 1_000).await?;
let cutoff = (chrono::Utc::now() - chrono::Duration::days(180))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
match sweep_tolerating_busy(
&pool,
"SELECT id FROM entries WHERE COALESCE(published, fetched_at) < ?1",
&cutoff,
"calm-sweep-test",
)
.await?
{
SweepOutcome::Completed(n) => {
assert_eq!(n, 3, "the uncontended sweep did not delete the fixture");
Ok(())
}
SweepOutcome::Contended => panic!(
"nothing was holding the write lock, yet the sweep reported \
contention — every test that skips on `Contended` is now \
skipping unconditionally"
),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_writer_gets_through_while_the_sweep_runs() -> Result<()> {
struct TempDb(std::path::PathBuf);
impl Drop for TempDb {
fn drop(&mut self) {
for suffix in ["", "-wal", "-shm"] {
std::fs::remove_file(format!("{}{suffix}", self.0.display())).ok();
}
}
}
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-sweeplock-{}.db", std::process::id()));
drop(TempDb(path.clone()));
let _tmp = TempDb(path.clone());
let url = format!("sqlite://{}", path.display());
let pool = init_url(&url).await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://lock.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
const BATCHES: i64 = 10;
const _: () = assert!(
BATCHES >= 5,
"the 50% ceiling assumes ~1/BATCHES; below 5 batches correct code false-fails",
);
let old = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let entries: Vec<NewEntry> = (0..(PRUNE_BATCH * BATCHES) as usize)
.map(|i| NewEntry {
guid: format!("lock-{i}"),
published: Some(old.clone()),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
const WRITER_BUSY_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(2);
let writer_pool = SqlitePoolOptions::new()
.min_connections(1)
.max_connections(1)
.connect_with(
SqliteConnectOptions::from_str(&url)?
.foreign_keys(true)
.busy_timeout(WRITER_BUSY_TIMEOUT)
.log_statements(tracing::log::LevelFilter::Debug),
)
.await?;
let done = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let writer_done = std::sync::Arc::clone(&done);
let writer = tokio::spawn(async move {
let mut outcomes: Vec<(std::time::Instant, bool)> = Vec::new();
let mut first_err: Option<(std::time::Instant, String)> = None;
while !writer_done.load(std::sync::atomic::Ordering::Relaxed) {
let started = std::time::Instant::now();
match grant_access(
&writer_pool,
&format!("did:plc:writer{}", outcomes.len()),
None,
"sweep-test",
None,
)
.await
{
Ok(()) => outcomes.push((started, true)),
Err(err) => {
outcomes.push((started, false));
if first_err.is_none() {
first_err = Some((started, format!("{err:#}")));
}
}
}
tokio::task::yield_now().await;
}
writer_pool.close().await;
(outcomes, first_err)
});
let hard_cutoff = (chrono::Utc::now() - chrono::Duration::days(180))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let t0 = std::time::Instant::now();
let outcome = sweep_tolerating_busy(
&pool,
"SELECT id FROM entries WHERE COALESCE(published, fetched_at) < ?1",
&hard_cutoff,
"sweep-lock-test",
)
.await?;
let sweep = t0.elapsed();
done.store(true, std::sync::atomic::Ordering::Relaxed);
let (outcomes, first_err) = writer.await?;
let deleted = match outcome {
SweepOutcome::Completed(n) => n,
SweepOutcome::Contended => {
eprintln!(
"sweep-lock test INCONCLUSIVE: the sweeper took SQLITE_BUSY; \
the hand-off assertions did not run"
);
return Ok(());
}
};
let sweep_end = t0 + sweep;
let inside: Vec<bool> = outcomes
.iter()
.filter(|(t, _)| *t >= t0 && *t < sweep_end)
.map(|(_, ok)| *ok)
.collect();
let attempts = inside.len();
let during = inside.iter().filter(|ok| **ok).count();
let max_refused_run = {
let (mut worst, mut run) = (0usize, 0usize);
for ok in &inside {
run = if *ok { 0 } else { run + 1 };
worst = worst.max(run);
}
worst
};
let why = first_err
.filter(|(t, _)| *t >= t0 && *t < sweep_end)
.map(|(_, e)| format!(" (first in-window write error: {e})"))
.unwrap_or_default();
assert_eq!(deleted as usize, entries.len());
let handoff_floor = PRUNE_BATCH_HANDOFF * BATCHES as u32;
assert!(
sweep > handoff_floor,
"the delete loop finished in {sweep:?}, under the {handoff_floor:?} that \
{BATCHES} batches of `PRUNE_BATCH_HANDOFF` alone would take — it is not \
handing the write lock over between batches at all"
);
assert!(
attempts >= 20,
"the writer only got {attempts} attempts inside a {sweep:?} delete loop \
— too few for the ratio below to mean anything. That is USUALLY an \
invalid measurement rather than a held lock, but note that a loop \
holding the lock throughout is itself one cause of a starved writer, \
so check {during} (landed) before concluding the test is at \
fault{why}"
);
assert!(
max_refused_run * 2 < attempts,
"the delete loop refused {max_refused_run} consecutive write attempts out \
of {attempts} ({during} landed) — a loop that hands the write lock over \
between batches refuses at most about one batch's worth in a row; one \
that holds the lock across them refuses nearly every attempt it sees{why}"
);
pool.close().await;
Ok(())
}
#[tokio::test]
async fn sparing_honours_every_did_not_just_one() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://shared.example/f.xml".to_string(),
..Default::default()
},
)
.await?;
let old = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let guids = [
"nobody-touched", "both-read-unstarred", "one-starred", "one-unread", ];
let entries: Vec<NewEntry> = guids
.iter()
.map(|g| NewEntry {
guid: (*g).to_string(),
published: Some(old.clone()),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
let id_of = |g: &'static str| {
let pool = pool.clone();
async move {
sqlx::query_scalar::<_, i64>("SELECT id FROM entries WHERE guid = ?1")
.bind(g)
.fetch_one(&pool)
.await
.unwrap()
}
};
let rows: [(&str, &'static str, i64, i64); 6] = [
("did:plc:a", "both-read-unstarred", 1, 0),
("did:plc:b", "both-read-unstarred", 1, 0),
("did:plc:a", "one-starred", 1, 0),
("did:plc:b", "one-starred", 1, 1),
("did:plc:a", "one-unread", 1, 0),
("did:plc:b", "one-unread", 0, 0),
];
for (did, guid, read, starred) in rows {
let id = id_of(guid).await;
sqlx::query(
"INSERT INTO entry_state (did, entry_id, read, starred, updated_at) \
VALUES (?1, ?2, ?3, ?4, '2026-01-01T00:00:00Z')",
)
.bind(did)
.bind(id)
.bind(read)
.bind(starred)
.execute(&pool)
.await?;
}
let deleted = prune_old_entries(&pool, 30, 0, 0).await?;
assert_eq!(
deleted, 2,
"expected the untouched and the all-read entries to go"
);
let left: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
.fetch_all(&pool)
.await?;
assert_eq!(
left,
vec!["one-starred".to_string(), "one-unread".to_string()],
"a second reader's star or unread mark must spare the SHARED entry"
);
Ok(())
}
#[tokio::test]
async fn the_cursor_scrub_drops_orphans_and_keeps_live_ids() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:race";
let feed_url = "https://race.example/f.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
..Default::default()
},
)
.await?;
insert_entries(
&pool,
feed_id,
&[
NewEntry {
guid: "live".to_string(),
..Default::default()
},
NewEntry {
guid: "doomed".to_string(),
..Default::default()
},
],
0,
)
.await?;
replace_sub_refs(&pool, did, &[feed_id]).await?;
let live_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'live'")
.fetch_one(&pool)
.await?;
let doomed_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'doomed'")
.fetch_one(&pool)
.await?;
upsert_cursor(
&pool,
&ReadCursor {
did: did.to_string(),
feed_url: feed_url.to_string(),
read_through: None,
read_ids: format!("[\"{doomed_id}\"]"),
unread_ids: "[]".to_string(),
dirty: false,
pds_created: false,
updated_at: now_rfc3339(),
},
)
.await?;
sqlx::query("DELETE FROM entries WHERE guid = 'doomed'")
.execute(&pool)
.await?;
mark_read(&pool, did, live_id, true).await?;
assert_eq!(prune_orphan_cursor_ids(&pool, None).await?, 1);
let cursor = get_cursor(&pool, did, feed_url).await?.expect("cursor");
let ids: Vec<String> = serde_json::from_str(&cursor.read_ids)?;
assert_eq!(
ids,
vec![live_id.to_string()],
"the scrub dropped a live id"
);
assert!(
!ids.contains(&doomed_id.to_string()),
"the orphaned id survived the scrub"
);
Ok(())
}
#[tokio::test]
async fn the_sweep_still_scrubs_orphaned_cursor_ids() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let did = "did:plc:scrub";
let feed_url = "https://scrub.example/f.xml";
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: feed_url.to_string(),
..Default::default()
},
)
.await?;
let old = (chrono::Utc::now() - chrono::Duration::days(400))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
insert_entries(
&pool,
feed_id,
&[NewEntry {
guid: "doomed".to_string(),
published: Some(old),
..Default::default()
}],
0,
)
.await?;
let doomed = entries_for_feed(&pool, did, feed_id).await;
drop(doomed);
let doomed_id: i64 = sqlx::query_scalar("SELECT id FROM entries WHERE guid = 'doomed'")
.fetch_one(&pool)
.await?;
upsert_cursor(
&pool,
&ReadCursor {
did: did.to_string(),
feed_url: feed_url.to_string(),
read_through: None,
read_ids: format!("[\"{doomed_id}\"]"),
unread_ids: "[]".to_string(),
dirty: false,
pds_created: false,
updated_at: now_rfc3339(),
},
)
.await?;
assert_eq!(prune_old_entries(&pool, 30, 180, 0).await?, 1);
let cursor = get_cursor(&pool, did, feed_url).await?.expect("cursor");
let ids: Vec<String> = serde_json::from_str(&cursor.read_ids)?;
assert!(
ids.is_empty(),
"the deleted entry's id survived in the cursor: {ids:?}"
);
assert!(cursor.dirty, "a rewritten cursor must be re-flushed");
Ok(())
}
#[tokio::test]
async fn prune_and_reclaim_drops_db_size() -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-prune-{}.db", std::process::id()));
let url = format!("sqlite://{}", path.display());
let pool = init_url(&url).await?;
let feed_id = upsert_feed(
&pool,
&NewFeed {
url: "https://bulk.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
let ancient = (chrono::Utc::now() - chrono::Duration::days(365))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let entries: Vec<NewEntry> = (0..2000)
.map(|i| NewEntry {
guid: format!("guid-{i}"),
content_html: Some("<p>".to_string() + &"x".repeat(400) + "</p>"),
published: Some(ancient.clone()),
fetched_at: Some(ancient.clone()),
..Default::default()
})
.collect();
insert_entries(&pool, feed_id, &entries, 0).await?;
let full = db_size_bytes(&pool).await?;
assert!(full > 0);
let deleted = prune_old_entries(&pool, 90, 3650, 0).await?;
assert_eq!(deleted, 2000);
reclaim(&pool).await?;
let after = db_size_bytes(&pool).await?;
assert!(
after < full,
"prune + reclaim must shrink db_size_bytes: {after} !< {full}"
);
drop(pool);
let _ = std::fs::remove_file(&path);
let _ = std::fs::remove_file(format!("{}-wal", path.display()));
let _ = std::fs::remove_file(format!("{}-shm", path.display()));
Ok(())
}
#[tokio::test]
async fn migrates_pre_intended_did_invite_codes_table() -> Result<()> {
let dir = std::env::temp_dir();
let path = dir.join(format!("fr-b1-{}.db", std::process::id()));
let url = format!("sqlite://{}", path.display());
let opts = SqliteConnectOptions::from_str(&url)?
.create_if_missing(true)
.foreign_keys(true);
let pool = SqlitePoolOptions::new()
.min_connections(1)
.max_connections(1)
.connect_with(opts)
.await?;
sqlx::query(
r#"CREATE TABLE invite_codes (
code TEXT PRIMARY KEY,
creator_did TEXT NOT NULL,
status TEXT NOT NULL,
invitee_did TEXT,
created_at INTEGER NOT NULL,
expires_at INTEGER NOT NULL,
redeemed_at INTEGER
);"#,
)
.execute(&pool)
.await?;
sqlx::query(
"INSERT INTO invite_codes (code, creator_did, status, created_at, expires_at) \
VALUES ('FEATHER-LEGACY00', 'did:plc:old', 'active', 1, 9999999999)",
)
.execute(&pool)
.await?;
let cols: Vec<String> = sqlx::query("PRAGMA table_info(invite_codes)")
.fetch_all(&pool)
.await?
.iter()
.map(|r| r.get::<String, _>("name"))
.collect();
assert!(
!cols.iter().any(|c| c == "intended_did"),
"pre-condition: legacy table must lack intended_did"
);
init_schema(&pool)
.await
.expect("init_schema on a pre-0.2.2 invite_codes table must not crash");
let cols: Vec<String> = sqlx::query("PRAGMA table_info(invite_codes)")
.fetch_all(&pool)
.await?
.iter()
.map(|r| r.get::<String, _>("name"))
.collect();
assert!(cols.iter().any(|c| c == "intended_did"));
let idx: Vec<String> = sqlx::query(
"SELECT name FROM sqlite_master WHERE type='index' AND tbl_name='invite_codes'",
)
.fetch_all(&pool)
.await?
.iter()
.map(|r| r.get::<String, _>("name"))
.collect();
assert!(idx.iter().any(|n| n == "idx_invite_codes_intended"));
assert!(idx.iter().any(|n| n == "idx_invite_codes_intended_active"));
init_schema(&pool)
.await
.expect("re-running init_schema must be idempotent");
let tables: Vec<String> =
sqlx::query_scalar("SELECT name FROM sqlite_master WHERE type = 'table'")
.fetch_all(&pool)
.await
.unwrap();
for table in ["oauth_state", "oauth_session", "oauth_nonce"] {
assert!(
tables.iter().any(|t| t == table),
"{table} is missing, so the rust backend would fail on its first request: {tables:?}"
);
}
let out = redeem_code(&pool, "FEATHER-LEGACY00", "did:plc:new", None, 100).await?;
assert_eq!(out, Ok(()));
drop(pool);
let _ = std::fs::remove_file(&path);
let _ = std::fs::remove_file(format!("{}-wal", path.display()));
let _ = std::fs::remove_file(format!("{}-shm", path.display()));
Ok(())
}
async fn upgrade_test_pool() -> Result<SqlitePool> {
let opts = SqliteConnectOptions::from_str("sqlite::memory:")?.foreign_keys(true);
Ok(SqlitePoolOptions::new()
.min_connections(1)
.max_connections(1)
.idle_timeout(None)
.max_lifetime(None)
.connect_with(opts)
.await?)
}
async fn schema_shape(pool: &SqlitePool) -> Result<std::collections::BTreeSet<String>> {
let mut shape = std::collections::BTreeSet::new();
let tables: Vec<String> = sqlx::query_scalar(
"SELECT name FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%'",
)
.fetch_all(pool)
.await?;
for t in tables {
for r in sqlx::query(
r#"SELECT name, type, "notnull", dflt_value, pk FROM pragma_table_info(?)"#,
)
.bind(&t)
.fetch_all(pool)
.await?
{
shape.insert(format!(
"column {t}.{} {} notnull={} default={:?} pk={}",
r.get::<String, _>("name"),
r.get::<String, _>("type"),
r.get::<i64, _>("notnull"),
r.get::<Option<String>, _>("dflt_value"),
r.get::<i64, _>("pk"),
));
}
for r in sqlx::query(r#"SELECT name, "unique", partial FROM pragma_index_list(?)"#)
.bind(&t)
.fetch_all(pool)
.await?
{
let name: String = r.get("name");
let cols: Vec<String> =
sqlx::query_scalar("SELECT name FROM pragma_index_info(?) ORDER BY seqno")
.bind(&name)
.fetch_all(pool)
.await?;
shape.insert(format!(
"index {t}.{name} unique={} partial={} ({})",
r.get::<i64, _>("unique"),
r.get::<i64, _>("partial"),
cols.join(", "),
));
}
}
Ok(shape)
}
async fn assert_upgrades_from(version: &str, fixture: &'static str) -> Result<()> {
let pool = upgrade_test_pool().await?;
sqlx::raw_sql(fixture).execute(&pool).await?;
let has_kind = |pool: SqlitePool| async move {
Ok::<_, anyhow::Error>(
sqlx::query_scalar::<_, i64>(
"SELECT count(*) FROM pragma_table_info('feeds') WHERE name = 'kind'",
)
.fetch_one(&pool)
.await?
== 1,
)
};
assert!(
!has_kind(pool.clone()).await?,
"pre-condition: a {version} feeds table has no kind column"
);
let publication = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
for u in ["https://example.com/feed.xml", publication] {
sqlx::query("INSERT INTO feeds (url) VALUES (?)")
.bind(u)
.execute(&pool)
.await?;
}
init_schema(&pool)
.await
.unwrap_or_else(|e| panic!("init_schema must upgrade a {version} database: {e:#}"));
let kinds: Vec<(String, String)> =
sqlx::query_as("SELECT url, kind FROM feeds ORDER BY id")
.fetch_all(&pool)
.await?;
assert_eq!(
kinds,
vec![
(
"https://example.com/feed.xml".to_string(),
"rss".to_string()
),
(publication.to_string(), "publication".to_string()),
],
"{version}: existing rows are back-filled from their URL"
);
let indexed: Vec<String> = sqlx::query_scalar(
"SELECT name FROM pragma_index_info('idx_feeds_kind') ORDER BY seqno",
)
.fetch_all(&pool)
.await?;
assert_eq!(
indexed,
vec!["kind".to_string()],
"{version}: idx_feeds_kind exists, on feeds(kind), after the column"
);
let fresh = upgrade_test_pool().await?;
init_schema(&fresh).await?;
let (want, got) = (schema_shape(&fresh).await?, schema_shape(&pool).await?);
assert!(
want == got,
"{version}: upgraded schema differs from a fresh one\n missing: {:#?}\n extra: {:#?}",
want.difference(&got).collect::<Vec<_>>(),
got.difference(&want).collect::<Vec<_>>(),
);
init_schema(&pool)
.await
.unwrap_or_else(|e| panic!("{version}: re-running init_schema failed: {e:#}"));
Ok(())
}
#[tokio::test]
async fn a_v0_3_8_database_upgrades_to_the_current_schema() -> Result<()> {
assert_upgrades_from(
"v0.3.8",
include_str!("../tests/fixtures/schema-v0.3.8.sql"),
)
.await
}
#[tokio::test]
async fn a_v0_2_0_database_upgrades_to_the_current_schema() -> Result<()> {
assert_upgrades_from(
"v0.2.0",
include_str!("../tests/fixtures/schema-v0.2.0.sql"),
)
.await
}
#[tokio::test]
async fn redeem_enforces_intended_did_binding() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let code = mint_code_for_did(&pool, "did:bot:fr", 3600, "did:plc:A").await?;
let stolen = redeem_code(&pool, &code, "did:plc:B", Some("thief.bsky"), 100).await?;
assert_eq!(stolen, Err(RedeemError::NotFound));
assert!(!has_beta_access(&pool, "did:plc:B").await?);
assert_eq!(count_active_codes(&pool).await?, 1, "code must stay active");
let ok = redeem_code(&pool, &code, "did:plc:A", Some("alice.bsky"), 100).await?;
assert_eq!(ok, Ok(()));
assert!(has_beta_access(&pool, "did:plc:A").await?);
let open = mint_code(&pool, "did:plc:admin", 3600).await?;
let anyone = redeem_code(&pool, &open, "did:plc:C", None, 100).await?;
assert_eq!(anyone, Ok(()));
assert!(has_beta_access(&pool, "did:plc:C").await?);
Ok(())
}
#[tokio::test]
async fn intended_active_partial_unique_index_blocks_double_mint() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
mint_code_for_did(&pool, "did:bot:fr", 3600, "did:plc:dup").await?;
let err = mint_code_for_did(&pool, "did:bot:fr", 3600, "did:plc:dup")
.await
.expect_err("second active mint for the same DID must fail the unique index");
assert!(
is_intended_active_conflict(&err),
"the conflict must be recognised so the web layer can recover: {err:?}"
);
assert!(find_active_code_for_did(&pool, "did:plc:dup")
.await?
.is_some());
let existing = find_active_code_for_did(&pool, "did:plc:dup")
.await?
.unwrap();
redeem_code(&pool, &existing, "did:plc:dup", None, 100).await??;
mint_code_for_did(&pool, "did:bot:fr", 3600, "did:plc:dup")
.await
.expect("a new mint is allowed after the prior one is redeemed");
sqlx::query(
"INSERT INTO invite_codes (code, creator_did, status, created_at, expires_at) \
VALUES ('FEATHER-DUPEKEY0', 'did:x', 'active', 1, 9999999999)",
)
.execute(&pool)
.await?;
let pk_err = sqlx::query(
"INSERT INTO invite_codes (code, creator_did, status, created_at, expires_at) \
VALUES ('FEATHER-DUPEKEY0', 'did:x', 'active', 1, 9999999999)",
)
.execute(&pool)
.await
.expect_err("duplicate PRIMARY KEY must error");
let as_anyhow = anyhow::Error::new(pk_err);
assert!(
!is_intended_active_conflict(&as_anyhow),
"a non-intended-index conflict must NOT be mistaken for the recover-able one"
);
Ok(())
}
#[tokio::test]
async fn purge_expires_orphaned_active_intended_code() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let code = mint_code_for_did(&pool, "did:bot:fr", 3600, "did:plc:leaver").await?;
assert_eq!(count_active_codes(&pool).await?, 1);
purge_did_data(&pool, "did:plc:leaver").await?;
assert_eq!(
count_active_codes(&pool).await?,
0,
"orphaned code must be expired by purge, not left active"
);
let intended: Option<String> =
sqlx::query("SELECT intended_did FROM invite_codes WHERE code = ?1")
.bind(&code)
.fetch_one(&pool)
.await?
.get("intended_did");
assert!(intended.is_none(), "intended_did must be NULLed");
Ok(())
}
#[tokio::test]
async fn network_stat_upserts_per_source() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let mut stat = NetworkStat {
key: ADOPTION_STAT_KEY.to_string(),
source: "https://relay1.us-west.bsky.network".to_string(),
value: 1,
truncated: false,
observed_at: "2026-08-12T00:00:00Z".to_string(),
};
record_network_stat(&pool, &stat).await?;
stat.value = 4;
stat.observed_at = "2026-08-13T00:00:00Z".to_string();
record_network_stat(&pool, &stat).await?;
let rows: i64 = sqlx::query("SELECT COUNT(*) AS n FROM network_stat")
.fetch_one(&pool)
.await?
.get("n");
assert_eq!(rows, 1, "the same relay must update, not duplicate");
let latest = latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.expect("a stat");
assert_eq!(latest.value, 4);
assert_eq!(latest.observed_at, "2026-08-13T00:00:00Z");
Ok(())
}
#[tokio::test]
async fn a_truncated_observation_never_lowers_a_stored_count() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
let mut stat = NetworkStat {
key: ADOPTION_STAT_KEY.to_string(),
source: "https://relay1.us-west.bsky.network".to_string(),
value: 2000,
truncated: false,
observed_at: "2026-08-13T00:00:00Z".to_string(),
};
record_network_stat(&pool, &stat).await?;
stat.value = 500;
stat.truncated = true;
stat.observed_at = "2026-08-14T00:00:00Z".to_string();
record_network_stat(&pool, &stat).await?;
let kept = latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.expect("a stat");
assert_eq!(kept.value, 2000, "a partial walk must not lower the count");
assert!(!kept.truncated, "and must not mark the kept row truncated");
assert_eq!(kept.observed_at, "2026-08-13T00:00:00Z");
stat.value = 3000;
record_network_stat(&pool, &stat).await?;
assert_eq!(
latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.expect("a stat")
.value,
3000
);
stat.value = 42;
stat.truncated = false;
record_network_stat(&pool, &stat).await?;
assert_eq!(
latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.expect("a stat")
.value,
42,
"a complete walk is authoritative even when it shrinks"
);
stat.truncated = true;
stat.observed_at = "2026-08-15T00:00:00Z".to_string();
record_network_stat(&pool, &stat).await?;
let kept = latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.expect("a stat");
assert_eq!(kept.value, 42);
assert!(
!kept.truncated,
"an equal truncated observation must not mark the kept row truncated"
);
assert_eq!(
kept.observed_at, "2026-08-14T00:00:00Z",
"the rejected observation must not have rewritten the row at all"
);
Ok(())
}
#[tokio::test]
async fn latest_network_stat_picks_the_max_across_sources() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
for (source, value, truncated) in [
("https://relay1.us-west.bsky.network", 2i64, false),
("https://relay1.us-east.bsky.network", 40i64, true),
] {
record_network_stat(
&pool,
&NetworkStat {
key: ADOPTION_STAT_KEY.to_string(),
source: source.to_string(),
value,
truncated,
observed_at: "2026-08-13T00:00:00Z".to_string(),
},
)
.await?;
}
let latest = latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.expect("a stat");
assert_eq!(latest.value, 40);
assert_eq!(latest.source, "https://relay1.us-east.bsky.network");
assert!(latest.truncated);
Ok(())
}
#[tokio::test]
async fn latest_network_stat_is_none_on_an_empty_table() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
assert!(latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.is_none());
Ok(())
}
#[tokio::test]
async fn latest_network_stat_ignores_other_keys() -> Result<()> {
let pool = init_url("sqlite::memory:").await?;
for (key, source, value) in [
(ADOPTION_STAT_KEY, "https://relay1.example", 40),
("some.other.metric", "https://relay1.example", 9_999),
] {
record_network_stat(
&pool,
&NetworkStat {
key: key.to_string(),
source: source.to_string(),
value,
truncated: false,
observed_at: now_rfc3339(),
},
)
.await?;
}
let latest = latest_network_stat(&pool, ADOPTION_STAT_KEY)
.await?
.expect("the adoption stat was recorded");
assert_eq!(
latest.value, 40,
"another key's value was returned as the adoption count"
);
Ok(())
}
async fn feed_polled(
pool: &SqlitePool,
url: &str,
last_polled: Option<&str>,
next_poll: Option<&str>,
) {
upsert_feed(
pool,
&NewFeed {
url: url.to_string(),
last_polled: last_polled.map(str::to_string),
next_poll: next_poll.map(str::to_string),
..Default::default()
},
)
.await
.unwrap();
}
#[tokio::test]
async fn poll_health_counts_tracked_recent_and_overdue() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let now = "2026-01-01T12:00:00Z";
let hour_ago = "2026-01-01T11:00:00Z";
feed_polled(
&pool,
"https://a.example/f",
Some("2026-01-01T11:50:00Z"),
Some("2026-01-01T12:50:00Z"),
)
.await;
feed_polled(
&pool,
"https://b.example/f",
Some("2026-01-01T09:00:00Z"),
Some("2026-01-01T10:00:00Z"),
)
.await;
feed_polled(&pool, "https://c.example/f", None, None).await;
let h = poll_health(&pool, now, hour_ago).await?;
assert_eq!(h.feeds_tracked, 3);
assert_eq!(
h.polled_last_hour, 1,
"only the 11:50 poll is within the hour"
);
assert_eq!(h.overdue, 2, "the stale feed and the never-polled one");
assert_eq!(
h.last_poll_secs_ago,
Some(600),
"most recent poll was 10 minutes ago"
);
assert_eq!(
h.oldest_poll_secs_ago, None,
"a never-polled feed must outrank any finite age"
);
assert_eq!(h.never_polled, 1);
sqlx::query("UPDATE feeds SET last_polled = ?1 WHERE last_polled IS NULL")
.bind("2026-01-01T09:00:00Z")
.execute(&pool)
.await?;
let h = poll_health(&pool, now, hour_ago).await?;
assert_eq!(h.never_polled, 0);
assert_eq!(h.oldest_poll_secs_ago, Some(10_800));
Ok(())
}
#[tokio::test]
async fn poll_health_ignores_unpollable_at_uri_rows() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let now = "2026-01-01T12:00:00Z";
let hour_ago = "2026-01-01T11:00:00Z";
feed_polled(
&pool,
"https://a.example/f",
Some("2026-01-01T11:50:00Z"),
Some("2026-01-01T12:50:00Z"),
)
.await;
feed_polled(
&pool,
"at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/3lab",
None,
None,
)
.await;
let h = poll_health(&pool, now, hour_ago).await?;
assert_eq!(
h.feeds_tracked, 1,
"an unpollable row was counted as tracked"
);
assert_eq!(h.overdue, 0, "an unpollable row was counted as overdue");
assert_eq!(
h.never_polled, 0,
"an unpollable row was counted as never polled"
);
assert_eq!(
h.oldest_poll_secs_ago,
Some(600),
"an unpollable row forced the oldest poll to `never`"
);
assert_eq!(h.polled_last_hour, 1);
Ok(())
}
#[tokio::test]
async fn failing_feeds_ignores_unpollable_at_uri_rows() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
for url in [
"https://broken.example/feed.xml",
"at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/3lab",
] {
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
bump_feed_errors(&pool, url, crate::feed::FailureKind::Fetch, "down").await?;
}
let failing = failing_feeds(&pool, 10).await?;
let urls: Vec<&str> = failing.iter().map(|f| f.url.as_str()).collect();
assert_eq!(
urls,
vec!["https://broken.example/feed.xml"],
"an unpollable row was listed as a failing feed"
);
Ok(())
}
#[tokio::test]
async fn the_at_uri_error_clearing_spares_a_row_that_has_been_polled() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let polled = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/polled";
let never = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/never";
for url in [polled, never] {
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
bump_feed_errors(&pool, url, crate::feed::FailureKind::Fetch, "down").await?;
}
sqlx::query("UPDATE feeds SET last_polled = '2026-01-01T00:00:00Z' WHERE url = ?1")
.bind(polled)
.execute(&pool)
.await?;
for boot in 1..=2 {
apply_migrations(&pool).await?;
let mut errors = std::collections::HashMap::new();
for url in [polled, never] {
let n: i64 =
sqlx::query_scalar("SELECT consecutive_errors FROM feeds WHERE url = ?1")
.bind(url)
.fetch_one(&pool)
.await?;
errors.insert(url, n);
}
assert_eq!(
errors[polled], 1,
"boot {boot} wiped a polled row's failure"
);
assert_eq!(
errors[never], 0,
"boot {boot} left a never-polled row failing"
);
}
Ok(())
}
#[test]
fn the_sql_kind_list_matches_the_rust_one() {
let expected = crate::feed::FeedKind::POLLABLE
.iter()
.map(|k| format!("'{}'", k.as_str()))
.collect::<Vec<_>>()
.join(", ");
assert_eq!(POLLABLE_KINDS_SQL, expected);
}
#[tokio::test]
async fn a_feed_row_records_its_kind_at_insert() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
for (url, want) in [
("https://real.example/feed.xml", crate::feed::FeedKind::Rss),
("http://real.example/feed.xml", crate::feed::FeedKind::Rss),
(
"at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab",
crate::feed::FeedKind::Publication,
),
] {
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
let got: String = sqlx::query_scalar("SELECT kind FROM feeds WHERE url = ?1")
.bind(url)
.fetch_one(&pool)
.await?;
assert_eq!(got, want.as_str(), "wrong kind recorded for {url}");
}
Ok(())
}
#[tokio::test]
async fn the_migration_backfills_kind_from_the_url() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
sqlx::query("DROP TABLE feeds").execute(&pool).await?;
sqlx::query(
"CREATE TABLE feeds (
id INTEGER PRIMARY KEY AUTOINCREMENT,
url TEXT NOT NULL UNIQUE,
title TEXT, site_url TEXT, etag TEXT, last_modified TEXT,
last_polled TEXT, next_poll TEXT,
consecutive_errors INTEGER NOT NULL DEFAULT 0,
last_error_kind TEXT, last_error TEXT
)",
)
.execute(&pool)
.await?;
for url in [
"https://real.example/feed.xml",
"at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab",
"At://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lac",
] {
sqlx::query("INSERT INTO feeds (url) VALUES (?1)")
.bind(url)
.execute(&pool)
.await?;
}
apply_migrations(&pool).await?;
let kinds: Vec<(String, String)> =
sqlx::query_as("SELECT url, kind FROM feeds ORDER BY url")
.fetch_all(&pool)
.await?;
let by_url: std::collections::HashMap<_, _> = kinds.into_iter().collect();
assert_eq!(by_url["https://real.example/feed.xml"], "rss");
assert_eq!(
by_url["at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab"],
"publication"
);
assert_eq!(
by_url["At://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lac"],
"unsupported",
"the back-fill must recognise a non-canonical spelling, and not poll it"
);
Ok(())
}
#[tokio::test]
async fn the_back_fill_corrects_a_kind_that_disagrees_with_the_url() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let at = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
for (url, wrong) in [
("https://real.example/feed.xml", "publication"),
(at, "rss"),
] {
sqlx::query("INSERT INTO feeds (url, kind) VALUES (?1, ?2)")
.bind(url)
.bind(wrong)
.execute(&pool)
.await?;
}
apply_migrations(&pool).await?;
let by_url: std::collections::HashMap<String, String> =
sqlx::query_as("SELECT url, kind FROM feeds")
.fetch_all(&pool)
.await?
.into_iter()
.collect();
assert_eq!(
by_url["https://real.example/feed.xml"], "rss",
"an http feed marked as a publication stayed one, and nothing polls it"
);
assert_eq!(by_url[at], "publication", "the at:// direction regressed");
Ok(())
}
#[tokio::test]
async fn a_row_taken_out_of_the_poller_loses_the_poll_state_it_cannot_use() -> anyhow::Result<()>
{
let pool = init_url("sqlite::memory:").await?;
let at = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/3lab";
sqlx::query(
"INSERT INTO feeds (url, kind, consecutive_errors, last_error_kind, last_error, \
next_poll, last_polled) \
VALUES (?1, 'rss', 7, 'fetch', 'connection refused', ?2, ?3)",
)
.bind(at)
.bind("2026-09-10T00:00:00Z")
.bind("2026-09-01T00:00:00Z")
.execute(&pool)
.await?;
apply_migrations(&pool).await?;
let (kind, errors, error_kind, error, next_poll): (
String,
i64,
Option<String>,
Option<String>,
Option<String>,
) = sqlx::query_as(
"SELECT kind, consecutive_errors, last_error_kind, last_error, next_poll \
FROM feeds WHERE url = ?1",
)
.bind(at)
.fetch_one(&pool)
.await?;
assert_eq!(kind, "unsupported", "the row was not reclassified at all");
assert_eq!(
(errors, error_kind, error, next_poll),
(0, None, None, None),
"a row the scheduler will never select again kept its backoff and failure history"
);
Ok(())
}
#[tokio::test]
async fn an_unreadable_feeds_row_does_not_stop_the_boot() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
sqlx::query("INSERT INTO feeds (url, kind) VALUES (X'ff41', 'rss')")
.execute(&pool)
.await?;
sqlx::query("INSERT INTO feeds (url, kind) VALUES (?1, 'publication')")
.bind("https://real.example/feed.xml")
.execute(&pool)
.await?;
apply_migrations(&pool).await?;
let corrected: String = sqlx::query_scalar("SELECT kind FROM feeds WHERE url = ?1")
.bind("https://real.example/feed.xml")
.fetch_one(&pool)
.await?;
assert_eq!(
corrected, "rss",
"one unreadable row aborted the pass before the readable ones were corrected"
);
let untouched: String =
sqlx::query_scalar("SELECT kind FROM feeds WHERE typeof(url) = 'blob'")
.fetch_one(&pool)
.await?;
assert_eq!(
untouched, "rss",
"a row we declined to classify was classified anyway"
);
Ok(())
}
#[tokio::test]
async fn a_re_upsert_re_derives_the_kind() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let url = "https://real.example/feed.xml";
let feed = NewFeed {
url: url.to_string(),
..Default::default()
};
upsert_feed(&pool, &feed).await?;
sqlx::query("UPDATE feeds SET kind = 'publication' WHERE url = ?1")
.bind(url)
.execute(&pool)
.await?;
upsert_feed(&pool, &feed).await?;
let kind: String = sqlx::query_scalar("SELECT kind FROM feeds WHERE url = ?1")
.bind(url)
.fetch_one(&pool)
.await?;
assert_eq!(
kind, "rss",
"a second subscription to the same URL kept the stale classification"
);
Ok(())
}
#[tokio::test]
async fn the_poller_and_the_pages_key_on_kind() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
upsert_feed(
&pool,
&NewFeed {
url: "https://looks-ordinary.example/feed.xml".to_string(),
..Default::default()
},
)
.await?;
sqlx::query("UPDATE feeds SET kind = 'unsupported' WHERE url LIKE 'https://looks%'")
.execute(&pool)
.await?;
let due = due_feeds(&pool, "2026-01-01T12:00:00Z", 10).await?;
assert!(due.is_empty(), "due_feeds read the URL, not the kind");
assert_eq!(
unpollable_feeds(&pool).await?,
1,
"unpollable_feeds read the URL"
);
let h = poll_health(&pool, "2026-01-01T12:00:00Z", "2026-01-01T11:00:00Z").await?;
assert_eq!(h.feeds_tracked, 0, "poll_health read the URL, not the kind");
Ok(())
}
#[tokio::test]
async fn a_mixed_case_at_uri_is_unpollable_on_both_sides() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let odd = "At://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
assert!(
!crate::feed::is_storable_feed_url(odd, true),
"a non-canonical spelling must not be storable"
);
upsert_feed(
&pool,
&NewFeed {
url: odd.to_string(),
..Default::default()
},
)
.await?;
let due = due_feeds(&pool, "2026-01-01T12:00:00Z", 10).await?;
assert!(
due.is_empty(),
"a row nothing can fetch was handed to the poller: {:?}",
due.iter().map(|f| &f.url).collect::<Vec<_>>()
);
bump_feed_errors(&pool, odd, crate::feed::FailureKind::Fetch, "refused").await?;
apply_migrations(&pool).await?;
let n: i64 = sqlx::query_scalar("SELECT consecutive_errors FROM feeds WHERE url = ?1")
.bind(odd)
.fetch_one(&pool)
.await?;
assert_eq!(n, 0, "the clearing skipped a mixed-case at-URI row");
Ok(())
}
#[tokio::test]
async fn the_ceiling_counts_unpollable_rows_and_they_are_countable() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
for url in [
"https://real.example/feed.xml",
"at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/3lab",
"At://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lac",
] {
upsert_feed(
&pool,
&NewFeed {
url: url.to_string(),
..Default::default()
},
)
.await?;
}
assert_eq!(
count_feeds(&pool).await?,
3,
"the ceiling must bound storage, so every row counts"
);
assert_eq!(
unpollable_feeds(&pool).await?,
2,
"both at-URI spellings are unpollable and must be countable"
);
Ok(())
}
#[tokio::test]
async fn poll_health_on_an_empty_instance_reports_no_polls() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let h = poll_health(&pool, "2026-01-01T12:00:00Z", "2026-01-01T11:00:00Z").await?;
assert_eq!(h.feeds_tracked, 0);
assert_eq!(h.last_poll_secs_ago, None);
assert_eq!(h.oldest_poll_secs_ago, None);
Ok(())
}
#[tokio::test]
async fn a_future_poll_timestamp_does_not_go_negative() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
feed_polled(
&pool,
"https://a.example/f",
Some("2026-01-01T13:00:00Z"),
None,
)
.await;
let h = poll_health(&pool, "2026-01-01T12:00:00Z", "2026-01-01T11:00:00Z").await?;
assert_eq!(h.last_poll_secs_ago, Some(0));
Ok(())
}
async fn aged_entry(pool: &SqlitePool, url: &str, days_old: i64) -> i64 {
let when = (chrono::Utc::now() - chrono::Duration::days(days_old))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
sqlx::query("INSERT INTO feeds (url) VALUES (?1) ON CONFLICT(url) DO NOTHING")
.bind("https://f.example/feed")
.execute(pool)
.await
.unwrap();
let feed_id: i64 = sqlx::query_scalar("SELECT id FROM feeds WHERE url = ?1")
.bind("https://f.example/feed")
.fetch_one(pool)
.await
.unwrap();
sqlx::query("INSERT INTO entries (feed_id, guid, url, title, published, fetched_at) VALUES (?1,?2,?3,'t',?4,?4)")
.bind(feed_id).bind(url).bind(url).bind(&when)
.execute(pool).await.unwrap();
sqlx::query_scalar("SELECT id FROM entries WHERE guid = ?1")
.bind(url)
.fetch_one(pool)
.await
.unwrap()
}
async fn mark(pool: &SqlitePool, entry_id: i64, read: i64, starred: i64) {
sqlx::query("INSERT INTO entry_state (did, entry_id, read, starred, updated_at) VALUES ('did:plc:x',?1,?2,?3,'2026-01-01T00:00:00Z')")
.bind(entry_id).bind(read).bind(starred)
.execute(pool).await.unwrap();
}
#[tokio::test]
async fn retention_keeps_starred_and_unread_entries() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let old_read = aged_entry(&pool, "old-read", 30).await;
let old_starred = aged_entry(&pool, "old-starred", 30).await;
let old_unread = aged_entry(&pool, "old-unread", 30).await;
let recent_read = aged_entry(&pool, "recent-read", 1).await;
mark(&pool, old_read, 1, 0).await;
mark(&pool, old_starred, 1, 1).await; mark(&pool, old_unread, 0, 0).await;
mark(&pool, recent_read, 1, 0).await;
let deleted = prune_old_entries(&pool, 14, 3650, 0).await?;
assert_eq!(deleted, 1, "only the old, read, unstarred entry should go");
let left: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
.fetch_all(&pool)
.await?;
assert_eq!(left, vec!["old-starred", "old-unread", "recent-read"]);
Ok(())
}
#[tokio::test]
async fn retention_evicts_entries_with_no_reader_state() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
aged_entry(&pool, "untouched-old", 30).await;
aged_entry(&pool, "untouched-new", 1).await;
assert_eq!(prune_old_entries(&pool, 14, 3650, 0).await?, 1);
Ok(())
}
#[tokio::test]
async fn a_recently_polled_feed_is_not_nudged_again() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let recent = "2026-01-01T11:59:00Z";
let stale_before = "2026-01-01T11:00:00Z";
sqlx::query("INSERT INTO feeds (url, last_polled, next_poll) VALUES (?1, ?2, ?3)")
.bind("https://fresh.example/f")
.bind(recent)
.bind("2026-01-01T12:59:00Z")
.execute(&pool)
.await?;
sqlx::query("INSERT INTO feeds (url, last_polled, next_poll) VALUES (?1, ?2, ?3)")
.bind("https://stale.example/f")
.bind("2026-01-01T06:00:00Z")
.bind("2026-01-01T07:00:00Z")
.execute(&pool)
.await?;
mark_feed_due(&pool, "https://fresh.example/f", stale_before).await?;
mark_feed_due(&pool, "https://stale.example/f", stale_before).await?;
let fresh: Option<String> =
sqlx::query_scalar("SELECT next_poll FROM feeds WHERE url = 'https://fresh.example/f'")
.fetch_one(&pool)
.await?;
let stale: Option<String> =
sqlx::query_scalar("SELECT next_poll FROM feeds WHERE url = 'https://stale.example/f'")
.fetch_one(&pool)
.await?;
assert!(
fresh.is_some(),
"a feed polled a minute ago was made due again — a reload loop is an \
amplification vector"
);
assert!(stale.is_none(), "a long-unpolled feed should be nudged");
Ok(())
}
#[tokio::test]
async fn a_never_polled_feed_is_nudged() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
sqlx::query("INSERT INTO feeds (url, last_polled, next_poll) VALUES (?1, NULL, ?2)")
.bind("https://new.example/f")
.bind("2026-01-01T12:59:00Z")
.execute(&pool)
.await?;
mark_feed_due(&pool, "https://new.example/f", "2026-01-01T11:00:00Z").await?;
let next: Option<String> =
sqlx::query_scalar("SELECT next_poll FROM feeds WHERE url = 'https://new.example/f'")
.fetch_one(&pool)
.await?;
assert!(next.is_none());
Ok(())
}
#[tokio::test]
async fn the_hard_ceiling_evicts_even_starred_and_unread() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let ancient_starred = aged_entry(&pool, "ancient-starred", 400).await;
let ancient_unread = aged_entry(&pool, "ancient-unread", 400).await;
let recent_starred = aged_entry(&pool, "recent-starred", 30).await;
mark(&pool, ancient_starred, 1, 1).await;
mark(&pool, ancient_unread, 0, 0).await;
mark(&pool, recent_starred, 1, 1).await;
prune_old_entries(&pool, 14, 180, 0).await?;
let left: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
.fetch_all(&pool)
.await?;
assert_eq!(
left,
vec!["recent-starred"],
"past the ceiling nothing is pinned — otherwise one reader can stall the poller \
for every reader"
);
Ok(())
}
#[tokio::test]
async fn the_per_feed_trim_spares_starred_entries() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let old_starred = aged_entry(&pool, "old-starred", 5).await;
mark(&pool, old_starred, 1, 1).await;
for i in 0..5 {
aged_entry(&pool, &format!("filler-{i}"), 1).await;
}
let feed_id: i64 = sqlx::query_scalar("SELECT id FROM feeds LIMIT 1")
.fetch_one(&pool)
.await?;
insert_entries(&pool, feed_id, &[], 2).await?;
let left: Vec<String> =
sqlx::query_scalar("SELECT guid FROM entries WHERE guid = 'old-starred'")
.fetch_all(&pool)
.await?;
assert_eq!(
left,
vec!["old-starred"],
"the per-feed trim evicted a starred entry"
);
Ok(())
}
#[tokio::test]
async fn the_trim_spares_the_newest_starred_entries_when_over_cap() -> anyhow::Result<()> {
let pool = init_url("sqlite::memory:").await?;
let mut ids = Vec::new();
for days_old in 1..=5 {
let id = aged_entry(&pool, &format!("starred-{days_old}"), days_old).await;
mark(&pool, id, 1, 1).await;
ids.push((days_old, id));
}
let feed_id: i64 = sqlx::query_scalar("SELECT feed_id FROM entries WHERE id = ?1")
.bind(ids[0].1)
.fetch_one(&pool)
.await?;
insert_entries(&pool, feed_id, &[], 2).await?;
let mut survivors: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries")
.fetch_all(&pool)
.await?;
survivors.sort();
assert_eq!(
survivors,
vec!["starred-1".to_string(), "starred-2".to_string()],
"the trim spared the wrong starred entries"
);
Ok(())
}
}