From c3cdefa5ea7cc6428def014ab87cadc0c4cdb99d Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 08:34:02 -0400 Subject: [PATCH] refactor(events): simplify backward pagination internals - Extract a shared fetch_run_events_page helper so the three client paging loops (full list, until, tail) no longer repeat the request/ convert/has_more skeleton; fold the tail loop's two descending-order checks into one and drop its redundant had_events flag. - Skip the latest-seq lookup in list_events_before_with_limit when the caller supplies a before_seq cursor, so a cold projection cache costs at most one full history scan per pagination session instead of one per page. - Remove the dead before_seq max(1) clamp and the passthrough order() accessor from EventListParams. - Document the CLI --tail 0 --follow seeding trick and the reader event_seq placeholder invariant. Co-Authored-By: Claude Fable 5 --- lib/apps/fabro-cli/src/commands/run/events.rs | 4 + .../fabro-server/src/server/handler/events.rs | 15 +- .../fabro-store/src/slate/run_store.rs | 55 +++---- lib/foundation/fabro-client/src/client.rs | 142 +++++++++--------- 4 files changed, 107 insertions(+), 109 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/events.rs b/lib/apps/fabro-cli/src/commands/run/events.rs index 4d23aa855..247a22c6f 100644 --- a/lib/apps/fabro-cli/src/commands/run/events.rs +++ b/lib/apps/fabro-cli/src/commands/run/events.rs @@ -44,6 +44,10 @@ pub(crate) async fn run( let events = match (args.tail, since_cutoff.is_none()) { (Some(tail), true) => { + // With --tail 0 --follow, fetch one event anyway so `last_seq` + // seeds the follow cursor at the true latest event instead of + // replaying the whole history; `apply_filters` drops it from + // the printed output. let tail = if args.follow { tail.max(1) } else { tail }; client.list_run_events_tail(&run_id, tail).await } diff --git a/lib/apps/fabro-server/src/server/handler/events.rs b/lib/apps/fabro-server/src/server/handler/events.rs index dcdc9e2e5..998dee98b 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -59,14 +59,6 @@ impl EventListParams { self.limit.unwrap_or(100).clamp(1, 1000) } - fn before_seq(&self) -> Option { - self.before_seq.map(|seq| seq.max(1)) - } - - fn order(&self) -> EventSequenceOrder { - self.order - } - fn cursor_error(&self) -> Option<&'static str> { match self.order { EventSequenceOrder::Asc if self.before_seq.is_some() => { @@ -238,19 +230,18 @@ async fn list_run_events( return ApiError::bad_request(detail).into_response(); } - let since_seq = params.since_seq(); let limit = params.limit(); match state.stores.runs.open_run_reader(&id).await { Ok(run_store) => { - let events = match params.order() { + let events = match params.order { EventSequenceOrder::Asc => { run_store - .list_events_from_with_limit(since_seq, limit) + .list_events_from_with_limit(params.since_seq(), limit) .await } EventSequenceOrder::Desc => { run_store - .list_events_before_with_limit(params.before_seq(), limit) + .list_events_before_with_limit(params.before_seq, limit) .await } }; diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index 06b0a8e7f..782663706 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -85,8 +85,10 @@ impl RunDatabase { run_summary_store: Arc>>, ) -> Result { let event_seq = if read_only { - // Readers never append, so they do not need to scan the full event - // history to recover the next write sequence. + // Placeholder: readers never append (append_event_envelope + // rejects read-only handles, and reader inners are never cached + // in active_runs), so skip the full event history scan that + // recovers the next write sequence. 1 } else { recover_next_seq(&db, keys::run_events_prefix(&run_id), keys::parse_event_seq).await? @@ -387,31 +389,11 @@ impl RunDatabase { before_seq: Option, limit: usize, ) -> Result> { - let latest_seq = match self - .inner - .shared_projection_cache - .last_seq(&self.inner.run_id) - .await - { - Some(seq) => seq, - None => recover_next_seq( - &self.inner.db, - keys::run_events_prefix(&self.inner.run_id), - keys::parse_event_seq, - ) - .await? - .saturating_sub(1), + let end_seq = match before_seq { + Some(seq) => u64::from(seq), + None => u64::from(self.latest_event_seq().await?) + 1, }; - if latest_seq == 0 { - return Ok(Vec::new()); - } - - let latest_exclusive = u64::from(latest_seq) + 1; - let end_seq = before_seq - .map_or(latest_exclusive, u64::from) - .min(latest_exclusive) - .max(1); - if end_seq == 1 { + if end_seq <= 1 { return Ok(Vec::new()); } @@ -431,6 +413,27 @@ impl RunDatabase { Ok(events) } + /// Latest appended event sequence, or 0 when the run has no events. + /// Served from the projection cache when warm; otherwise recovered by + /// scanning the event history. + async fn latest_event_seq(&self) -> Result { + match self + .inner + .shared_projection_cache + .last_seq(&self.inner.run_id) + .await + { + Some(seq) => Ok(seq), + None => Ok(recover_next_seq( + &self.inner.db, + keys::run_events_prefix(&self.inner.run_id), + keys::parse_event_seq, + ) + .await? + .saturating_sub(1)), + } + } + pub async fn get_event(&self, seq: u32) -> Result> { get_event(&self.inner.db, &self.inner.run_id, seq).await } diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index 53af127b3..531b9f2ba 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -1534,28 +1534,14 @@ impl Client { let mut all_events = Vec::new(); loop { - let response = self - .send_api(|client| async move { - let mut request = client.list_run_events().id(run_id.to_string()); - if let Some(seq) = next_since_seq.and_then(non_zero_u64_from_u32) { - request = request.since_seq(seq); - } - if let Some(limit) = limit.and_then(non_zero_u64_from_usize) { - request = request.limit(limit); - } - request.send().await - }) - .await?; - let parsed = response.into_inner(); - let page_events = parsed - .data - .into_iter() - .map(convert_type::<_, EventEnvelope>) - .collect::>>()?; + let page = EventPageCursor::Ascending { + since_seq: next_since_seq, + }; + let (page_events, has_more) = self.fetch_run_events_page(run_id, page, limit).await?; let next_page_since_seq = page_events.last().map(|event| event.seq.saturating_add(1)); all_events.extend(page_events); - if limit.is_some() || !parsed.meta.has_more || next_page_since_seq.is_none() { + if limit.is_some() || !has_more || next_page_since_seq.is_none() { break; } next_since_seq = next_page_since_seq; @@ -1579,52 +1565,34 @@ impl Client { let fetch_target = max_events.max(2); let mut before_seq = None; let mut descending_events: Vec = Vec::new(); - while descending_events.len() < fetch_target { + loop { let remaining = fetch_target - descending_events.len(); - let response = self - .send_api(|client| async move { - let mut request = client - .list_run_events() - .id(run_id.to_string()) - .order(types::ListRunEventsOrder::Desc) - .limit(remaining.min(1000) as u64); - if let Some(seq) = before_seq.and_then(non_zero_u64_from_u32) { - request = request.before_seq(seq); - } - request.send().await - }) + let (page_events, has_more) = self + .fetch_run_events_page( + run_id, + EventPageCursor::Descending { before_seq }, + Some(remaining), + ) .await?; - let parsed = response.into_inner(); - let page_events = parsed - .data - .into_iter() - .map(convert_type::<_, EventEnvelope>) - .collect::>>()?; - let page_is_descending = page_events - .windows(2) - .all(|events| events[0].seq > events[1].seq); - let continues_descending = descending_events + let keeps_descending = descending_events .last() - .zip(page_events.first()) - .is_none_or(|(previous, next)| previous.seq > next.seq); - if !page_is_descending || !continues_descending { + .into_iter() + .chain(&page_events) + .is_sorted_by(|previous, next| previous.seq > next.seq); + if !keeps_descending { + // An older server ignored the order parameter and returned + // ascending history; fetch everything and slice the tail. let mut events = self.list_run_events(run_id, None, None).await?; let tail_start = events.len().saturating_sub(max_events); return Ok(events.split_off(tail_start)); } - let next_before_seq = page_events.last().map(|event| event.seq); - let had_events = !page_events.is_empty(); + before_seq = page_events.last().map(|event| event.seq); descending_events.extend(page_events); - if descending_events.len() >= fetch_target - || !parsed.meta.has_more - || !had_events - || next_before_seq.is_none() - { + if descending_events.len() >= fetch_target || !has_more || before_seq.is_none() { break; } - before_seq = next_before_seq; } descending_events.reverse(); @@ -1646,28 +1614,16 @@ impl Client { let mut all_events = Vec::new(); while all_events.len() < max_events { let remaining = max_events - all_events.len(); - let response = self - .send_api(|client| async move { - let mut request = client - .list_run_events() - .id(run_id.to_string()) - .limit(remaining.min(1000) as u64); - if let Some(seq) = next_since_seq.and_then(non_zero_u64_from_u32) { - request = request.since_seq(seq); - } - request.send().await - }) + let page = EventPageCursor::Ascending { + since_seq: next_since_seq, + }; + let (page_events, has_more) = self + .fetch_run_events_page(run_id, page, Some(remaining)) .await?; - let parsed = response.into_inner(); - let page_events = parsed - .data - .into_iter() - .map(convert_type::<_, EventEnvelope>) - .collect::>>()?; let next_page_since_seq = page_events.last().map(|event| event.seq.saturating_add(1)); all_events.extend(page_events); - if !parsed.meta.has_more || next_page_since_seq.is_none() { + if !has_more || next_page_since_seq.is_none() { break; } next_since_seq = next_page_since_seq; @@ -1676,6 +1632,44 @@ impl Client { Ok(all_events) } + async fn fetch_run_events_page( + &self, + run_id: &RunId, + cursor: EventPageCursor, + limit: Option, + ) -> Result<(Vec, bool)> { + let response = self + .send_api(|client| async move { + let mut request = client.list_run_events().id(run_id.to_string()); + match cursor { + EventPageCursor::Ascending { since_seq } => { + if let Some(seq) = since_seq.and_then(non_zero_u64_from_u32) { + request = request.since_seq(seq); + } + } + EventPageCursor::Descending { before_seq } => { + request = request.order(types::ListRunEventsOrder::Desc); + if let Some(seq) = before_seq.and_then(non_zero_u64_from_u32) { + request = request.before_seq(seq); + } + } + } + let page_limit = limit.map(|limit| limit.min(1000)); + if let Some(limit) = page_limit.and_then(non_zero_u64_from_usize) { + request = request.limit(limit); + } + request.send().await + }) + .await?; + let parsed = response.into_inner(); + let events = parsed + .data + .into_iter() + .map(convert_type::<_, EventEnvelope>) + .collect::>>()?; + Ok((events, parsed.meta.has_more)) + } + pub async fn attach_run_events( &self, run_id: &RunId, @@ -2153,6 +2147,12 @@ pub fn apply_bearer_token_auth( Ok(builder.default_headers(headers)) } +#[derive(Clone, Copy)] +enum EventPageCursor { + Ascending { since_seq: Option }, + Descending { before_seq: Option }, +} + fn non_zero_u64_from_u32(value: u32) -> Option { NonZeroU64::new(u64::from(value)) }