From ae292e41eb7da0b9196110f18dfa6984b5aee07d Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Wed, 26 Aug 2026 11:29:59 -0400 Subject: [PATCH] fix(pds): one backwards walk, and it says when it did not reach the bottom MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `state_events_for` and `list_statuses` were the same walk twice, and neither copy reported hitting its page cap — so the cap read as the end of the collection, and a `merged` status one page past it read as `open`. Both are `records_about` now, which returns `complete`: `pr close`/`reopen` refuses a state it cannot settle, and the issue listing says so and carries on. Change-Id: If377168fa95485aeb5ff16d171a922c9b2a05ab6 --- plan/module-layout.md | 31 ++-- src/clients/atproto/pds.rs | 341 ++++++++++++++++++++++++++++++++++++- src/cmd/issue/read.rs | 82 +++------ src/cmd/pr/write.rs | 123 ++++++------- 4 files changed, 448 insertions(+), 129 deletions(-) diff --git a/plan/module-layout.md b/plan/module-layout.md index 80eb70f..0649c63 100644 --- a/plan/module-layout.md +++ b/plan/module-layout.md @@ -31,17 +31,6 @@ call sites are named wrappers whose doc comments say what each one's nested ## What it needs -- [ ] One backwards walk over a subject's records, not two. - `cmd/issue/read.rs`'s `state_events` and `cmd/pr/write.rs`'s - `list_statuses` are the same function: page `listRecords` newest-first, - keep the records naming this subject, and stop once a page ends below - the subject's own rkey, because a record about a thing cannot predate - it. They differ in the NSID, the subject field (`issue`/`pull`), the - state field (`state`/`status`), the struct built, and a 50-page cap - spelled twice under two names. The comparison is strict in both, for the - same unobvious reason — Tangled's backfill wrote state records under the - *subject's* own key and those must still be seen — and a rule that - subtle stated twice is a rule that will be fixed once - [ ] Not to be merged with it: `newest_state` and `state_of` look like the same pair and are deliberately not. The issue half orders on a parsed instant because it merges the author's records with the acting account's @@ -67,6 +56,26 @@ call sites are named wrappers whose doc comments say what each one's nested ## Done +- [x] One backwards walk over a subject's records, not two. + `cmd/issue/read.rs`'s `state_events_for` and `cmd/pr/write.rs`'s + `list_statuses` were the same function — page `listRecords` + newest-first, keep the records naming this subject, stop once a page + ends below the subject's own rkey — differing in the NSID, the subject + field, the state field, the struct built, and a 50-page cap spelled + twice under two names. They are `pds::records_about` now, and the + strict comparison that keeps a backfilled record written under the + subject's own key visible is stated once, in `page_ended_below`, with + the test that would fail if somebody tidied it to `<=`. + The move turned up what the duplication was hiding. Neither copy said + whether it had reached the bottom, so both read the page cap as the end + of the collection — and the absence of a record is not nothing here, it + is what "still open" is read off. A `merged` status one page past the + cap therefore read as `open`, and `pr close` would write over it + without the merged-is-terminal refusal ever firing. `records_about` + reports `complete`; the write path refuses on a short walk and the + issue listing says so and carries on, which is the difference between a + state that decides a write and a state that fills a column + - [x] `html/` — a folder for the documents atgc serves a browser, sibling to `term/`, and `art.rs` at the root for the field of DNA both media draw. The OAuth callback pages left `clients/atproto/oauth/`, where markup and diff --git a/src/clients/atproto/pds.rs b/src/clients/atproto/pds.rs index 8fd2a8d..ccacbe6 100644 --- a/src/clients/atproto/pds.rs +++ b/src/clients/atproto/pds.rs @@ -80,6 +80,186 @@ pub const MAX_PAGES: usize = 5; /// which is every listing under a filter it cannot see through. pub const MAX_PAGES_COMPLETE: usize = 100; +/// How many pages to walk when looking backwards for records *about* +/// something. +/// +/// [`records_about`] stops on its own long before this for any real account, +/// because the subject's own record key is a floor. The cap is only here so +/// that a PDS handing back a cursor forever could not spin. +pub const MAX_PAGES_ABOUT: usize = 50; + +/// Whether a page of records ended below `floor`, meaning nothing older can +/// still be about the subjects being asked about. +/// +/// **The comparison is strict, and that is not a detail.** Tangled's own +/// backfill migration wrote state records under the *subject's* own record +/// key, so a record whose key equals the floor exists in the wild and has to +/// be seen. `<=` here would skip exactly those. +/// +/// A page whose last record carries no readable key compares as empty, which +/// is below every floor and stops the walk. Deliberately the conservative +/// direction: a page this cannot read is not evidence that there is more +/// worth crawling for. +fn page_ended_below(records: &[serde_json::Value], floor: &str) -> bool { + let oldest = records + .last() + .and_then(|r| r["uri"].as_str()) + .and_then(|u| u.rsplit('/').next()) + .unwrap_or_default(); + if oldest < floor { + crate::logging::debug::log(format!( + "stopping at rkey {oldest}, which predates the oldest subject key asked \ + about ({floor})" + )); + return true; + } + false +} + +/// Every record in `collection` naming one of `subjects` in its +/// `subject_field`, grouped by the subject at-uri it named. +/// +/// `listRecords` has no filter, so this pages and picks. The picking costs +/// nothing once a page is in hand — the round trips are the whole expense — +/// which is why a batch of subjects is asked about in one walk rather than +/// one walk each. A listing of thirty issues used to run this thirty times, +/// for pages that were very often the same bytes. +/// +/// What keeps the walk bounded is not [`MAX_PAGES_ABOUT`]. Record keys are +/// TIDs and `listRecords` returns them newest first, and a record *about* a +/// thing cannot have been written before the thing itself, so once a page +/// ends below the oldest subject key asked about there is nothing left to +/// find. Per subject that floor is its own key; for a batch it has to be the +/// minimum, or the walk would stop while a later subject's records were +/// still ahead of it. [`page_ended_below`] is the comparison, and says why +/// it is strict. +/// +/// `subjects` is `(subject at-uri, subject record key)` pairs. Records come +/// back raw, since what each caller wants out of one differs and the walk is +/// the part that was worth having once. +pub async fn records_about( + pds: &str, + did: &str, + collection: &str, + subject_field: &str, + subjects: &[(&str, &str)], +) -> Result { + let walk = walk_about(subjects, subject_field, |cursor| async move { + list_records(pds, did, collection, cursor.as_deref()).await + }) + .await?; + crate::logging::debug::log(format!( + "{} page(s) of {collection} from {did} covered {} subject(s)", + walk.pages, + subjects.len() + )); + if !walk.complete { + crate::logging::debug::log(format!( + "{MAX_PAGES_ABOUT} page(s) of {collection} from {did} went by without reaching \ + the oldest subject asked about: the walk is short" + )); + } + Ok(walk) +} + +/// The walk itself, over any source of pages. +/// +/// Split from [`records_about`] so the stopping rules can be proven without a +/// socket — which matters more here than the split usually does, because the +/// rule that decides `complete` is about pages that *did not arrive*, and the +/// mock PDS in `tests/support/` serves this collection oldest-first where a +/// real one serves it newest-first. A flow test against it would exercise the +/// walk backwards. See `plan/testing.md`. +async fn walk_about( + subjects: &[(&str, &str)], + subject_field: &str, + mut fetch: F, +) -> Result +where + F: FnMut(Option) -> Fut, + Fut: Future, Option)>>, +{ + let mut found: std::collections::HashMap> = + std::collections::HashMap::new(); + let Some(floor) = subjects.iter().map(|(_, rkey)| *rkey).min() else { + return Ok(RecordsAbout { + by_subject: found, + complete: true, + pages: 0, + }); + }; + let wanted: std::collections::HashSet<&str> = subjects.iter().map(|(uri, _)| *uri).collect(); + let mut cursor: Option = None; + let mut pages = 0; + // The walk is complete when it *ran out of reasons to keep going*: it + // reached the floor, or the collection ended. Falling out of the loop + // with neither having happened is the page cap, and the difference + // matters — see [`RecordsAbout::complete`]. + let mut complete = false; + + for _ in 0..MAX_PAGES_ABOUT { + let (records, next) = fetch(cursor.clone()).await?; + pages += 1; + if records.is_empty() { + complete = true; + break; + } + let done = page_ended_below(&records, floor); + for record in records { + let Some(subject) = record["value"][subject_field].as_str().map(str::to_string) else { + continue; + }; + if !wanted.contains(subject.as_str()) { + continue; + } + found.entry(subject).or_default().push(record); + } + if done { + complete = true; + break; + } + match next { + Some(c) => { + cursor = Some(c); + } + None => { + complete = true; + break; + } + } + } + Ok(RecordsAbout { + by_subject: found, + complete, + pages, + }) +} + +/// What one backwards walk found, and whether it saw the whole of it. +pub struct RecordsAbout { + /// The records naming each subject, by that subject's at-uri. A subject + /// nothing was written about is absent rather than empty. + pub by_subject: std::collections::HashMap>, + /// Whether the walk stopped because there was nothing left to find. + /// + /// **A short walk and an empty answer look identical, and they mean + /// opposite things.** The absence of a state record is itself an answer — + /// it is what "still open" is read off — so a caller that infers a + /// subject's state from what came back may only do so over a complete + /// walk. Both callers used to drop this on the floor and treat the page + /// cap as the end of the collection, which on the pull side meant a + /// `merged` record beyond the cap read as `open`, and the refusal that + /// exists to keep a merge from being overwritten never fired. + /// + /// `false` is rare and not hypothetical: it takes an account with + /// [`MAX_PAGES_ABOUT`] pages of this one collection written since the + /// subject, which is a busy repo and an old pull rather than anything + /// exotic. + pub complete: bool, + /// How many pages the walk actually asked for, for the debug log. + pages: usize, +} + /// One bounded GET at a PDS, logged, retried, with nothing read off it yet. /// /// The URL is [`crate::clients::xrpc::endpoint`]'s, which is where this @@ -522,7 +702,9 @@ pub fn rkey(record: &serde_json::Value) -> &str { #[cfg(test)] mod tests { - use super::{MAX_PAGES, Reach, page, record, rkey}; + use super::{ + MAX_PAGES, MAX_PAGES_ABOUT, Reach, page, page_ended_below, record, rkey, walk_about, + }; use serde_json::json; const DID: &str = "did:plc:nlzmjyfv6loqtxyzvdcznwgf"; @@ -661,4 +843,161 @@ mod tests { assert_eq!(rkey(&json!({ "uri": "" })), ""); assert_eq!(rkey(&json!({})), "?"); } + + /// A record written under the subject's *own* key is about the subject + /// and must still be seen. + /// + /// This is the whole reason the floor comparison is strict, and it is not + /// hypothetical: Tangled's backfill migration wrote pull status records + /// under the pull's own record key. The rule used to be spelled out in + /// prose in two modules and enforced by two copies of one `<`; it is one + /// function now, and this is the test that would fail if somebody + /// "tidied" it into `<=`. + #[test] + fn a_page_ending_exactly_at_the_floor_does_not_stop_the_walk() { + let floor = "3msg7wllcqc2b"; + let at = + |rkey: &str| json!({ "uri": format!("at://{DID}/sh.tangled.repo.pull.status/{rkey}") }); + + assert!( + !page_ended_below(&[at("3mtaaaaaaaaaa"), at(floor)], floor), + "a page ending at the subject's own key stopped the walk" + ); + // One key older than the floor, and there is nothing left to find. + assert!(page_ended_below(&[at(floor), at("3ma000000000")], floor)); + // A page this cannot read stops rather than crawls on. + assert!(page_ended_below(&[json!({})], floor)); + // And no page at all is not evidence of anything either way; the + // caller has already stopped on an empty page before asking. + assert!(page_ended_below(&[], floor)); + } + + // -- the backwards walk ------------------------------------------------ + + /// One page of `n` status-shaped records, keys descending from `from`, + /// which is the order a real `listRecords` serves them in. + fn page_of(subject: &str, from: usize, n: usize) -> Vec { + (0..n) + .map(|i| { + let key = format!("3k{:08}", from - i); + json!({ + "uri": format!("at://{DID}/sh.tangled.repo.pull.status/{key}"), + "value": { "pull": subject, "status": "open" }, + }) + }) + .collect() + } + + /// A walk that reaches the floor is complete, and stops there. + #[tokio::test] + async fn a_walk_that_reaches_the_floor_is_complete() { + let subject = "at://x/p/one"; + let mut asked = 0; + let walk = walk_about(&[(subject, "3k00000100")], "pull", |_cursor| { + asked += 1; + async move { + // Page one ends at 3k00000150, still above the floor; page + // two ends below it. + let from = if asked == 1 { 200 } else { 140 }; + Ok((page_of(subject, from, 50), Some("more".to_string()))) + } + }) + .await + .expect("no transport to fail"); + + assert!(walk.complete); + assert_eq!( + walk.pages, 2, + "it kept going past the page that did not reach" + ); + assert_eq!(walk.by_subject[subject].len(), 100); + } + + /// A collection that simply ends is complete too, floor or no floor. + #[tokio::test] + async fn a_collection_that_runs_out_is_complete() { + let subject = "at://x/p/one"; + let walk = walk_about(&[(subject, "3k00000001")], "pull", |_cursor| async move { + Ok((page_of(subject, 200, 3), None)) + }) + .await + .expect("no transport to fail"); + + assert!( + walk.complete, + "a cursorless page is the end of the collection" + ); + assert_eq!(walk.pages, 1); + } + + /// **The one this exists for.** A walk that runs out of pages before + /// reaching the floor is *not* complete, and says so. + /// + /// The records that did arrive are still handed back — they are real — + /// but a caller reading a state off the ones that are *missing* has to + /// know it never got to the bottom. Both callers used to be told nothing + /// here and read the cap as the end of the collection, which on the pull + /// side meant a `merged` record past the cap read as `open`. + #[tokio::test] + async fn a_walk_that_runs_out_of_pages_is_not_complete() { + let subject = "at://x/p/one"; + let walk = walk_about(&[(subject, "3k00000001")], "pull", |_cursor| async move { + // Every page ends far above the floor and always offers another. + Ok((page_of(subject, 900_000, 2), Some("more".to_string()))) + }) + .await + .expect("no transport to fail"); + + assert!(!walk.complete, "the cap is not the end of the collection"); + assert_eq!(walk.pages, MAX_PAGES_ABOUT, "it spent the whole cap"); + assert_eq!(walk.by_subject[subject].len(), MAX_PAGES_ABOUT * 2); + } + + /// A batch's floor is the *oldest* subject asked about, not the newest. + /// + /// Stopping at the newest would end the walk while an older subject's + /// records were still ahead of it — the failure that makes one walk for + /// a page of rows different from thirty walks of one. + #[tokio::test] + async fn a_batch_stops_at_the_oldest_subject_not_the_newest() { + let old = "at://x/p/old"; + let new = "at://x/p/new"; + let mut asked = 0; + let walk = walk_about( + &[(new, "3k00000900"), (old, "3k00000100")], + "pull", + |_cursor| { + asked += 1; + async move { + let subject = if asked == 1 { new } else { old }; + // Page two ends at 3k00000096, just below the older + // subject's own key and therefore below the floor. + let from = if asked == 1 { 950 } else { 105 }; + Ok((page_of(subject, from, 10), Some("more".to_string()))) + } + }, + ) + .await + .expect("no transport to fail"); + + assert!(walk.complete); + assert!( + walk.by_subject.contains_key(old), + "the walk stopped before the older subject's records: {:?}", + walk.by_subject.keys().collect::>() + ); + } + + /// Records naming something nobody asked about are dropped, not filed. + #[tokio::test] + async fn a_record_about_another_subject_is_not_collected() { + let wanted = "at://x/p/one"; + let walk = walk_about(&[(wanted, "3k00000001")], "pull", |_cursor| async move { + Ok((page_of("at://x/p/other", 5, 3), None)) + }) + .await + .expect("no transport to fail"); + + assert!(walk.by_subject.is_empty(), "{:?}", walk.by_subject); + } } diff --git a/src/cmd/issue/read.rs b/src/cmd/issue/read.rs index b056388..e6e1b2f 100644 --- a/src/cmd/issue/read.rs +++ b/src/cmd/issue/read.rs @@ -52,7 +52,6 @@ use crate::term::column::{day, ellipsize, pad_to}; use anyhow::Result; use jacquard::types::string::Datetime; use std::collections::HashMap; -use std::collections::HashSet; use std::str::FromStr; /// Tangled's website, where every `view:` link points. The address and its @@ -368,12 +367,6 @@ impl State { } } -/// A PDS will not return more than this many records in one page, and the -/// walk below stops on its own well before this many pages for any real -/// account. The cap is only here so a PDS that returned a cursor forever -/// could not spin. -const MAX_STATE_PAGES: usize = 50; - /// One `listRecords` row read as the issue it is about and the event it is, /// or `None` for a record that names no issue. /// @@ -427,64 +420,35 @@ async fn state_events_for( did: &str, issues: &[(&str, &str)], ) -> Result>> { - let mut found: HashMap> = HashMap::new(); - let Some(floor) = issues.iter().map(|(_, rkey)| *rkey).min() else { - return Ok(found); - }; - let wanted: HashSet<&str> = issues.iter().map(|(uri, _)| *uri).collect(); let pds = crate::clients::atproto::did::pds_or_fail(did).await?; - let mut cursor: Option = None; - let mut pages = 0; - - for _ in 0..MAX_STATE_PAGES { - let (records, next) = crate::clients::atproto::pds::list_records( - &pds, - did, - ISSUE_STATE_NSID, - cursor.as_deref(), - ) - .await?; - pages += 1; - if records.is_empty() { - break; - } - - for record in &records { - let Some((issue, event)) = state_event(record) else { - continue; - }; - if !wanted.contains(issue.as_str()) { - continue; - } - found.entry(issue).or_default().push(event); - } - - let oldest_on_page = records - .last() - .and_then(|r| r["uri"].as_str()) - .and_then(|u| u.rsplit('/').next()) - .unwrap_or_default(); - if oldest_on_page < floor { - crate::logging::debug::log(format!( - "stopping at rkey {oldest_on_page}, which predates the oldest issue key \ - asked about ({floor})" - )); - break; - } - match next { - Some(c) => cursor = Some(c), - None => break, + let walk = + crate::clients::atproto::pds::records_about(&pds, did, ISSUE_STATE_NSID, "issue", issues) + .await?; + // A short walk is said out loud here rather than returned. A listing + // reads a state to *show* it, and a row that reads `?` because the walk + // ran out of pages is a row this can explain; the write paths, where the + // same absence decides whether a record is appended, refuse instead. + if !walk.complete { + crate::term::say::warning!( + Pds, + "the state walk hit its page cap before reaching the oldest issue asked \ + about, so a state below it reads `?` rather than the state it has" + ); + } + let mut found: HashMap> = HashMap::new(); + for (issue, records) in walk.by_subject { + let events: Vec = records + .iter() + .filter_map(state_event) + .map(|(_, event)| event) + .collect(); + if !events.is_empty() { + found.insert(issue, events); } } - crate::logging::debug::log(format!( - "issue states: {} page(s) of {ISSUE_STATE_NSID} from {did} covered {} issue(s)", - pages, - issues.len() - )); Ok(found) } -/// The same question about one issue, which is what `issue view` asks. async fn state_events(did: &str, issue_uri: &str, issue_rkey: &str) -> Result> { Ok(state_events_for(did, &[(issue_uri, issue_rkey)]) .await? diff --git a/src/cmd/pr/write.rs b/src/cmd/pr/write.rs index 9449182..70a72c3 100644 --- a/src/cmd/pr/write.rs +++ b/src/cmd/pr/write.rs @@ -1449,12 +1449,6 @@ fn state_of(statuses: &[Status]) -> PullState { .unwrap_or(PullState::Open) } -/// A PDS will not return more than this many records in one page, and paging -/// stops on its own well before this many pages for any real account — see -/// [`list_statuses`]. The cap is only here so a PDS that returned a cursor -/// forever could not spin. -const MAX_STATUS_PAGES: usize = 50; - /// Every status record in `did`'s PDS that is about `pull_uri`. /// /// listRecords has no filter, so this pages and picks. What keeps that bounded @@ -1463,58 +1457,47 @@ const MAX_STATUS_PAGES: usize = 50; /// once a page ends below the pull's own key there is nothing left to find. /// The comparison is strict, because Tangled's own backfill migration wrote /// status records under the *pull's* record key, and those must still be seen. -async fn list_statuses(did: &str, pull_uri: &str, pull_rkey: &str) -> Result> { +async fn list_statuses(did: &str, pull_uri: &str, pull_rkey: &str) -> Result { let pds = crate::clients::atproto::did::pds_or_fail(did).await?; - let mut found = Vec::new(); - let mut cursor: Option = None; - - for _ in 0..MAX_STATUS_PAGES { - let (records, next) = crate::clients::atproto::pds::list_records( - &pds, - did, - crate::lexicon::tangled::PULL_STATUS_NSID, - cursor.as_deref(), - ) - .await?; - if records.is_empty() { - break; - } - - for record in &records { - if record["value"]["pull"].as_str() != Some(pull_uri) { - continue; - } - found.push(Status { - uri: record["uri"].as_str().unwrap_or_default().to_string(), - created_at: record["value"]["createdAt"] - .as_str() - .unwrap_or_default() - .to_string(), - status: record["value"]["status"] - .as_str() - .unwrap_or_default() - .to_string(), - }); - } + let walk = crate::clients::atproto::pds::records_about( + &pds, + did, + crate::lexicon::tangled::PULL_STATUS_NSID, + "pull", + &[(pull_uri, pull_rkey)], + ) + .await?; + let found = walk + .by_subject + .get(pull_uri) + .map(|records| { + records + .iter() + .map(|record| Status { + uri: record["uri"].as_str().unwrap_or_default().to_string(), + created_at: record["value"]["createdAt"] + .as_str() + .unwrap_or_default() + .to_string(), + status: record["value"]["status"] + .as_str() + .unwrap_or_default() + .to_string(), + }) + .collect() + }) + .unwrap_or_default(); + Ok(StatusWalk { + statuses: found, + complete: walk.complete, + }) +} - let oldest_on_page = records - .last() - .and_then(|r| r["uri"].as_str()) - .and_then(|u| u.rsplit('/').next()) - .unwrap_or_default() - .to_string(); - if oldest_on_page.as_str() < pull_rkey { - crate::logging::debug::log(format!( - "stopping at rkey {oldest_on_page}, which predates the pull's own key {pull_rkey}" - )); - break; - } - match next { - Some(c) => cursor = Some(c), - None => break, - } - } - Ok(found) +/// One account's status records about a pull, and whether they are all of +/// them. +struct StatusWalk { + statuses: Vec, + complete: bool, } /// Whether an account may write a status record about this pull that Tangled @@ -1909,9 +1892,33 @@ async fn set_state(args: StateArgs, wanted: PullState) -> Result<()> { } let uri = target.uri(); - let mut statuses = list_statuses(&target.author, &uri, &target.rkey).await?; + let walk = list_statuses(&target.author, &uri, &target.rkey).await?; + let mut statuses = walk.statuses; + let mut complete = walk.complete; if acting != target.author { - statuses.extend(list_statuses(&acting, &uri, &target.rkey).await?); + let mine = list_statuses(&acting, &uri, &target.rkey).await?; + statuses.extend(mine.statuses); + complete &= mine.complete; + } + // **A state read off a short walk is not a state.** Every decision below + // turns on the *absence* of a record — "already closed, nothing written" + // and the refusal that keeps a merge from being overwritten are both read + // off what did not come back — and a walk that ran out of pages before + // reaching this pull produces exactly the same absence as a pull nobody + // has acted on. Acting on it would mean writing `closed` over a `merged` + // that is sitting one page past the cap, which is the one write this + // command refuses outright when it can see it. + if !complete { + bail!( + "the status walk hit its page cap before reaching {}, so this pull's current \ + state cannot be settled\n\ + refusing to {verb} it: a state record beyond the cap — a merge, most of all — \ + is invisible from here, and writing over a merge is the one thing this \ + command refuses outright when it can see it\n\ + {verb} it through tangled.org, which reads its own index rather than walking \ + the records", + target.rkey, + ); } let current = state_of(&statuses); crate::logging::debug::log(format!( -- 2.51.2