use std::collections::BTreeMap;
use tracing::{info, warn};
use crate::lexicon::ReadState;
use crate::store::{self, ReadCursor};
use crate::AppState;
pub async fn flush_did(state: &AppState, did: &str) -> anyhow::Result<()> {
let cursors = store::dirty_cursors(&state.db, did).await?;
if cursors.is_empty() {
return Ok(());
}
let mut batch: BTreeMap<String, (ReadState, ReadCursor)> = BTreeMap::new();
for cursor in cursors {
let cursor = compact_if_large(state, did, cursor).await;
let rkey = read_state_rkey(&cursor.feed_url);
let record = read_state_record(&cursor);
batch.insert(rkey, (record, cursor));
}
let ops: Vec<(String, ReadState, bool)> = batch
.iter()
.map(|(rkey, (record, cursor))| (rkey.clone(), record.clone(), cursor.pds_created))
.collect();
state.repo().flush_read_states(did, &ops).await?;
let flushed = ops.len();
for (_rkey, (_record, cursor)) in batch {
if !cursor.pds_created {
if let Err(err) = store::mark_cursor_pds_created(&state.db, did, &cursor.feed_url).await
{
warn!(%did, feed = %cursor.feed_url, %err, "failed to mark cursor pds_created");
}
}
if let Err(err) =
store::clear_cursor_dirty(&state.db, did, &cursor.feed_url, &cursor.updated_at).await
{
warn!(%did, feed = %cursor.feed_url, %err, "failed to clear cursor dirty flag");
}
}
info!(%did, feeds = flushed, "read-state flusher: flushed dirty cursors");
Ok(())
}
const COMPACT_READ_IDS_THRESHOLD: usize = ReadState::MAX_IDS / 2;
async fn compact_if_large(state: &AppState, did: &str, cursor: ReadCursor) -> ReadCursor {
if parse_id_array(&cursor.read_ids).len() < COMPACT_READ_IDS_THRESHOLD {
return cursor;
}
match store::compact_cursor(&state.db, did, &cursor.feed_url).await {
Ok(Some(watermark)) => {
match store::get_cursor(&state.db, did, &cursor.feed_url).await {
Ok(Some(fresh)) => {
info!(
%did,
feed = %cursor.feed_url,
%watermark,
before = parse_id_array(&cursor.read_ids).len(),
after = parse_id_array(&fresh.read_ids).len(),
"read-state compacted into readThrough"
);
fresh
}
Ok(None) => cursor,
Err(err) => {
warn!(%err, %did, feed = %cursor.feed_url, "could not re-read a compacted cursor");
cursor
}
}
}
Ok(None) => cursor,
Err(err) => {
warn!(%err, %did, feed = %cursor.feed_url, "read-state compaction failed; flushing uncompacted");
cursor
}
}
}
fn read_state_record(cursor: &ReadCursor) -> ReadState {
let read_ids = parse_id_array(&cursor.read_ids);
let unread_ids = parse_id_array(&cursor.unread_ids);
let mut record = ReadState::new(
&cursor.feed_url,
cursor.read_through.clone(),
&cursor.updated_at,
);
record.read_ids = cap(read_ids, ReadState::MAX_IDS);
record.unread_ids = cap(unread_ids, ReadState::MAX_IDS);
record
}
fn parse_id_array(raw: &str) -> Vec<String> {
if raw.trim().is_empty() {
return Vec::new();
}
match serde_json::from_str::<Vec<serde_json::Value>>(raw) {
Ok(vals) => vals
.into_iter()
.map(|v| match v {
serde_json::Value::String(s) => s,
other => other.to_string(),
})
.collect(),
Err(err) => {
warn!(%err, raw, "read-state flusher: unparseable id array; treating as empty");
Vec::new()
}
}
}
fn cap(mut ids: Vec<String>, max: usize) -> Vec<String> {
if ids.len() > max {
let drop = ids.len() - max;
warn!(
dropped = drop,
kept = max,
"read-state id set exceeded the lexicon cap even after compaction; \
the oldest marks will not sync"
);
ids.drain(0..drop);
}
ids
}
pub fn read_state_rkey(feed_url: &str) -> String {
format!("rs-{:016x}", fnv1a_64(feed_url.as_bytes()))
}
pub fn fnv1a_64(bytes: &[u8]) -> u64 {
const OFFSET: u64 = 0xcbf2_9ce4_8422_2325;
const PRIME: u64 = 0x0000_0100_0000_01b3;
let mut hash = OFFSET;
for &b in bytes {
hash ^= b as u64;
hash = hash.wrapping_mul(PRIME);
}
hash
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn rkey_is_stable_and_valid() {
let a = read_state_rkey("https://example.com/feed.xml");
let b = read_state_rkey("https://example.com/feed.xml");
assert_eq!(a, b, "rkey must be deterministic");
assert_ne!(a, read_state_rkey("https://other.example/feed.xml"));
assert!(crate::atproto::is_valid_rkey(&a), "{a:?}");
}
#[test]
fn parse_id_array_tolerates_shapes() {
assert_eq!(parse_id_array(""), Vec::<String>::new());
assert_eq!(parse_id_array("[]"), Vec::<String>::new());
assert_eq!(parse_id_array(r#"["a","b"]"#), vec!["a", "b"]);
assert_eq!(parse_id_array("[1,2,3]"), vec!["1", "2", "3"]);
assert_eq!(parse_id_array("not json"), Vec::<String>::new());
}
#[test]
fn cap_keeps_tail_within_bound() {
let ids: Vec<String> = (0..10).map(|i| i.to_string()).collect();
let capped = cap(ids, 3);
assert_eq!(capped, vec!["7", "8", "9"]);
}
#[test]
fn read_state_record_applies_the_id_cap() {
let ids: Vec<String> = (0..ReadState::MAX_IDS + 5).map(|i| i.to_string()).collect();
let json = serde_json::to_string(&ids).unwrap();
let cursor = crate::store::ReadCursor {
did: "did:plc:x".into(),
feed_url: "https://example.com/feed.xml".into(),
read_through: None,
read_ids: json.clone(),
unread_ids: json,
dirty: true,
pds_created: false,
updated_at: "2026-07-12T00:00:00Z".into(),
};
let rec = read_state_record(&cursor);
assert_eq!(
rec.read_ids.len(),
ReadState::MAX_IDS,
"read_ids not capped"
);
assert_eq!(
rec.unread_ids.len(),
ReadState::MAX_IDS,
"unread_ids not capped"
);
}
#[test]
fn record_maps_cursor_fields() {
let cursor = ReadCursor {
did: "did:plc:abc".into(),
feed_url: "https://example.com/feed.xml".into(),
read_through: Some("2026-07-12T00:00:00Z".into()),
read_ids: r#"["10","11"]"#.into(),
unread_ids: "[]".into(),
dirty: true,
pds_created: false,
updated_at: "2026-07-12T01:00:00Z".into(),
};
let rec = read_state_record(&cursor);
assert_eq!(rec.feed_url, "https://example.com/feed.xml");
assert_eq!(rec.read_through.as_deref(), Some("2026-07-12T00:00:00Z"));
assert_eq!(rec.read_ids, vec!["10", "11"]);
assert!(rec.unread_ids.is_empty());
assert_eq!(rec.updated_at, "2026-07-12T01:00:00Z");
}
#[test]
fn read_through_omitted_when_local_unset() {
let cursor = ReadCursor {
did: "did:plc:abc".into(),
feed_url: "https://example.com/feed.xml".into(),
read_through: None,
read_ids: r#"["42"]"#.into(),
unread_ids: "[]".into(),
dirty: true,
pds_created: false,
updated_at: "2026-07-12T01:00:00Z".into(),
};
let rec = read_state_record(&cursor);
assert_eq!(
rec.read_through, None,
"no local water-mark => readThrough absent (backlog not implicitly read)"
);
assert_eq!(rec.read_ids, vec!["42"]);
let json = serde_json::to_value(&rec).expect("serialize");
assert!(json.get("readThrough").is_none());
}
#[test]
fn read_through_present_when_local_high_water_mark_exists() {
let cursor = ReadCursor {
did: "did:plc:abc".into(),
feed_url: "https://example.com/feed.xml".into(),
read_through: Some("2026-07-11T00:00:00Z".into()),
read_ids: "[]".into(),
unread_ids: "[]".into(),
dirty: true,
pds_created: false,
updated_at: "2026-07-12T01:00:00Z".into(),
};
let rec = read_state_record(&cursor);
assert_eq!(rec.read_through.as_deref(), Some("2026-07-11T00:00:00Z"));
}
#[test]
fn flush_with_only_read_ids_sets_no_read_through() {
let cursor = ReadCursor {
did: "did:plc:abc".into(),
feed_url: "https://example.com/feed.xml".into(),
read_through: None,
read_ids: r#"["100","101","102"]"#.into(),
unread_ids: "[]".into(),
dirty: true,
pds_created: false,
updated_at: "2026-07-12T02:00:00Z".into(),
};
let rec = read_state_record(&cursor);
assert_eq!(rec.read_through, None);
assert_eq!(rec.read_ids, vec!["100", "101", "102"]);
let json = serde_json::to_value(&rec).expect("serialize");
assert!(json.get("readThrough").is_none());
assert_eq!(json["readIds"], serde_json::json!(["100", "101", "102"]));
}
}