From 1e1cfaa81bcc41578329c200e0e3fd2e7efb2e6f Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Tue, 18 Aug 2026 21:04:18 -0400 Subject: [PATCH] perf(stack): read a stack's members back together, not one at a time `resubmit` and `merge` each read every member with one `.await` per loop iteration, and every read is a `getRecord` plus a patch blob download, so a ten-member stack was ten round trips before either could plan. Both go through `read_members`, `buffered` and so still in chain order. Co-Authored-By: Claude Opus 5 (1M context) --- TODO.md | 13 +++++- src/cmd/stack/write.rs | 89 +++++++++++++++++++++++++++++++++++------- 2 files changed, 87 insertions(+), 15 deletions(-) diff --git a/TODO.md b/TODO.md index 721e889..15c6b5e 100644 --- a/TODO.md +++ b/TODO.md @@ -1218,7 +1218,18 @@ warn where the stack command is better. instant the running total passes `max_bytes` — a bound enforced against bytes actually received rather than a header the sender controls -- [ ] Concurrent member reads are what's left of that review batch — the +- [x] Concurrent member reads, the last of that review batch. `resubmit` and + `merge` both read every member with one `.await` per loop iteration, + and each read is a `getRecord` plus a blob download — the latest + round's patch, where the change-id lives and whose bytes decide whether + anything changed — so a ten-member stack was ten round trips end to end + before either command could plan anything. Both go through + `read_members` now, `buffered` at `MEMBER_READ_CONCURRENCY`. Ordered + rather than unordered, because both callers build a `Vec` indexed by + position in the chain and the chain's order is the whole subject; + `reordering_the_branch_relinks_the_chain_onto_the_new_order` is what + holds that, its three members all reading at once at this concurrency +- [x] The rest of that batch had already landed — the rest landed: the listing-chain-states preamble is one loader now (`complete_rows`/`own_chain` in `stack/mod.rs`), the blob-fetch path is shared (`pds::blob_bounded`, used by `review` and the stack), diff --git a/src/cmd/stack/write.rs b/src/cmd/stack/write.rs index 3e9809b..fe71151 100644 --- a/src/cmd/stack/write.rs +++ b/src/cmd/stack/write.rs @@ -1265,6 +1265,52 @@ async fn old_member( }) } +/// How many members are read back at once. +/// +/// Each read is a `getRecord` and a blob download — the latest round's +/// patch, which is where a member's change-id lives and whose bytes decide +/// whether anything changed. `resubmit` and `merge` both did one per loop +/// iteration with an `.await` on it, so a ten-member stack was ten round +/// trips end to end before either could plan anything. +/// +/// Bounded, for `images.rs`'s reason: a stack should not open as many +/// simultaneous streams against one PDS as it happens to have members. +/// *Ordered* — `buffered`, not `buffer_unordered` — because both callers +/// build a `Vec` whose index is the member's position in the chain, and the +/// chain's order is the whole subject of these commands. +const MEMBER_READ_CONCURRENCY: usize = 4; + +/// Read every member of `members` back, in order. +/// +/// `state_of` is a pure lookup against the listing already in hand, so it is +/// evaluated per member here rather than being another thing to thread +/// through. +/// The order this returns is the chain's, and +/// `reordering_the_branch_relinks_the_chain_onto_the_new_order` in +/// `tests/stack_flows.rs` is what holds it: its three members all read at +/// once at this concurrency, and it asserts the exact `dependentOn` chain +/// that comes out. +async fn read_members<'a>( + pds: &str, + did: &str, + members: impl Iterator, +) -> Result> { + use futures_util::stream::{self, StreamExt, TryStreamExt}; + let members: Vec<_> = members.collect(); + crate::logging::debug::log(format!( + "stack: reading {} member(s) back, {MEMBER_READ_CONCURRENCY} at a time", + members.len() + )); + stream::iter( + members + .into_iter() + .map(|(member, state)| async move { old_member(pds, did, member, &state).await }), + ) + .buffered(MEMBER_READ_CONCURRENCY) + .try_collect() + .await +} + /// Reconcile the stack's records with the rewritten branch. /// /// Tangled's own semantics, so a stack is interchangeable between atgc and @@ -1342,11 +1388,15 @@ async fn resubmit_inner( // settled: `repo_rows` already resolves the fresh-stack case — no // status record anywhere, but the acting account owns the repo and // authored the pull, so its own PDS was a complete answer — to open. - let mut old: Vec = Vec::new(); - for member in &chain.members { - let uri = member["uri"].as_str().unwrap_or_default(); - old.push(old_member(&pds, &me, member, &super::state_of(&rows, uri)).await?); - } + let old: Vec = read_members( + &pds, + &me, + chain.members.iter().map(|member| { + let uri = member["uri"].as_str().unwrap_or_default(); + (*member, super::state_of(&rows, uri)) + }), + ) + .await?; // The trap this pair of commands is most dangerous in, asked before the // rewrite rather than after it. A commit with no change-id gets @@ -2122,15 +2172,26 @@ pub(in crate::cmd) async fn merge(args: MergeArgs) -> Result<()> { bail!("{title} ({uri}) is {state}; a stack merges bottom-up through open pulls only"); } }; - let mut patches: Vec = Vec::new(); - let mut landing: Vec = Vec::new(); - for &index in &selected { - let member = &chain.members[index]; - let uri = member["uri"].as_str().unwrap_or_default().to_string(); - let old = old_member(&pds, &me, member, "open").await?; - patches.push(old.latest_patch); - landing.push(uri); - } + // Every state here is "open": `selected` is exactly the members the + // refusal above let through, and that is what it checks. + let read = read_members( + &pds, + &me, + selected + .iter() + .map(|&index| (chain.members[index], "open".to_string())), + ) + .await?; + let landing: Vec = selected + .iter() + .map(|&index| { + chain.members[index]["uri"] + .as_str() + .unwrap_or_default() + .to_string() + }) + .collect(); + let patches: Vec = read.into_iter().map(|old| old.latest_patch).collect(); if landing.is_empty() { bail!("everything at or below position {through} is already merged; nothing to do"); } -- 2.51.2