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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-07-24 08:34:02 -04:00
parent bb1afae363
commit c3cdefa5ea
No known key found for this signature in database
4 changed files with 107 additions and 109 deletions

View file

@ -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
}

View file

@ -59,14 +59,6 @@ impl EventListParams {
self.limit.unwrap_or(100).clamp(1, 1000)
}
fn before_seq(&self) -> Option<u32> {
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
}
};

View file

@ -85,8 +85,10 @@ impl RunDatabase {
run_summary_store: Arc<OnceLock<Arc<RunSummaryStore>>>,
) -> Result<Self> {
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<u32>,
limit: usize,
) -> Result<Vec<EventEnvelope>> {
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<u32> {
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<Option<EventEnvelope>> {
get_event(&self.inner.db, &self.inner.run_id, seq).await
}

View file

@ -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::<Result<Vec<EventEnvelope>>>()?;
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<EventEnvelope> = 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::<Result<Vec<EventEnvelope>>>()?;
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::<Result<Vec<EventEnvelope>>>()?;
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<usize>,
) -> Result<(Vec<EventEnvelope>, 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::<Result<Vec<EventEnvelope>>>()?;
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<u32> },
Descending { before_seq: Option<u32> },
}
fn non_zero_u64_from_u32(value: u32) -> Option<NonZeroU64> {
NonZeroU64::new(u64::from(value))
}