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}