Skip to main content

feather_reader/
readstate.rs

1//! Read-state flushing — turning dirty local cursors into `readState` records
2//! in the user's own PDS.
3//!
4//! **In the LIB, not the binary's `scheduler`, because two callers need it.**
5//! The background flusher is one. Sign-out is the other: it must flush before
6//! revoking the session, or the reads are stranded with nothing able to send
7//! them (#117). `scheduler` keeps the *scheduling* — the interval loop and the
8//! DID selection; this module owns the domain logic.
9
10use std::collections::{BTreeMap, HashSet};
11
12use tracing::{info, warn};
13
14use crate::atproto::AtProtoError;
15use crate::lexicon::ReadState;
16use crate::store::{self, ReadCursor};
17use crate::AppState;
18
19/// Flush a single DID's dirty cursors in one batched `applyWrites`, then clear
20/// the `dirty` flag for the cursors that were included.
21pub async fn flush_did(state: &AppState, did: &str) -> anyhow::Result<()> {
22    let cursors = store::dirty_cursors(&state.db, did).await?;
23    if cursors.is_empty() {
24        return Ok(());
25    }
26
27    // Build (rkey, ReadState) pairs, deduping on rkey so two rows that hash to
28    // the same feed-key don't produce two ops in one batch (applyWrites rejects
29    // duplicate writes to the same key). Deterministic order for stable batches.
30    let mut batch: Batch = BTreeMap::new();
31    for cursor in cursors {
32        // **Compact before capping.**
33        //
34        // `read_ids` grows one id per article read and is bounded only by
35        // `max_entries_per_feed` (2000), while `cap` below truncates the record
36        // at `ReadState::MAX_IDS` (1000) keeping the TAIL. Past 1000 read
37        // articles in one feed the oldest read-state silently stopped syncing,
38        // and those articles came back UNREAD in every other atproto reader —
39        // the one thing the shared lexicon exists to prevent.
40        //
41        // `store::compact_cursor` folds the covered ids into the `read_through`
42        // high-water-mark, which is the field that exists for exactly this and
43        // was never being computed. Done here rather than on every mark-read
44        // because this is the moment the size actually matters, and it is per
45        // dirty cursor per flush rather than per click.
46        let cursor = compact_if_large(state, did, cursor).await;
47        let rkey = read_state_rkey(&cursor.feed_url);
48        let record = read_state_record(&cursor);
49        batch.insert(rkey, (record, cursor));
50    }
51
52    // ONE applyWrites batch for all of this DID's dirty feeds — sent as several
53    // calls past the PDS's per-call limits (#240), in rkey order.
54    let ops = batch_ops(&batch);
55    if let Err(err) = state.repo().flush_read_states(did, &ops).await {
56        // **What landed is settled first, whatever went wrong after it.** A
57        // split batch commits call by call, so a failure in call 2 can follow a
58        // call 1 that created its records. Left dirty with `pds_created` false,
59        // those cursors went out again next round as `#create` on keys that now
60        // exist, the PDS refused, the run stopped at that call — and because the
61        // order is fixed, every cursor sorted after them starved, every round.
62        settle_landed(state, did, &mut batch, &err).await;
63
64        // **A flag that disagrees with the PDS wedged this DID forever (#241).**
65        //
66        // `pds_created` picks create-vs-update and is learned only from a
67        // success. A success whose answer was lost, or a fresh or restored
68        // database against a repo that already holds these stable rkeys, leaves
69        // it false over a record that exists; a record deleted elsewhere leaves
70        // it true over one that does not. Either way one op fails, applyWrites
71        // is atomic, the whole call fails — and the next flush sends the same
72        // batch. Nothing ever corrected the flag.
73        //
74        // So on a failure that COULD be that, ask the PDS once what exists, and
75        // retry what did not land, once, if the answer changes anything. Never
76        // in a loop: a retry that fails again is returned like any other
77        // failure, and the corrected flags it leaves behind make the next round
78        // an ordinary flush.
79        //
80        // Not ALSO proactively, on each DID's first flush after startup: that
81        // costs every healthy DID a listing per restart to save a wedged one a
82        // single failed applyWrites, since this path converges inside the same
83        // flush — and it would not cover a success lost mid-process anyway.
84        if !may_be_existence_mismatch(&err) {
85            return Err(err);
86        }
87        let corrected = reconcile_pds_created(state, did, &mut batch).await;
88        match corrected {
89            Ok(0) => return Err(err),
90            Ok(fixed) => {
91                info!(%did, fixed, "read-state flusher: pds_created disagreed with the PDS; reconciled, retrying once");
92            }
93            Err(list_err) => {
94                warn!(%did, err = %list_err, "read-state flusher: could not list readState records to reconcile");
95                return Err(err);
96            }
97        }
98        let ops = batch_ops(&batch);
99        if let Err(retry_err) = state.repo().flush_read_states(did, &ops).await {
100            // The retry can split and part-land too.
101            settle_landed(state, did, &mut batch, &retry_err).await;
102            // **The reason goes IN the message, not only under it.** Both
103            // callers log `%err`, which prints an anyhow context's outermost
104            // message alone, so a bare "failed again" dropped the PDS's status
105            // and error name from the log. Switching the callers to `{:#}`
106            // instead would print the reason twice on every flush failure:
107            // the Rust client's rejection, and `ApplyWritesIncomplete`, both
108            // already repeat their cause in their own message. The chain is
109            // kept, so `ApplyWritesIncomplete::of` still reads the progress.
110            let context =
111                format!("read-state flush failed again after reconciling pds_created: {retry_err}");
112            return Err(retry_err.context(context));
113        }
114    }
115
116    let flushed = batch.len();
117    for (_rkey, (_record, cursor)) in batch {
118        settle(state, did, &cursor).await;
119    }
120
121    info!(%did, feeds = flushed, "read-state flusher: flushed dirty cursors");
122    Ok(())
123}
124
125/// Record that `cursor`'s write landed: mark its PDS record as created (so
126/// future flushes emit an update), then clear `dirty` but ONLY if its
127/// `updated_at` still matches the snapshot that was flushed. A mark-read that
128/// landed DURING the in-flight PDS write bumped `updated_at` and re-dirtied the
129/// row; the conditional clear leaves that row dirty so its new reads re-flush
130/// next round instead of being silently dropped.
131async fn settle(state: &AppState, did: &str, cursor: &ReadCursor) {
132    // Flip the created flag first: the record now exists in the PDS regardless
133    // of whether the dirty-clear below is a no-op due to a concurrent bump.
134    if !cursor.pds_created {
135        if let Err(err) = store::mark_cursor_pds_created(&state.db, did, &cursor.feed_url).await {
136            warn!(%did, feed = %cursor.feed_url, %err, "failed to mark cursor pds_created");
137        }
138    }
139    if let Err(err) =
140        store::clear_cursor_dirty(&state.db, did, &cursor.feed_url, &cursor.updated_at).await
141    {
142        // The PDS write already landed; a failure to clear the local flag
143        // just means we harmlessly re-flush this cursor next round.
144        warn!(%did, feed = %cursor.feed_url, %err, "failed to clear cursor dirty flag");
145    }
146}
147
148/// Settle, and take out of `batch`, the cursors a failed flush DID write.
149///
150/// A chunked `applyWrites` stops at the first failing call, so what landed is
151/// a prefix of the ops — and the ops are `batch` in its own (rkey) order, so
152/// it is a prefix of `batch` too. [`crate::atproto::ApplyWritesIncomplete`]
153/// says how long. The failed call's own writes are in doubt and stay in the
154/// batch, dirty: a later round, or the reconcile, finds out which landed.
155async fn settle_landed(state: &AppState, did: &str, batch: &mut Batch, err: &anyhow::Error) {
156    let landed = crate::atproto::ApplyWritesIncomplete::of(err).map_or(0, |p| p.landed);
157    let keys: Vec<String> = batch.keys().take(landed).cloned().collect();
158    for key in keys {
159        if let Some((_record, cursor)) = batch.remove(&key) {
160            settle(state, did, &cursor).await;
161        }
162    }
163    if landed > 0 {
164        info!(%did, landed, "read-state flusher: a split flush failed part-way; settled what landed");
165    }
166}
167
168/// One flush's cursors, by rkey: the record to write and the row it came from.
169type Batch = BTreeMap<String, (ReadState, ReadCursor)>;
170
171/// The `(rkey, record, pds_created)` ops for a batch.
172///
173/// Each op carries whether its PDS record already exists: a not-yet-created
174/// cursor becomes an applyWrites#create (not an #update, which would error and,
175/// since applyWrites is atomic-per-repo, drop the whole DID batch on a feed's
176/// first flush). All create + update ops ride ONE batch.
177fn batch_ops(batch: &Batch) -> Vec<(String, ReadState, bool)> {
178    batch
179        .iter()
180        .map(|(rkey, (record, cursor))| (rkey.clone(), record.clone(), cursor.pds_created))
181        .collect()
182}
183
184/// Whether a failed flush may have been refused for a create/update mismatch —
185/// `#create` at an rkey that exists, or `#update` at one that does not.
186///
187/// **Matched on the structured rejection, and it cannot be narrower than
188/// this.** Both backends surface a PDS refusal as [`AtProtoError::Xrpc`] with
189/// the status and error name; anything else — a transport failure, a missing
190/// session, a body we could not read — is not a refusal at all, and is not
191/// reconciled.
192///
193/// What the reference PDS sends for a mismatch was read from its source, not
194/// guessed: `applyWrites` checks nothing per op without a `swapRecord`, so the
195/// collision surfaces in `@atproto/repo`'s MST, whose `add` throws `There is
196/// already a value at key` and whose `update` throws `Could not find a record
197/// with key` — plain `Error`s, which xrpc-server answers as **500
198/// `InternalServerError`**, message replaced by "Internal Server Error". There
199/// is no more specific signal to match. A mismatch-shaped 500 is therefore
200/// ambiguous by construction, which is why the reconcile retries only when the
201/// listing actually finds a flag to correct.
202///
203/// A 400 or 409 naming a conflict is what a PDS that checks explicitly would
204/// send (`InvalidSwap` is the reference's own name for a swap mismatch), so
205/// those are included. NOT included: 401/403 (auth — a listing would fail the
206/// same way), 429, and 502/503/504, which say the request may not have been
207/// processed at all — a reconcile there would add a repo walk per DID per
208/// round to a PDS that is already struggling. If such a write did land, the
209/// next round's create meets the record and reconciles then.
210fn may_be_existence_mismatch(err: &anyhow::Error) -> bool {
211    // **The chunked write's own cause is walked too, explicitly.** Every flush
212    // error now arrives wrapped in `ApplyWritesIncomplete`, whose `source()`
213    // continues from its cause's SOURCE (so `{:#}` does not print the PDS's
214    // message twice). On the sidecar client the cause IS the `AtProtoError`,
215    // so `err.chain()` skips it — measured: every sidecar reconcile test went
216    // red on the merge with #240 until this was added.
217    let wrapped = crate::atproto::ApplyWritesIncomplete::of(err).map(|p| p.cause().chain());
218    err.chain()
219        .chain(wrapped.into_iter().flatten())
220        .any(|cause| match cause.downcast_ref::<AtProtoError>() {
221            Some(AtProtoError::Xrpc { status, error, .. }) => match status.as_u16() {
222                500 => error == "InternalServerError",
223                400 => matches!(
224                    error.as_str(),
225                    "InvalidRequest" | "InvalidSwap" | "RecordNotFound"
226                ),
227                409 => true,
228                _ => false,
229            },
230            _ => false,
231        })
232}
233
234/// Set each batch cursor's `pds_created` to whether its record exists on the
235/// PDS, locally and in memory. Returns how many flags changed.
236///
237/// **One listing of the collection, not a `getRecord` per rkey.** The sidecar
238/// has no `get` action, so per-rkey reads would mean new surface on a backend
239/// being retired; and a listing answers for every cursor in the batch at once,
240/// which matters when the batch is hundreds of feeds — including ones a
241/// partially-landed chunked batch (#240) created. It is the existing bounded
242/// walk: a repo past its page cap REFUSES rather than answering short, and a
243/// refusal here just returns the original error, which is the pre-#241
244/// behaviour rather than a wrong flag.
245///
246/// Raw records, not [`crate::repo::Repo::list_read_states`]: that skips a record
247/// it cannot parse, and a skipped record would read as a missing one and send
248/// `#create` straight back into the collision.
249///
250/// The corrected flag is persisted even if the retry then fails, because it is
251/// what the PDS holds either way — and it makes the next round a plain flush.
252async fn reconcile_pds_created(
253    state: &AppState,
254    did: &str,
255    batch: &mut Batch,
256) -> anyhow::Result<usize> {
257    let existing: HashSet<String> = state
258        .repo()
259        .list_all_records(did, crate::lexicon::nsid::READ_STATE)
260        .await?
261        .iter()
262        .filter_map(|record| record.rkey().map(str::to_string))
263        .collect();
264    let mut fixed = 0;
265    for (rkey, (_record, cursor)) in batch.iter_mut() {
266        let exists = existing.contains(rkey);
267        if cursor.pds_created == exists {
268            continue;
269        }
270        cursor.pds_created = exists;
271        fixed += 1;
272        if let Err(err) =
273            store::set_cursor_pds_created(&state.db, did, &cursor.feed_url, exists).await
274        {
275            // The in-memory flag still drives the retry. The row keeps the
276            // stale flag — the success path only marks rows it believes it
277            // CREATED — so the next flush meets the same mismatch and
278            // reconciles again: one more listing, not a wedge.
279            warn!(%did, feed = %cursor.feed_url, %err, "failed to record reconciled pds_created");
280        }
281    }
282    Ok(fixed)
283}
284
285/// `read_ids` length at which a cursor is compacted before flushing.
286///
287/// Half of [`ReadState::MAX_IDS`], so compaction happens well before the cap
288/// truncates anything, and the common cursor — a handful of ids — never pays for
289/// the two extra queries.
290const COMPACT_READ_IDS_THRESHOLD: usize = ReadState::MAX_IDS / 2;
291
292/// Fold covered ids into the `read_through` water-mark when the exception set has
293/// grown enough to matter, and return the rewritten cursor.
294///
295/// On ANY failure this returns the cursor it was given. Flushing an uncompacted
296/// cursor is the behaviour that shipped for months — a compaction problem must
297/// not become a read-state-sync problem.
298async fn compact_if_large(state: &AppState, did: &str, cursor: ReadCursor) -> ReadCursor {
299    if parse_id_array(&cursor.read_ids).len() < COMPACT_READ_IDS_THRESHOLD {
300        return cursor;
301    }
302    match store::compact_cursor(&state.db, did, &cursor.feed_url).await {
303        Ok(Some(watermark)) => {
304            // Re-read: `compact_cursor` rewrote the row, and the flusher's
305            // conditional dirty-clear compares `updated_at` against the version
306            // it flushed. Carrying the pre-compaction snapshot forward would
307            // clear a flag for a row that has since changed.
308            match store::get_cursor(&state.db, did, &cursor.feed_url).await {
309                Ok(Some(fresh)) => {
310                    info!(
311                        %did,
312                        feed = %cursor.feed_url,
313                        %watermark,
314                        before = parse_id_array(&cursor.read_ids).len(),
315                        after = parse_id_array(&fresh.read_ids).len(),
316                        "read-state compacted into readThrough"
317                    );
318                    fresh
319                }
320                Ok(None) => cursor,
321                Err(err) => {
322                    warn!(%err, %did, feed = %cursor.feed_url, "could not re-read a compacted cursor");
323                    cursor
324                }
325            }
326        }
327        Ok(None) => cursor,
328        Err(err) => {
329            warn!(%err, %did, feed = %cursor.feed_url, "read-state compaction failed; flushing uncompacted");
330            cursor
331        }
332    }
333}
334
335/// Turn a local [`ReadCursor`] row into the PDS [`ReadState`] lexicon record.
336///
337/// The store keeps `read_ids` / `unread_ids` as JSON arrays of ids; the lexicon
338/// wants string arrays. `read_through` is optional both locally AND in the
339/// record: when the cursor has no local high-water-mark we pass `None` so the
340/// record OMITS `readThrough` entirely. This is the conservative behaviour —
341/// `readThrough` is a "everything seen/published `<=` this is read" water-mark,
342/// so synthesizing a flush-time (`≈ now`) value for a cursor that has none would
343/// assert the whole unread backlog is read. With `None` only the explicit
344/// `read_ids` mark entries read. Both id-sets are capped at [`ReadState::MAX_IDS`]
345/// to respect the lexicon bound.
346fn read_state_record(cursor: &ReadCursor) -> ReadState {
347    let read_ids = parse_id_array(&cursor.read_ids);
348    let unread_ids = parse_id_array(&cursor.unread_ids);
349
350    // Do NOT synthesize a water-mark from `updated_at`: an unset local
351    // `read_through` means "no high-water-mark", which the record represents by
352    // omitting `readThrough` (None), not by back-dating it to flush time.
353    let mut record = ReadState::new(
354        &cursor.feed_url,
355        cursor.read_through.clone(),
356        &cursor.updated_at,
357    );
358    record.read_ids = cap(read_ids, ReadState::MAX_IDS);
359    record.unread_ids = cap(unread_ids, ReadState::MAX_IDS);
360    record
361}
362
363/// Parse a stored JSON id-array into `Vec<String>`, tolerating both string and
364/// numeric ids (the store keeps entry ids). A malformed/empty value yields an
365/// empty set rather than an error — read-state must never fail to flush over a
366/// cosmetic parse issue.
367fn parse_id_array(raw: &str) -> Vec<String> {
368    if raw.trim().is_empty() {
369        return Vec::new();
370    }
371    match serde_json::from_str::<Vec<serde_json::Value>>(raw) {
372        Ok(vals) => vals
373            .into_iter()
374            .map(|v| match v {
375                serde_json::Value::String(s) => s,
376                other => other.to_string(),
377            })
378            .collect(),
379        Err(err) => {
380            warn!(%err, raw, "read-state flusher: unparseable id array; treating as empty");
381            Vec::new()
382        }
383    }
384}
385
386/// Truncate a set to `max`, keeping the most recent (tail) ids — the lexicon's
387/// hard cap, and the LAST line of defence rather than the only one.
388///
389/// This used to be the only one, under the stated assumption that "the exception
390/// sets are expected to stay well under the cap in normal use". Against a 2000
391/// entry per-feed ceiling and one id per article read, that did not hold: past
392/// 1000 read articles in a feed this silently dropped the oldest read-state, and
393/// those articles came back UNREAD in every other atproto reader.
394///
395/// [`compact_if_large`] now folds covered ids into `read_through` before a cursor
396/// gets here, so reaching this truncation means compaction could not advance the
397/// water-mark — which happens only when the feed's oldest entry is genuinely
398/// unread. Losing the tail is still wrong in that case, but it is now a rare
399/// shape rather than the ordinary consequence of reading a busy feed.
400fn cap(mut ids: Vec<String>, max: usize) -> Vec<String> {
401    if ids.len() > max {
402        let drop = ids.len() - max;
403        warn!(
404            dropped = drop,
405            kept = max,
406            "read-state id set exceeded the lexicon cap even after compaction; \
407             the oldest marks will not sync"
408        );
409        ids.drain(0..drop);
410    }
411    ids
412}
413
414/// Derive the deterministic, stable rkey for a feed's read-state record from its
415/// URL, so there is exactly **one record per feed** (a fixed key, not a fresh tid
416/// per flush).
417///
418/// atproto record keys must match `[A-Za-z0-9._~:-]{1,512}` (and not be `.`/`..`).
419/// A lowercase-hex FNV-1a-64 digest of the feed URL satisfies that, is stable
420/// across restarts and instances, and collides only on genuine hash collision
421/// (astronomically unlikely at feed scale; the flusher additionally dedups by
422/// rkey within a batch as a belt-and-braces guard).
423///
424/// The stable rkey is what makes create-then-update work: a feed's FIRST flush
425/// emits an `applyWrites#create` at this key (tracked by `read_cursor.pds_created`)
426/// and every subsequent flush an `#update` at the same key, so there is exactly
427/// one record per feed and the first flush never fails on a missing record.
428pub fn read_state_rkey(feed_url: &str) -> String {
429    format!("rs-{:016x}", fnv1a_64(feed_url.as_bytes()))
430}
431
432/// FNV-1a 64-bit — a tiny, dependency-free stable hash for the feed-key.
433///
434/// `pub` because `scheduler::jittered` seeds its per-feed jitter from the same
435/// hash. Two unrelated uses of one generic utility; exported rather than copied,
436/// since two drifting implementations of a stable key hash would be worse than
437/// the slightly odd home.
438pub fn fnv1a_64(bytes: &[u8]) -> u64 {
439    const OFFSET: u64 = 0xcbf2_9ce4_8422_2325;
440    const PRIME: u64 = 0x0000_0100_0000_01b3;
441    let mut hash = OFFSET;
442    for &b in bytes {
443        hash ^= b as u64;
444        hash = hash.wrapping_mul(PRIME);
445    }
446    hash
447}
448
449#[cfg(test)]
450pub(crate) mod tests {
451    use super::*;
452
453    #[test]
454    fn rkey_is_stable_and_valid() {
455        let a = read_state_rkey("https://example.com/feed.xml");
456        let b = read_state_rkey("https://example.com/feed.xml");
457        assert_eq!(a, b, "rkey must be deterministic");
458        assert_ne!(a, read_state_rkey("https://other.example/feed.xml"));
459        // Valid atproto rkey: charset, length and the reserved names, as the
460        // one shared rule states them.
461        assert!(crate::atproto::is_valid_rkey(&a), "{a:?}");
462    }
463    #[test]
464    fn parse_id_array_tolerates_shapes() {
465        assert_eq!(parse_id_array(""), Vec::<String>::new());
466        assert_eq!(parse_id_array("[]"), Vec::<String>::new());
467        assert_eq!(parse_id_array(r#"["a","b"]"#), vec!["a", "b"]);
468        assert_eq!(parse_id_array("[1,2,3]"), vec!["1", "2", "3"]);
469        assert_eq!(parse_id_array("not json"), Vec::<String>::new());
470    }
471    #[test]
472    fn cap_keeps_tail_within_bound() {
473        let ids: Vec<String> = (0..10).map(|i| i.to_string()).collect();
474        let capped = cap(ids, 3);
475        assert_eq!(capped, vec!["7", "8", "9"]);
476    }
477    /// **The cap is APPLIED, not just correct.** `cap` had its own unit test
478    /// and every record-building test used 1–3 ids, so `read_state_record`
479    /// could stop calling it with the suite green — and the flusher would
480    /// publish id arrays past the lexicon's bound, the PDS would reject the
481    /// atomic batch, and every feed's read-state would stop syncing.
482    #[test]
483    fn read_state_record_applies_the_id_cap() {
484        let ids: Vec<String> = (0..ReadState::MAX_IDS + 5).map(|i| i.to_string()).collect();
485        let json = serde_json::to_string(&ids).unwrap();
486        let cursor = crate::store::ReadCursor {
487            did: "did:plc:x".into(),
488            feed_url: "https://example.com/feed.xml".into(),
489            read_through: None,
490            read_ids: json.clone(),
491            unread_ids: json,
492            dirty: true,
493            pds_created: false,
494            updated_at: "2026-07-12T00:00:00Z".into(),
495        };
496        let rec = read_state_record(&cursor);
497        assert_eq!(
498            rec.read_ids.len(),
499            ReadState::MAX_IDS,
500            "read_ids not capped"
501        );
502        assert_eq!(
503            rec.unread_ids.len(),
504            ReadState::MAX_IDS,
505            "unread_ids not capped"
506        );
507    }
508
509    #[test]
510    fn record_maps_cursor_fields() {
511        let cursor = ReadCursor {
512            did: "did:plc:abc".into(),
513            feed_url: "https://example.com/feed.xml".into(),
514            read_through: Some("2026-07-12T00:00:00Z".into()),
515            read_ids: r#"["10","11"]"#.into(),
516            unread_ids: "[]".into(),
517            dirty: true,
518            pds_created: false,
519            updated_at: "2026-07-12T01:00:00Z".into(),
520        };
521        let rec = read_state_record(&cursor);
522        assert_eq!(rec.feed_url, "https://example.com/feed.xml");
523        assert_eq!(rec.read_through.as_deref(), Some("2026-07-12T00:00:00Z"));
524        assert_eq!(rec.read_ids, vec!["10", "11"]);
525        assert!(rec.unread_ids.is_empty());
526        assert_eq!(rec.updated_at, "2026-07-12T01:00:00Z");
527    }
528    #[test]
529    fn read_through_omitted_when_local_unset() {
530        // A cursor with no local high-water-mark must NOT synthesize one from
531        // `updated_at` (≈ now) — doing so would mark the whole backlog read. The
532        // record omits `readThrough` (None) so only explicit read_ids apply.
533        let cursor = ReadCursor {
534            did: "did:plc:abc".into(),
535            feed_url: "https://example.com/feed.xml".into(),
536            read_through: None,
537            read_ids: r#"["42"]"#.into(),
538            unread_ids: "[]".into(),
539            dirty: true,
540            pds_created: false,
541            updated_at: "2026-07-12T01:00:00Z".into(),
542        };
543        let rec = read_state_record(&cursor);
544        assert_eq!(
545            rec.read_through, None,
546            "no local water-mark => readThrough absent (backlog not implicitly read)"
547        );
548        // The explicit read_ids still carry through.
549        assert_eq!(rec.read_ids, vec!["42"]);
550        // Serialized form must not carry a readThrough field at all.
551        let json = serde_json::to_value(&rec).expect("serialize");
552        assert!(json.get("readThrough").is_none());
553    }
554    #[test]
555    fn read_through_present_when_local_high_water_mark_exists() {
556        // A real high-water-mark IS written through unchanged.
557        let cursor = ReadCursor {
558            did: "did:plc:abc".into(),
559            feed_url: "https://example.com/feed.xml".into(),
560            read_through: Some("2026-07-11T00:00:00Z".into()),
561            read_ids: "[]".into(),
562            unread_ids: "[]".into(),
563            dirty: true,
564            pds_created: false,
565            updated_at: "2026-07-12T01:00:00Z".into(),
566        };
567        let rec = read_state_record(&cursor);
568        assert_eq!(rec.read_through.as_deref(), Some("2026-07-11T00:00:00Z"));
569    }
570    #[test]
571    fn flush_with_only_read_ids_sets_no_read_through() {
572        // The core F1 guarantee: a flush whose cursor carries only explicit
573        // read_ids (and no water-mark) emits a record WITHOUT readThrough, so the
574        // user's PDS never asserts the backlog is read.
575        let cursor = ReadCursor {
576            did: "did:plc:abc".into(),
577            feed_url: "https://example.com/feed.xml".into(),
578            read_through: None,
579            read_ids: r#"["100","101","102"]"#.into(),
580            unread_ids: "[]".into(),
581            dirty: true,
582            pds_created: false,
583            updated_at: "2026-07-12T02:00:00Z".into(),
584        };
585        let rec = read_state_record(&cursor);
586        assert_eq!(rec.read_through, None);
587        assert_eq!(rec.read_ids, vec!["100", "101", "102"]);
588        let json = serde_json::to_value(&rec).expect("serialize");
589        assert!(json.get("readThrough").is_none());
590        assert_eq!(json["readIds"], serde_json::json!(["100", "101", "102"]));
591    }
592
593    // ── #241: a pds_created flag that disagrees with the PDS ─────────────────
594    //
595    // Driven end to end through `flush_did` against a STATEFUL fake repo,
596    // because the bug is the interaction: the flag picks create-vs-update, the
597    // PDS refuses the whole atomic batch, and nothing ever corrects the flag.
598    // The fixed-body servers in `net::tests` answer every request identically,
599    // so they cannot model "the second applyWrites succeeds because the listing
600    // changed what was sent".
601
602    use std::sync::{Arc, Mutex};
603
604    use crate::metrics::Backend;
605
606    pub(crate) const DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
607
608    /// An applyWrites failure the fake answers with.
609    #[derive(Clone, Copy, Debug)]
610    struct Fail {
611        status: u16,
612        error: &'static str,
613    }
614
615    /// What the reference PDS answers to `#create` on an existing rkey AND to
616    /// `#update` on a missing one: `@atproto/repo`'s MST `add`/`update` throw a
617    /// plain `Error`, which xrpc-server turns into a 500 whose message it strips.
618    const MISMATCH: Fail = Fail {
619        status: 500,
620        error: "InternalServerError",
621    };
622
623    /// One account's `readState` collection, with the reference PDS's
624    /// create/update semantics and atomicity, and counters for what was asked.
625    #[derive(Default)]
626    pub(crate) struct FakeRepo {
627        pub(crate) records: BTreeMap<String, serde_json::Value>,
628        pub(crate) apply_calls: usize,
629        /// listRecords requests with no cursor — one per walk, however many
630        /// pages that walk turns out to need.
631        pub(crate) list_walks: usize,
632        /// Fail every applyWrites with this, whatever its ops.
633        always_fail: Option<Fail>,
634        /// Land the first N ops of the next applyWrites, then answer 503 — a
635        /// later chunk failing after an earlier one committed (#240).
636        partial_then_503: Option<usize>,
637        /// Answer every listRecords with a 503.
638        list_fails: bool,
639        /// Drop the connection, unanswered, on this applyWrites call (1-based)
640        /// — a transport failure after the earlier calls committed.
641        pub(crate) drop_call: Option<usize>,
642    }
643
644    /// Not an answer: the fake hangs up instead of replying.
645    const HANG_UP: Fail = Fail {
646        status: 0,
647        error: "connection dropped",
648    };
649
650    impl FakeRepo {
651        fn op(w: &serde_json::Value) -> (String, String, serde_json::Value) {
652            // The sidecar spells the action out; XRPC tags the union member.
653            let action = w["action"].as_str().map(str::to_string).unwrap_or_else(|| {
654                w["$type"]
655                    .as_str()
656                    .and_then(|t| t.rsplit('#').next())
657                    .unwrap_or_default()
658                    .to_string()
659            });
660            assert_eq!(w["collection"], crate::lexicon::nsid::READ_STATE);
661            let rkey = w["rkey"].as_str().expect("readState ops carry an rkey");
662            (action, rkey.to_string(), w["value"].clone())
663        }
664
665        fn apply(&mut self, writes: &[serde_json::Value]) -> Result<(), Fail> {
666            self.apply_calls += 1;
667            if self.drop_call == Some(self.apply_calls) {
668                return Err(HANG_UP);
669            }
670            // The reference PDS's per-call cap, which #240 chunks under.
671            if writes.len() > crate::atproto::APPLY_WRITES_MAX_OPS {
672                return Err(Fail {
673                    status: 400,
674                    error: "InvalidRequest",
675                });
676            }
677            if let Some(fail) = self.always_fail {
678                return Err(fail);
679            }
680            if let Some(landed) = self.partial_then_503.take() {
681                for w in writes.iter().take(landed) {
682                    let (_, rkey, value) = Self::op(w);
683                    self.records.insert(rkey, value);
684                }
685                return Err(Fail {
686                    status: 503,
687                    error: "PartitionUnavailable",
688                });
689            }
690            // All or nothing, like the reference: every op is checked before
691            // any lands.
692            for w in writes {
693                let (action, rkey, _) = Self::op(w);
694                let exists = self.records.contains_key(&rkey);
695                if (action == "create" && exists) || (action == "update" && !exists) {
696                    return Err(MISMATCH);
697                }
698            }
699            for w in writes {
700                let (_, rkey, value) = Self::op(w);
701                self.records.insert(rkey, value);
702            }
703            Ok(())
704        }
705
706        fn page(&mut self, limit: Option<usize>, cursor: Option<&str>) -> serde_json::Value {
707            if cursor.is_none() {
708                self.list_walks += 1;
709            }
710            let limit = limit.unwrap_or(50);
711            let after: Vec<_> = self
712                .records
713                .iter()
714                .filter(|(rkey, _)| cursor.is_none_or(|c| rkey.as_str() > c))
715                .collect();
716            let page: Vec<_> = after.iter().take(limit).collect();
717            let records: Vec<serde_json::Value> = page
718                .iter()
719                .map(|(rkey, value)| {
720                    serde_json::json!({
721                        "uri": format!("at://{DID}/{}/{rkey}", crate::lexicon::nsid::READ_STATE),
722                        "cid": "bafyreigh2akiscaildc",
723                        "value": value,
724                    })
725                })
726                .collect();
727            let mut body = serde_json::json!({ "records": records });
728            if after.len() > limit {
729                body["cursor"] = serde_json::json!(page.last().unwrap().0);
730            }
731            body
732        }
733    }
734
735    /// Abandon the request mid-flight, so the client sees the connection close
736    /// with no response — a transport failure, not an answer. Unwinding the
737    /// handler drops hyper's connection task; `resume_unwind` skips the panic
738    /// hook, so it does not read as a test failure in the output. The caller
739    /// releases the fake's lock first, or the unwind would poison it.
740    fn hang_up() -> axum::response::Response {
741        std::panic::resume_unwind(Box::new("fake PDS hung up"))
742    }
743
744    /// Serve `fake` as BOTH a sidecar (`/internal/repo`) and a PDS (`/xrpc/*`),
745    /// so one fixture drives either backend. Returns the sidecar base URL and
746    /// the PDS audience a session should carry.
747    async fn serve_fake(fake: Arc<Mutex<FakeRepo>>) -> (String, String) {
748        use axum::http::StatusCode;
749        use axum::response::IntoResponse as _;
750
751        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
752        let addr = listener.local_addr().unwrap();
753        let port = addr.port();
754        // Unique per server: the override table is process-wide.
755        let host = format!("pds-{port}.readstate.test");
756        crate::net::test_host_override(&host, addr);
757
758        let app = axum::Router::new().fallback(move |req: axum::extract::Request| {
759            let fake = Arc::clone(&fake);
760            async move {
761                let (parts, body) = req.into_parts();
762                let raw = axum::body::to_bytes(body, usize::MAX).await.unwrap();
763                let query: std::collections::HashMap<String, String> =
764                    url::form_urlencoded::parse(parts.uri.query().unwrap_or("").as_bytes())
765                        .into_owned()
766                        .collect();
767                let reply = |status: u16, body: serde_json::Value| {
768                    (StatusCode::from_u16(status).unwrap(), axum::Json(body)).into_response()
769                };
770                let mut fake = fake.lock().unwrap();
771                match parts.uri.path() {
772                    "/internal/repo" => {
773                        let req: serde_json::Value = serde_json::from_slice(&raw).unwrap();
774                        match req["action"].as_str() {
775                            Some("list") => {
776                                let page = fake.page(
777                                    req["limit"].as_u64().map(|l| l as usize),
778                                    req["cursor"].as_str(),
779                                );
780                                if fake.list_fails {
781                                    return reply(
782                                        503,
783                                        serde_json::json!({ "ok": false, "error": "PartitionUnavailable", "status": 503 }),
784                                    );
785                                }
786                                reply(200, serde_json::json!({ "ok": true, "data": page }))
787                            }
788                            Some("applyWrites") => {
789                                match fake.apply(req["writes"].as_array().unwrap()) {
790                                    Err(f) if f.status == HANG_UP.status => {
791                                        drop(fake);
792                                        hang_up()
793                                    }
794                                    Ok(()) => {
795                                        reply(200, serde_json::json!({ "ok": true, "data": {} }))
796                                    }
797                                    // The sidecar's envelope, carrying the PDS's
798                                    // status and error name through.
799                                    Err(f) => reply(
800                                        f.status,
801                                        serde_json::json!({
802                                            "ok": false,
803                                            "error": f.error,
804                                            "message": "Internal Server Error",
805                                            "status": f.status,
806                                        }),
807                                    ),
808                                }
809                            }
810                            other => panic!("unexpected sidecar action {other:?}"),
811                        }
812                    }
813                    "/xrpc/com.atproto.repo.listRecords" => {
814                        let page = fake.page(
815                            query.get("limit").and_then(|l| l.parse().ok()),
816                            query.get("cursor").map(String::as_str),
817                        );
818                        if fake.list_fails {
819                            return reply(503, serde_json::json!({ "error": "PartitionUnavailable" }));
820                        }
821                        reply(200, page)
822                    }
823                    "/xrpc/com.atproto.repo.applyWrites" => {
824                        let req: serde_json::Value = serde_json::from_slice(&raw).unwrap();
825                        match fake.apply(req["writes"].as_array().unwrap()) {
826                            Err(f) if f.status == HANG_UP.status => {
827                                drop(fake);
828                                hang_up()
829                            }
830                            Ok(()) => reply(200, serde_json::json!({ "results": [] })),
831                            Err(f) => reply(
832                                f.status,
833                                serde_json::json!({
834                                    "error": f.error,
835                                    "message": "Internal Server Error",
836                                }),
837                            ),
838                        }
839                    }
840                    other => panic!("unexpected request to {other}"),
841                }
842            }
843        });
844        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
845        (format!("http://{addr}"), format!("http://{host}:{port}"))
846    }
847
848    /// An `AppState` on `backend`, pointed at the fake — the sidecar through its
849    /// internal URL, the Rust client through a live session whose `aud` is it.
850    pub(crate) async fn state_on(backend: Backend, fake: &Arc<Mutex<FakeRepo>>) -> AppState {
851        let (sidecar, aud) = serve_fake(Arc::clone(fake)).await;
852        let db = store::init_url("sqlite::memory:").await.unwrap();
853        let state = AppState::new(
854            crate::config::Config {
855                repo_backend: backend,
856                public_url: "http://localhost:8080".into(),
857                sidecar: crate::config::SidecarConfig {
858                    public_url: sidecar.clone(),
859                    internal_url: sidecar,
860                    internal_secret: "test-secret".into(),
861                },
862                oauth: crate::config::OauthConfig {
863                    // Per test, never the relative default — see `repo::tests`.
864                    key_path: std::env::temp_dir().join(format!(
865                        "fr-readstate-oauth-key-{}-{:p}.json",
866                        std::process::id(),
867                        &db as *const _
868                    )),
869                    encryption_key: Some("a".repeat(43)),
870                    ..crate::config::OauthConfig::default()
871                },
872                ..crate::config::Config::default()
873            },
874            db,
875        )
876        .unwrap();
877        if backend == Backend::Rust {
878            let runtime = state.oauth.as_deref().expect("oauth runtime");
879            crate::oauth::store::put_session(
880                &state.db,
881                &runtime.codec,
882                &crate::oauth::store::OAuthSession {
883                    sub: DID.into(),
884                    issuer: "https://auth.invalid".into(),
885                    aud,
886                    dpop_key_jwk: crate::oauth::keys::SigningKey::generate("session-dpop")
887                        .to_jwk_json()
888                        .unwrap(),
889                    access_token: "at".into(),
890                    refresh_token: "rt".into(),
891                    token_type: "DPoP".into(),
892                    granted_scope: "atproto".into(),
893                    expires_at: Some(store::now_unix() + 3600),
894                },
895            )
896            .await
897            .unwrap();
898        }
899        state
900    }
901
902    pub(crate) fn feed(i: usize) -> String {
903        format!("https://f{i}.example/feed.xml")
904    }
905
906    /// A dirty cursor for feed `i`, as a mark-read leaves it. `updated_at`
907    /// differs per call so the conditional dirty-clear sees a new snapshot.
908    pub(crate) async fn mark_read(state: &AppState, i: usize, id: &str) {
909        static TICK: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(0);
910        let tick = TICK.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
911        store::upsert_cursor(
912            &state.db,
913            &ReadCursor {
914                did: DID.into(),
915                feed_url: feed(i),
916                read_through: None,
917                read_ids: format!("[\"{id}\"]"),
918                unread_ids: "[]".into(),
919                dirty: true,
920                pds_created: false,
921                updated_at: format!(
922                    "2026-10-04T{:02}:{:02}:{:02}Z",
923                    tick / 3600 % 24,
924                    tick / 60 % 60,
925                    tick % 60
926                ),
927            },
928        )
929        .await
930        .unwrap();
931    }
932
933    /// A record already in the fake repo for feed `i`, from some earlier flush.
934    fn existing(fake: &Arc<Mutex<FakeRepo>>, i: usize) {
935        let record = ReadState::new(feed(i), None, "2026-01-01T00:00:00Z");
936        fake.lock().unwrap().records.insert(
937            read_state_rkey(&feed(i)),
938            serde_json::to_value(record).unwrap(),
939        );
940    }
941
942    pub(crate) async fn cursor(state: &AppState, i: usize) -> ReadCursor {
943        store::get_cursor(&state.db, DID, &feed(i))
944            .await
945            .unwrap()
946            .expect("cursor row")
947    }
948
949    fn read_ids_on_pds(fake: &Arc<Mutex<FakeRepo>>, i: usize) -> serde_json::Value {
950        fake.lock().unwrap().records[&read_state_rkey(&feed(i))]["readIds"].clone()
951    }
952
953    const BACKENDS: [Backend; 2] = [Backend::Sidecar, Backend::Rust];
954
955    /// **(1) A lost success response converges.** The applyWrites that created
956    /// the record landed, but its answer never arrived, so `pds_created` is still
957    /// false and every later flush sends `#create` at an rkey that exists. Before
958    /// the fix that batch failed forever and took the DID's whole read-state
959    /// sync with it.
960    #[tokio::test]
961    async fn a_lost_create_response_converges_on_the_next_flush() {
962        for backend in BACKENDS {
963            let fake = Arc::new(Mutex::new(FakeRepo::default()));
964            let state = state_on(backend, &fake).await;
965            existing(&fake, 1);
966            mark_read(&state, 1, "42").await;
967
968            flush_did(&state, DID)
969                .await
970                .unwrap_or_else(|e| panic!("{backend:?}: the flush stayed wedged: {e:#}"));
971
972            let c = cursor(&state, 1).await;
973            assert!(c.pds_created, "{backend:?}: pds_created was not corrected");
974            assert!(!c.dirty, "{backend:?}: the cursor is still dirty");
975            assert_eq!(
976                read_ids_on_pds(&fake, 1),
977                serde_json::json!(["42"]),
978                "{backend:?}"
979            );
980            {
981                let f = fake.lock().unwrap();
982                assert_eq!(
983                    f.list_walks, 1,
984                    "{backend:?}: expected exactly one reconcile"
985                );
986                assert_eq!(
987                    f.apply_calls, 2,
988                    "{backend:?}: the failed batch, then the retry"
989                );
990            }
991
992            // Steady state again: the next flush is a plain update, no listing.
993            mark_read(&state, 1, "43").await;
994            flush_did(&state, DID).await.expect("steady-state flush");
995            let f = fake.lock().unwrap();
996            assert_eq!(f.list_walks, 1, "{backend:?}: a healthy flush reconciled");
997            assert_eq!(f.apply_calls, 3, "{backend:?}");
998        }
999    }
1000
1001    /// **(2) A fresh or restored database against a repo that already holds
1002    /// the records.** Several feeds among records of others, more than one
1003    /// listing page of them, so the reconcile has to walk.
1004    #[tokio::test]
1005    async fn a_fresh_database_converges_against_existing_records() {
1006        for backend in BACKENDS {
1007            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1008            let state = state_on(backend, &fake).await;
1009            for i in 0..150 {
1010                existing(&fake, i);
1011            }
1012            for i in [3, 77, 149] {
1013                mark_read(&state, i, "7").await;
1014            }
1015
1016            flush_did(&state, DID)
1017                .await
1018                .unwrap_or_else(|e| panic!("{backend:?}: {e:#}"));
1019
1020            for i in [3, 77, 149] {
1021                let c = cursor(&state, i).await;
1022                assert!(
1023                    c.pds_created && !c.dirty,
1024                    "{backend:?}: feed {i} did not converge"
1025                );
1026                assert_eq!(
1027                    read_ids_on_pds(&fake, i),
1028                    serde_json::json!(["7"]),
1029                    "{backend:?}"
1030                );
1031            }
1032            let f = fake.lock().unwrap();
1033            assert_eq!(f.list_walks, 1, "{backend:?}");
1034            assert_eq!(f.apply_calls, 2, "{backend:?}");
1035            assert_eq!(
1036                f.records.len(),
1037                150,
1038                "{backend:?}: a record was duplicated or lost"
1039            );
1040        }
1041    }
1042
1043    /// **(3) The mirror: `pds_created` is set but the record was deleted
1044    /// elsewhere** (another client, a repo reset). `#update` on a missing rkey
1045    /// fails the same way, and converges by creating it.
1046    #[tokio::test]
1047    async fn a_deleted_record_converges_by_creating_it() {
1048        for backend in BACKENDS {
1049            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1050            let state = state_on(backend, &fake).await;
1051            mark_read(&state, 1, "9").await;
1052            store::mark_cursor_pds_created(&state.db, DID, &feed(1))
1053                .await
1054                .unwrap();
1055
1056            flush_did(&state, DID)
1057                .await
1058                .unwrap_or_else(|e| panic!("{backend:?}: {e:#}"));
1059
1060            let c = cursor(&state, 1).await;
1061            assert!(c.pds_created && !c.dirty, "{backend:?}");
1062            assert_eq!(
1063                read_ids_on_pds(&fake, 1),
1064                serde_json::json!(["9"]),
1065                "{backend:?}"
1066            );
1067            let f = fake.lock().unwrap();
1068            assert_eq!((f.list_walks, f.apply_calls), (1, 2), "{backend:?}");
1069        }
1070    }
1071
1072    /// **(4) Part of a batch landed.** With applyWrites chunked (#240), a later
1073    /// chunk can fail after an earlier one committed — exactly this issue's
1074    /// state for every cursor in the landed chunk. That first failure is a 503,
1075    /// which is NOT reconciled; the next flush meets the half-landed batch, mixed
1076    /// with a flag that is wrong the other way, and converges in one reconcile.
1077    #[tokio::test]
1078    async fn a_partially_landed_batch_converges() {
1079        for backend in BACKENDS {
1080            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1081            let state = state_on(backend, &fake).await;
1082            for i in 1..=4 {
1083                mark_read(&state, i, "5").await;
1084            }
1085            // Feed 4 claims a record that is not there.
1086            store::mark_cursor_pds_created(&state.db, DID, &feed(4))
1087                .await
1088                .unwrap();
1089            fake.lock().unwrap().partial_then_503 = Some(2);
1090
1091            flush_did(&state, DID)
1092                .await
1093                .expect_err("the 503 must surface as a failure");
1094            {
1095                let f = fake.lock().unwrap();
1096                assert_eq!(
1097                    f.list_walks, 0,
1098                    "{backend:?}: a 503 is not a mismatch and must not reconcile"
1099                );
1100                assert_eq!(f.records.len(), 2, "{backend:?}: fixture");
1101            }
1102
1103            flush_did(&state, DID)
1104                .await
1105                .unwrap_or_else(|e| panic!("{backend:?}: {e:#}"));
1106            for i in 1..=4 {
1107                let c = cursor(&state, i).await;
1108                assert!(c.pds_created && !c.dirty, "{backend:?}: feed {i}");
1109                assert_eq!(
1110                    read_ids_on_pds(&fake, i),
1111                    serde_json::json!(["5"]),
1112                    "{backend:?}"
1113                );
1114            }
1115            let f = fake.lock().unwrap();
1116            assert_eq!(f.list_walks, 1, "{backend:?}");
1117            assert_eq!(f.apply_calls, 3, "{backend:?}: 503, mismatch, retry");
1118        }
1119    }
1120
1121    /// **(5) An unrelated failure behaves exactly as before:** the error comes
1122    /// back, the cursors stay dirty, and nothing is listed. A reconcile on every
1123    /// outage would add a repo walk per DID per minute to a PDS that is already
1124    /// struggling.
1125    #[tokio::test]
1126    async fn an_unrelated_failure_is_returned_without_a_reconcile() {
1127        for backend in BACKENDS {
1128            for fail in [
1129                Fail {
1130                    status: 503,
1131                    error: "PartitionUnavailable",
1132                },
1133                Fail {
1134                    status: 502,
1135                    error: "UpstreamFailure",
1136                },
1137                Fail {
1138                    status: 401,
1139                    error: "AuthRequired",
1140                },
1141                Fail {
1142                    status: 429,
1143                    error: "RateLimitExceeded",
1144                },
1145            ] {
1146                let fake = Arc::new(Mutex::new(FakeRepo::default()));
1147                let state = state_on(backend, &fake).await;
1148                // A real mismatch is present too, so a reconcile WOULD find
1149                // something to fix — what is under test is that it is not tried.
1150                existing(&fake, 1);
1151                mark_read(&state, 1, "1").await;
1152                fake.lock().unwrap().always_fail = Some(fail);
1153
1154                flush_did(&state, DID)
1155                    .await
1156                    .expect_err("the failure must be returned");
1157
1158                let c = cursor(&state, 1).await;
1159                assert!(c.dirty, "{backend:?} {fail:?}: the reads were dropped");
1160                assert!(!c.pds_created, "{backend:?} {fail:?}: the flag moved");
1161                let f = fake.lock().unwrap();
1162                assert_eq!(
1163                    f.list_walks, 0,
1164                    "{backend:?} {fail:?}: reconciled an unrelated failure"
1165                );
1166                assert_eq!(
1167                    f.apply_calls, 1,
1168                    "{backend:?} {fail:?}: retried an unrelated failure"
1169                );
1170            }
1171        }
1172    }
1173
1174    /// **(6) At most one reconcile per flush — never a loop.** The PDS refuses
1175    /// every write with the mismatch-shaped error, so the retry fails too: the
1176    /// flush must give up after one listing and one retry, keep the cursor
1177    /// dirty, and still record the truth it learned.
1178    #[tokio::test]
1179    async fn a_flush_reconciles_at_most_once() {
1180        for backend in BACKENDS {
1181            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1182            let state = state_on(backend, &fake).await;
1183            existing(&fake, 1);
1184            mark_read(&state, 1, "1").await;
1185            fake.lock().unwrap().always_fail = Some(MISMATCH);
1186
1187            let err = flush_did(&state, DID)
1188                .await
1189                .expect_err("a retry that fails again is a failure");
1190            // Both callers log `%err` — the plain Display, which for an anyhow
1191            // context is the OUTERMOST message only. The PDS's reason for the
1192            // second refusal must be in it, or the log says "failed again"
1193            // and nothing about why.
1194            let shown = err.to_string();
1195            assert!(
1196                shown.contains("after reconciling")
1197                    && shown.contains("500")
1198                    && shown.contains("InternalServerError"),
1199                "{backend:?}: the logged line lost the PDS's reason: {shown}"
1200            );
1201
1202            let c = cursor(&state, 1).await;
1203            assert!(c.dirty, "{backend:?}: the reads were dropped");
1204            assert!(
1205                c.pds_created,
1206                "{backend:?}: the truth the listing found was not kept"
1207            );
1208            let f = fake.lock().unwrap();
1209            assert_eq!(f.list_walks, 1, "{backend:?}: reconciled more than once");
1210            assert_eq!(f.apply_calls, 2, "{backend:?}: retried more than once");
1211        }
1212    }
1213
1214    /// **Which refusals earn a reconcile, pinned per arm.** The end-to-end tests
1215    /// above use the reference PDS's 500 and four unrelated statuses; the 400
1216    /// and 409 arms, and the rule that only a STRUCTURED rejection counts, are
1217    /// asserted here.
1218    #[test]
1219    fn only_a_structured_conflict_shaped_rejection_may_be_a_mismatch() {
1220        let xrpc = |status: u16, error: &str| -> anyhow::Error {
1221            AtProtoError::Xrpc {
1222                status: reqwest::StatusCode::from_u16(status).unwrap(),
1223                error: error.into(),
1224                message: None,
1225            }
1226            .into()
1227        };
1228        for (status, error) in [
1229            (500, "InternalServerError"),
1230            (400, "InvalidRequest"),
1231            (400, "InvalidSwap"),
1232            (400, "RecordNotFound"),
1233            (409, "Conflict"),
1234        ] {
1235            assert!(
1236                may_be_existence_mismatch(&xrpc(status, error)),
1237                "{status} {error} should reconcile"
1238            );
1239            // Through a context layer, as the Rust client wraps it.
1240            assert!(
1241                may_be_existence_mismatch(&xrpc(status, error).context("applyWrites failed")),
1242                "{status} {error} behind a context should reconcile"
1243            );
1244        }
1245        for (status, error) in [
1246            (500, "Unknown"),
1247            (400, "ExpiredToken"),
1248            (401, "AuthRequired"),
1249            (403, "CollectionNotAllowed"),
1250            (404, "SessionNotFound"),
1251            (429, "RateLimitExceeded"),
1252            (502, "UpstreamFailure"),
1253            (503, "StoreUnavailable"),
1254            (504, "UpstreamTimeout"),
1255        ] {
1256            assert!(
1257                !may_be_existence_mismatch(&xrpc(status, error)),
1258                "{status} {error} must not reconcile"
1259            );
1260        }
1261        // Words are not a rejection: only the structured error counts.
1262        assert!(!may_be_existence_mismatch(&anyhow::anyhow!(
1263            "applyWrites failed: status 500 (InternalServerError)"
1264        )));
1265    }
1266
1267    /// **A reconcile that cannot list gives up, and says what the WRITE said.**
1268    /// The listing is the only evidence a retry would be different; without it
1269    /// the batch is not resent, the cursor stays dirty, and the error returned
1270    /// is the applyWrites failure, not the listing's.
1271    #[tokio::test]
1272    async fn a_reconcile_that_cannot_list_returns_the_original_error() {
1273        for backend in BACKENDS {
1274            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1275            let state = state_on(backend, &fake).await;
1276            existing(&fake, 1);
1277            mark_read(&state, 1, "1").await;
1278            fake.lock().unwrap().list_fails = true;
1279
1280            let err = flush_did(&state, DID)
1281                .await
1282                .expect_err("an unreconciled mismatch is still a failure");
1283            assert!(
1284                may_be_existence_mismatch(&err),
1285                "{backend:?}: returned the listing's error, not the write's: {err:#}"
1286            );
1287            let c = cursor(&state, 1).await;
1288            assert!(c.dirty && !c.pds_created, "{backend:?}");
1289            let f = fake.lock().unwrap();
1290            assert_eq!((f.list_walks, f.apply_calls), (1, 1), "{backend:?}");
1291        }
1292    }
1293
1294    /// **A 500 the listing cannot explain is not retried.** The reference PDS
1295    /// reports a mismatch as a bare 500, so a 500 earns one listing — but when
1296    /// every flag already matches the repo, resending the identical batch could
1297    /// only fail the same way, so the original error comes back unchanged.
1298    #[tokio::test]
1299    async fn a_500_with_no_mismatch_is_not_retried() {
1300        for backend in BACKENDS {
1301            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1302            let state = state_on(backend, &fake).await;
1303            // pds_created=false and no record: the flag is already right.
1304            mark_read(&state, 1, "1").await;
1305            fake.lock().unwrap().always_fail = Some(MISMATCH);
1306
1307            flush_did(&state, DID)
1308                .await
1309                .expect_err("the 500 must be returned");
1310
1311            let c = cursor(&state, 1).await;
1312            assert!(c.dirty && !c.pds_created, "{backend:?}");
1313            let f = fake.lock().unwrap();
1314            assert_eq!((f.list_walks, f.apply_calls), (1, 1), "{backend:?}");
1315        }
1316    }
1317
1318    // ── #240 × #241: a flush split into several applyWrites calls ────────────
1319
1320    /// Feeds `0..n` in the order the flush sends them: the batch is keyed by
1321    /// rkey, so a split's first call holds the lowest rkeys.
1322    pub(crate) fn send_order(n: usize) -> Vec<usize> {
1323        let mut order: Vec<usize> = (0..n).collect();
1324        order.sort_by_key(|&i| read_state_rkey(&feed(i)));
1325        order
1326    }
1327
1328    /// **(a) Call 1 lands, call 2 hits a transport failure.** Call 1's creates
1329    /// are committed, so its cursors must be marked created and clean NOW. Left
1330    /// as they were, the next round re-sends them as `#create`, the PDS refuses,
1331    /// and — the calls stop at the first failure, in a fixed rkey order — every
1332    /// cursor sorted after them starves, every round. A transport failure is not
1333    /// a mismatch, so nothing is listed; the next flush sends only what did not
1334    /// land, plus a re-dirtied landed cursor as an UPDATE.
1335    #[tokio::test]
1336    async fn a_split_flush_keeps_what_landed_when_a_later_call_fails() {
1337        for backend in BACKENDS {
1338            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1339            let state = state_on(backend, &fake).await;
1340            for i in 0..250 {
1341                mark_read(&state, i, "1").await;
1342            }
1343            fake.lock().unwrap().drop_call = Some(2);
1344
1345            let err = flush_did(&state, DID)
1346                .await
1347                .expect_err("the second call failed");
1348            // What callers log (`%err`) says how far the split got.
1349            assert!(
1350                err.to_string().contains("200 of 250 writes had landed"),
1351                "{backend:?}: the logged line lost the progress: {err}"
1352            );
1353
1354            let order = send_order(250);
1355            let (landed, rest) = order.split_at(crate::atproto::APPLY_WRITES_MAX_OPS);
1356            for &i in landed {
1357                let c = cursor(&state, i).await;
1358                assert!(
1359                    c.pds_created && !c.dirty,
1360                    "{backend:?}: feed {i} landed in call 1 but was not settled"
1361                );
1362            }
1363            for &i in rest {
1364                let c = cursor(&state, i).await;
1365                assert!(
1366                    c.dirty && !c.pds_created,
1367                    "{backend:?}: feed {i} never landed"
1368                );
1369            }
1370            {
1371                let f = fake.lock().unwrap();
1372                assert_eq!(f.records.len(), 200, "{backend:?}: fixture");
1373                assert_eq!((f.list_walks, f.apply_calls), (0, 2), "{backend:?}");
1374            }
1375
1376            // A landed feed is read again: it must go as #update now — the
1377            // fake refuses a #create on an existing key.
1378            mark_read(&state, landed[0], "2").await;
1379            flush_did(&state, DID)
1380                .await
1381                .unwrap_or_else(|e| panic!("{backend:?}: the remainder did not flush: {e:#}"));
1382            for i in 0..250 {
1383                let c = cursor(&state, i).await;
1384                assert!(c.pds_created && !c.dirty, "{backend:?}: feed {i}");
1385            }
1386            assert_eq!(
1387                read_ids_on_pds(&fake, landed[0]),
1388                serde_json::json!(["2"]),
1389                "{backend:?}"
1390            );
1391            let f = fake.lock().unwrap();
1392            assert_eq!(f.records.len(), 250, "{backend:?}");
1393            assert_eq!(
1394                (f.list_walks, f.apply_calls),
1395                (0, 3),
1396                "{backend:?}: 51 writes are one call, and nothing needed listing"
1397            );
1398        }
1399    }
1400
1401    /// **(b) Call 1 lands, call 2 meets records that already exist.** The
1402    /// mismatch arrives wrapped in the chunked write's progress error, and must
1403    /// still be recognised: one listing, the flags corrected, and ONE retry of
1404    /// what did not land — not of the 200 that did.
1405    #[tokio::test]
1406    async fn a_mismatch_in_a_later_call_reconciles_and_retries_only_the_rest() {
1407        for backend in BACKENDS {
1408            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1409            let state = state_on(backend, &fake).await;
1410            let order = send_order(250);
1411            for &i in &order[240..] {
1412                existing(&fake, i);
1413            }
1414            for i in 0..250 {
1415                mark_read(&state, i, "3").await;
1416            }
1417
1418            flush_did(&state, DID)
1419                .await
1420                .unwrap_or_else(|e| panic!("{backend:?}: did not converge: {e:#}"));
1421
1422            for i in 0..250 {
1423                let c = cursor(&state, i).await;
1424                assert!(c.pds_created && !c.dirty, "{backend:?}: feed {i}");
1425                assert_eq!(
1426                    read_ids_on_pds(&fake, i),
1427                    serde_json::json!(["3"]),
1428                    "{backend:?}: feed {i}"
1429                );
1430            }
1431            let f = fake.lock().unwrap();
1432            assert_eq!(f.records.len(), 250, "{backend:?}");
1433            assert_eq!(f.list_walks, 1, "{backend:?}: expected one reconcile");
1434            assert_eq!(
1435                f.apply_calls, 3,
1436                "{backend:?}: call 1, call 2 refused, then one retry of the 50 left"
1437            );
1438        }
1439    }
1440
1441    /// **The classifier sees through the chunked write's progress error.**
1442    /// `ApplyWritesIncomplete::source` continues from its cause's SOURCE, so a
1443    /// cause that IS the `AtProtoError` — the sidecar client's shape — does not
1444    /// appear in `err.chain()` at all. Both shapes, through the real wrapper.
1445    #[tokio::test]
1446    async fn a_mismatch_is_recognised_inside_a_chunked_write_error() {
1447        let op = crate::atproto::WriteOp::Delete {
1448            collection: crate::lexicon::nsid::READ_STATE.into(),
1449            rkey: "rs-0".into(),
1450        };
1451        let mismatch = || AtProtoError::Xrpc {
1452            status: reqwest::StatusCode::INTERNAL_SERVER_ERROR,
1453            error: "InternalServerError".into(),
1454            message: None,
1455        };
1456        // Sidecar shape: the cause is the AtProtoError itself.
1457        let bare = crate::atproto::apply_writes_chunked(std::slice::from_ref(&op), |_| async {
1458            Err(mismatch().into())
1459        })
1460        .await
1461        .unwrap_err();
1462        assert!(
1463            crate::atproto::ApplyWritesIncomplete::of(&bare).is_some(),
1464            "fixture: not wrapped"
1465        );
1466        assert!(may_be_existence_mismatch(&bare), "sidecar shape: {bare:#}");
1467        // Rust-client shape: the AtProtoError behind a context.
1468        let wrapped = crate::atproto::apply_writes_chunked(std::slice::from_ref(&op), |_| async {
1469            Err(anyhow::Error::new(mismatch()).context("applyWrites failed"))
1470        })
1471        .await
1472        .unwrap_err();
1473        assert!(
1474            may_be_existence_mismatch(&wrapped),
1475            "rust shape: {wrapped:#}"
1476        );
1477        // And an unrelated failure, wrapped, is still unrelated.
1478        let transport =
1479            crate::atproto::apply_writes_chunked(std::slice::from_ref(&op), |_| async {
1480                Err(anyhow::anyhow!("connection reset"))
1481            })
1482            .await
1483            .unwrap_err();
1484        assert!(!may_be_existence_mismatch(&transport));
1485    }
1486
1487    /// **The retry can part-land too.** Call 1 is refused outright (it meets
1488    /// records that exist), the reconcile corrects them, and the retry — 250
1489    /// writes, so two calls — loses its second call's connection. The retry's
1490    /// landed prefix must be settled exactly like the first attempt's, or its
1491    /// 190 fresh creates go out as `#create` again next round.
1492    #[tokio::test]
1493    async fn a_retry_that_part_lands_settles_what_it_landed() {
1494        for backend in BACKENDS {
1495            let fake = Arc::new(Mutex::new(FakeRepo::default()));
1496            let state = state_on(backend, &fake).await;
1497            let order = send_order(250);
1498            for &i in &order[..10] {
1499                existing(&fake, i);
1500            }
1501            for i in 0..250 {
1502                mark_read(&state, i, "4").await;
1503            }
1504            // Call 1 refused (mismatch), call 2 is the retry's first, call 3
1505            // the retry's second.
1506            fake.lock().unwrap().drop_call = Some(3);
1507
1508            let err = flush_did(&state, DID)
1509                .await
1510                .expect_err("the retry's second call failed");
1511            assert!(
1512                err.to_string().contains("after reconciling"),
1513                "{backend:?}: {err}"
1514            );
1515
1516            let (landed, rest) = order.split_at(crate::atproto::APPLY_WRITES_MAX_OPS);
1517            for &i in landed {
1518                let c = cursor(&state, i).await;
1519                assert!(
1520                    c.pds_created && !c.dirty,
1521                    "{backend:?}: feed {i} landed in the retry but was not settled"
1522                );
1523            }
1524            for &i in rest {
1525                let c = cursor(&state, i).await;
1526                assert!(c.dirty && !c.pds_created, "{backend:?}: feed {i}");
1527            }
1528            let f = fake.lock().unwrap();
1529            assert_eq!((f.list_walks, f.apply_calls), (1, 3), "{backend:?}");
1530        }
1531    }
1532}