From bb1afae3637aee7462b7f7b3dbd25b5f63259829 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 07:23:23 -0400 Subject: [PATCH 1/5] feat(events): add backward cursor pagination --- docs/public/api-reference/fabro-api.yaml | 42 +++- lib/apps/fabro-cli/src/commands/run/events.rs | 12 +- .../fabro-server/src/server/handler/events.rs | 88 ++++++-- lib/apps/fabro-server/src/server/tests.rs | 91 +++++++++ .../fabro-store/src/slate/projection_cache.rs | 9 + .../fabro-store/src/slate/run_store.rs | 139 ++++++++++++- lib/foundation/fabro-client/src/client.rs | 189 ++++++++++++++++++ .../src/api/run-internals-api.ts | 44 +++- 8 files changed, 574 insertions(+), 40 deletions(-) diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 519ee4188..3af6d7e3c 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -2838,11 +2838,17 @@ paths: operationId: listRunEvents tags: [Run Internals] summary: List Run Events - description: Returns a paginated JSON list of stored run events. + description: | + Returns a paginated JSON list of stored run events. Ascending order + uses `since_seq` as an inclusive cursor. Descending order uses + `before_seq` as an exclusive cursor and starts at the newest event + when `before_seq` is omitted. parameters: - $ref: "#/components/parameters/RunId" - $ref: "#/components/parameters/SinceSeq" - $ref: "#/components/parameters/EventLimit" + - $ref: "#/components/parameters/BeforeSeq" + - $ref: "#/components/parameters/EventOrder" responses: "200": description: Paginated list of run events @@ -2850,6 +2856,15 @@ paths: application/json: schema: $ref: "#/components/schemas/PaginatedEventList" + "400": + description: Invalid cursor and order combination + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" "404": description: Run not found headers: @@ -5853,6 +5868,31 @@ components: default: 1 example: 42 + BeforeSeq: + name: before_seq + in: query + required: false + description: | + Exclusive upper event sequence cursor for descending order. Omit on + the first descending request to start from the newest event. + schema: + type: integer + minimum: 1 + example: 42 + + EventOrder: + name: order + in: query + required: false + description: | + Event sequence order. `since_seq` is valid only with `asc`; + `before_seq` is valid only with `desc`. + schema: + type: string + enum: [asc, desc] + default: asc + example: desc + EventLimit: name: limit in: query diff --git a/lib/apps/fabro-cli/src/commands/run/events.rs b/lib/apps/fabro-cli/src/commands/run/events.rs index 40cb62a52..4d23aa855 100644 --- a/lib/apps/fabro-cli/src/commands/run/events.rs +++ b/lib/apps/fabro-cli/src/commands/run/events.rs @@ -42,10 +42,14 @@ pub(crate) async fn run( None => None, }; - let events = client - .list_run_events(&run_id, None, None) - .await - .context("Failed to list server-backed run events")?; + let events = match (args.tail, since_cutoff.is_none()) { + (Some(tail), true) => { + let tail = if args.follow { tail.max(1) } else { tail }; + client.list_run_events_tail(&run_id, tail).await + } + _ => client.list_run_events(&run_id, None, None).await, + } + .context("Failed to list server-backed run events")?; let last_seq = events.last().map_or(0, |event| event.seq); let all_lines = events .iter() diff --git a/lib/apps/fabro-server/src/server/handler/events.rs b/lib/apps/fabro-server/src/server/handler/events.rs index b8508bd5c..dcdc9e2e5 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -30,12 +30,24 @@ pub(super) fn routes() -> Router> { .route("/runs/{id}/attach", get(attach_run_events)) } +#[derive(Clone, Copy, Default, PartialEq, Eq, serde::Deserialize)] +#[serde(rename_all = "snake_case")] +enum EventSequenceOrder { + #[default] + Asc, + Desc, +} + #[derive(serde::Deserialize)] pub(crate) struct EventListParams { #[serde(default)] - since_seq: Option, + since_seq: Option, #[serde(default)] - limit: Option, + before_seq: Option, + #[serde(default)] + order: EventSequenceOrder, + #[serde(default)] + limit: Option, } impl EventListParams { @@ -46,6 +58,26 @@ impl EventListParams { pub(crate) fn limit(&self) -> usize { 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() => { + Some("before_seq requires order=desc.") + } + EventSequenceOrder::Desc if self.since_seq.is_some() => { + Some("since_seq cannot be combined with order=desc; use before_seq instead.") + } + _ => None, + } + } } #[derive(serde::Deserialize)] @@ -202,29 +234,43 @@ async fn list_run_events( State(state): State>, Query(params): Query, ) -> Response { + if let Some(detail) = params.cursor_error() { + 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) => match run_store - .list_events_from_with_limit(since_seq, limit) - .await - { - Ok(mut events) => { - let has_more = events.len() > limit; - events.truncate(limit); - Json(PaginatedEventList { - data: events, - meta: PaginationMeta { - has_more, - total: None, - }, - }) - .into_response() + Ok(run_store) => { + let events = match params.order() { + EventSequenceOrder::Asc => { + run_store + .list_events_from_with_limit(since_seq, limit) + .await + } + EventSequenceOrder::Desc => { + run_store + .list_events_before_with_limit(params.before_seq(), limit) + .await + } + }; + match events { + Ok(mut events) => { + let has_more = events.len() > limit; + events.truncate(limit); + Json(PaginatedEventList { + data: events, + meta: PaginationMeta { + has_more, + total: None, + }, + }) + .into_response() + } + Err(err) => ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) + .into_response(), } - Err(err) => { - ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() - } - }, + } Err(_) => ApiError::not_found("Run not found.").into_response(), } } diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 454455be0..dec926ffc 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -9828,6 +9828,97 @@ async fn list_run_events_returns_paginated_json() { assert!(body["meta"]["has_more"].is_boolean()); } +#[tokio::test] +async fn list_run_events_descends_from_latest_with_exclusive_cursor() { + let state = test_app_state(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = RunId::new(); + create_durable_run_with_events(&state, run_id, &[ + workflow_event::Event::RunRunnable { + source: fabro_types::RunRunnableSource::StartRequested, + actor: None, + }, + workflow_event::Event::RunStarting, + workflow_event::Event::RunRunning, + ]) + .await; + + let response = app + .clone() + .oneshot( + Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/events?order=desc&limit=2"))) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + let body = response_json!(response, StatusCode::OK).await; + let seqs = body["data"] + .as_array() + .unwrap() + .iter() + .map(|event| event["seq"].as_u64().unwrap()) + .collect::>(); + assert_eq!(seqs, vec![4, 3]); + assert_eq!(body["meta"]["has_more"], true); + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri(api(&format!( + "/runs/{run_id}/events?order=desc&before_seq=3&limit=2" + ))) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + let body = response_json!(response, StatusCode::OK).await; + let seqs = body["data"] + .as_array() + .unwrap() + .iter() + .map(|event| event["seq"].as_u64().unwrap()) + .collect::>(); + assert_eq!(seqs, vec![2, 1]); + assert_eq!(body["meta"]["has_more"], false); +} + +#[tokio::test] +async fn list_run_events_rejects_cursor_for_opposite_order() { + let app = crate::test_support::build_test_router(test_app_state()); + let run_id = RunId::new(); + let cases = [ + ( + format!("/runs/{run_id}/events?order=desc&since_seq=2"), + "since_seq cannot be combined with order=desc; use before_seq instead.", + ), + ( + format!("/runs/{run_id}/events?before_seq=2"), + "before_seq requires order=desc.", + ), + ]; + + for (path, expected_detail) in cases { + let response = app + .clone() + .oneshot( + Request::builder() + .method("GET") + .uri(api(&path)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + let body = response_json!(response, StatusCode::BAD_REQUEST).await; + assert_eq!(body["errors"][0]["detail"], expected_detail); + } +} + #[tokio::test] async fn append_run_event_rejects_run_id_mismatch() { let state = test_app_state(); diff --git a/lib/components/fabro-store/src/slate/projection_cache.rs b/lib/components/fabro-store/src/slate/projection_cache.rs index 73f8aaf2e..c5bf4b2e3 100644 --- a/lib/components/fabro-store/src/slate/projection_cache.rs +++ b/lib/components/fabro-store/src/slate/projection_cache.rs @@ -171,6 +171,15 @@ impl RunProjectionCache { .map(|entry| state.with_children_count(entry)) } + pub(crate) async fn last_seq(&self, run_id: &RunId) -> Option { + self.state + .lock() + .await + .entries + .get(run_id) + .map(|entry| entry.last_seq) + } + pub(crate) async fn get_summary(&self, run_id: &RunId, now: DateTime) -> Option { let mut entry = { let state = self.state.lock().await; diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index 83b0a087a..06b0a8e7f 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -379,6 +379,58 @@ impl RunDatabase { list_events_from_with_limit(&self.inner.db, &self.inner.run_id, start_seq, limit).await } + /// Returns up to `limit + 1` events immediately before `before_seq` in + /// descending sequence order. Omitting `before_seq` starts at the newest + /// event. The extra item lets callers compute `has_more`. + pub async fn list_events_before_with_limit( + &self, + 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), + }; + 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 { + return Ok(Vec::new()); + } + + let window_size = u64::try_from(limit.saturating_add(1)).unwrap_or(u64::MAX); + let start_seq = + u32::try_from(end_seq.saturating_sub(window_size).max(1)).unwrap_or(u32::MAX); + let end_seq = u32::try_from(end_seq).ok(); + let mut events = list_events_in_range_with_limit( + &self.inner.db, + &self.inner.run_id, + start_seq, + end_seq, + limit, + ) + .await?; + events.reverse(); + Ok(events) + } + pub async fn get_event(&self, seq: u32) -> Result> { get_event(&self.inner.db, &self.inner.run_id, seq).await } @@ -568,13 +620,35 @@ async fn list_events_from_with_limit( where R: DbRead + Sync, { + list_events_in_range_with_limit(db, run_id, start_seq, None, limit).await +} + +async fn list_events_in_range_with_limit( + db: &R, + run_id: &RunId, + start_seq: u32, + end_seq: Option, + limit: usize, +) -> Result> +where + R: DbRead + Sync, +{ + if end_seq.is_some_and(|end_seq| end_seq <= start_seq) { + return Ok(Vec::new()); + } + let event_prefix = keys::run_events_prefix(run_id); let max_events = limit.saturating_add(1); // Seek to the page cursor and decode only the requested page plus the // sentinel used to compute `has_more`. - let mut iter = db - .scan(keys::run_event_seq_prefix(run_id, start_seq)..) - .await?; + let start_key = keys::run_event_seq_prefix(run_id, start_seq); + let mut iter = match end_seq { + Some(end_seq) => { + let end_key = keys::run_event_seq_prefix(run_id, end_seq); + db.scan(start_key..end_key).await? + } + None => db.scan(start_key..).await?, + }; let mut events = Vec::new(); while events.len() < max_events { let Some(entry) = iter.next().await? else { @@ -916,6 +990,65 @@ mod tests { assert_eq!(seqs, vec![3]); } + #[tokio::test] + async fn list_events_before_with_limit_returns_newest_events_and_sentinel() { + let run = fresh_run().await; + let run_id = run.run_id(); + for idx in 1..=5 { + run.append_event(&stage_prompt_payload(&run_id, idx, Some("alpha"))) + .await + .unwrap(); + } + + let events = run.list_events_before_with_limit(None, 2).await.unwrap(); + + let seqs: Vec = events.iter().map(|event| event.seq).collect(); + assert_eq!(seqs, vec![6, 5, 4]); + } + + #[tokio::test] + async fn list_events_before_with_limit_does_not_read_older_history() { + let run = fresh_run().await; + let run_id = run.run_id(); + for idx in 1..=5 { + run.append_event(&stage_prompt_payload(&run_id, idx, Some("alpha"))) + .await + .unwrap(); + } + run.inner + .db + .put(keys::run_event_key(&run_id, 2, 0), b"invalid json") + .await + .unwrap(); + + let events = run.list_events_before_with_limit(None, 2).await.unwrap(); + + let seqs: Vec = events.iter().map(|event| event.seq).collect(); + assert_eq!(seqs, vec![6, 5, 4]); + } + + #[tokio::test] + async fn list_events_before_with_limit_uses_exclusive_cursor() { + let run = fresh_run().await; + let run_id = run.run_id(); + for idx in 1..=5 { + run.append_event(&stage_prompt_payload(&run_id, idx, Some("alpha"))) + .await + .unwrap(); + } + + let events = run.list_events_before_with_limit(Some(5), 2).await.unwrap(); + + let seqs: Vec = events.iter().map(|event| event.seq).collect(); + assert_eq!(seqs, vec![4, 3, 2]); + assert!( + run.list_events_before_with_limit(Some(1), 2) + .await + .unwrap() + .is_empty() + ); + } + #[tokio::test] async fn list_events_for_stage_returns_only_matching_events_in_seq_order() { let run = fresh_run().await; diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index be21fe2dd..53af127b3 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -1564,6 +1564,74 @@ impl Client { Ok(all_events) } + /// Returns the newest `max_events` in ascending sequence order. + pub async fn list_run_events_tail( + &self, + run_id: &RunId, + max_events: usize, + ) -> Result> { + if max_events == 0 { + return Ok(Vec::new()); + } + + // Fetch two events when the caller asks for one so an older server + // that silently ignores the new order parameter can be detected. + 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 { + 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 + }) + .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 + .last() + .zip(page_events.first()) + .is_none_or(|(previous, next)| previous.seq > next.seq); + if !page_is_descending || !continues_descending { + 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(); + descending_events.extend(page_events); + if descending_events.len() >= fetch_target + || !parsed.meta.has_more + || !had_events + || next_before_seq.is_none() + { + break; + } + before_seq = next_before_seq; + } + + descending_events.reverse(); + let tail_start = descending_events.len().saturating_sub(max_events); + Ok(descending_events.split_off(tail_start)) + } + pub async fn list_run_events_until( &self, run_id: &RunId, @@ -2148,6 +2216,17 @@ mod tests { } } + fn run_event_json(run_id: &RunId, seq: u32) -> serde_json::Value { + json!({ + "seq": seq, + "event": "run.running", + "id": format!("evt-{seq}"), + "run_id": run_id, + "ts": "2026-07-24T12:00:00Z", + "properties": {}, + }) + } + #[cfg(unix)] #[tokio::test] async fn refresh_access_token_allows_plain_http_targets() { @@ -2280,6 +2359,116 @@ mod tests { assert!(models.is_empty()); } + #[tokio::test] + async fn list_run_events_tail_pages_backward_and_returns_ascending() { + let server = MockServer::start_async().await; + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let newest_page = server + .mock_async(|when, then| { + when.method(GET) + .path(format!("/api/v1/runs/{run_id}/events")) + .query_param("order", "desc") + .query_param("limit", "5"); + then.status(200) + .header("Content-Type", "application/json") + .json_body(json!({ + "data": [ + run_event_json(&run_id, 6), + run_event_json(&run_id, 5), + run_event_json(&run_id, 4), + ], + "meta": { "has_more": true }, + })); + }) + .await; + let older_page = server + .mock_async(|when, then| { + when.method(GET) + .path(format!("/api/v1/runs/{run_id}/events")) + .query_param("order", "desc") + .query_param("before_seq", "4") + .query_param("limit", "2"); + then.status(200) + .header("Content-Type", "application/json") + .json_body(json!({ + "data": [ + run_event_json(&run_id, 3), + run_event_json(&run_id, 2), + ], + "meta": { "has_more": true }, + })); + }) + .await; + + let client = Client::new_no_proxy(&server.url("")).unwrap(); + let events = client.list_run_events_tail(&run_id, 5).await.unwrap(); + + newest_page.assert_async().await; + older_page.assert_async().await; + let seqs = events + .into_iter() + .map(|event| event.seq) + .collect::>(); + assert_eq!(seqs, vec![2, 3, 4, 5, 6]); + } + + #[tokio::test] + async fn list_run_events_tail_falls_back_when_server_ignores_descending_order() { + let server = MockServer::start_async().await; + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let unsupported_descending_page = server + .mock_async(|when, then| { + when.method(GET) + .path(format!("/api/v1/runs/{run_id}/events")) + .query_param("order", "desc") + .query_param("limit", "3"); + then.status(200) + .header("Content-Type", "application/json") + .json_body(json!({ + "data": [ + run_event_json(&run_id, 1), + run_event_json(&run_id, 2), + run_event_json(&run_id, 3), + ], + "meta": { "has_more": true }, + })); + }) + .await; + let full_history = server + .mock_async(|when, then| { + when.method(GET) + .path(format!("/api/v1/runs/{run_id}/events")) + .query_param_missing("order") + .query_param_missing("before_seq") + .query_param_missing("since_seq") + .query_param_missing("limit"); + then.status(200) + .header("Content-Type", "application/json") + .json_body(json!({ + "data": [ + run_event_json(&run_id, 1), + run_event_json(&run_id, 2), + run_event_json(&run_id, 3), + run_event_json(&run_id, 4), + run_event_json(&run_id, 5), + ], + "meta": { "has_more": false }, + })); + }) + .await; + + let client = Client::new_no_proxy(&server.url("")).unwrap(); + let events = client.list_run_events_tail(&run_id, 3).await.unwrap(); + + unsupported_descending_page.assert_async().await; + full_history.assert_async().await; + let seqs = events + .into_iter() + .map(|event| event.seq) + .collect::>(); + assert_eq!(seqs, vec![3, 4, 5]); + } + #[tokio::test] async fn test_provider_credentials_posts_api_key() { let server = MockServer::start_async().await; diff --git a/lib/packages/fabro-api-client/src/api/run-internals-api.ts b/lib/packages/fabro-api-client/src/api/run-internals-api.ts index f10deda30..85dfe6632 100644 --- a/lib/packages/fabro-api-client/src/api/run-internals-api.ts +++ b/lib/packages/fabro-api-client/src/api/run-internals-api.ts @@ -472,15 +472,17 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config }; }, /** - * Returns a paginated JSON list of stored run events. + * Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. * @summary List Run Events * @param {string} id Unique run identifier (ULID). * @param {number} [sinceSeq] First event sequence number to include. * @param {number} [limit] Maximum number of events to return. + * @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event. + * @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - listRunEvents: async (id: string, sinceSeq?: number, limit?: number, options: RawAxiosRequestConfig = {}): Promise => { + listRunEvents: async (id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options: RawAxiosRequestConfig = {}): Promise => { // verify required parameter 'id' is not null or undefined assertParamExists('listRunEvents', 'id', id) const localVarPath = `/api/v1/runs/{id}/events` @@ -510,6 +512,14 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config localVarQueryParameter['limit'] = limit; } + if (beforeSeq !== undefined) { + localVarQueryParameter['before_seq'] = beforeSeq; + } + + if (order !== undefined) { + localVarQueryParameter['order'] = order; + } + localVarHeaderParameter['Accept'] = 'application/json'; setSearchParams(localVarUrlObj, localVarQueryParameter); @@ -1037,16 +1047,18 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, /** - * Returns a paginated JSON list of stored run events. + * Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. * @summary List Run Events * @param {string} id Unique run identifier (ULID). * @param {number} [sinceSeq] First event sequence number to include. * @param {number} [limit] Maximum number of events to return. + * @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event. + * @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - async listRunEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { - const localVarAxiosArgs = await localVarAxiosParamCreator.listRunEvents(id, sinceSeq, limit, options); + async listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.listRunEvents(id, sinceSeq, limit, beforeSeq, order, options); const localVarOperationServerIndex = configuration?.serverIndex ?? 0; const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.listRunEvents']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); @@ -1278,16 +1290,18 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b return localVarFp.listRunArtifacts(id, options).then((request) => request(axios, basePath)); }, /** - * Returns a paginated JSON list of stored run events. + * Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. * @summary List Run Events * @param {string} id Unique run identifier (ULID). * @param {number} [sinceSeq] First event sequence number to include. * @param {number} [limit] Maximum number of events to return. + * @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event. + * @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - listRunEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig): AxiosPromise { - return localVarFp.listRunEvents(id, sinceSeq, limit, options).then((request) => request(axios, basePath)); + listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.listRunEvents(id, sinceSeq, limit, beforeSeq, order, options).then((request) => request(axios, basePath)); }, /** * Returns the ordered list of stages in a run\'s workflow graph with their current status and timing. Stages are bounded by the workflow graph size, typically fewer than 20. @@ -1499,16 +1513,18 @@ export class RunInternalsApi extends BaseAPI { } /** - * Returns a paginated JSON list of stored run events. + * Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. * @summary List Run Events * @param {string} id Unique run identifier (ULID). * @param {number} [sinceSeq] First event sequence number to include. * @param {number} [limit] Maximum number of events to return. + * @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event. + * @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - public listRunEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig) { - return RunInternalsApiFp(this.configuration).listRunEvents(id, sinceSeq, limit, options).then((request) => request(this.axios, this.basePath)); + public listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).listRunEvents(id, sinceSeq, limit, beforeSeq, order, options).then((request) => request(this.axios, this.basePath)); } /** @@ -1611,3 +1627,9 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).writeRunBlob(id, body, options).then((request) => request(this.axios, this.basePath)); } } + +export const ListRunEventsOrderEnum = { + ASC: 'asc', + DESC: 'desc' +} as const; +export type ListRunEventsOrderEnum = typeof ListRunEventsOrderEnum[keyof typeof ListRunEventsOrderEnum]; From c3cdefa5ea7cc6428def014ab87cadc0c4cdb99d Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 08:34:02 -0400 Subject: [PATCH 2/5] 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)) } From b886f826224aae5d058ebf4acca710f74d9601d7 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 09:50:40 -0400 Subject: [PATCH 3/5] Clamp backward pagination end bound to the event key-order limit Event keys zero-pad seq to six digits, so an exclusive end bound past MAX_EVENT_SEQ formatted as a seven-digit prefix that sorts before real event keys, producing an inverted scan range. This made the newest page come back empty once a run reached MAX_EVENT_SEQ, and let a client supplied before_seq beyond MAX_EVENT_SEQ garble the range. Clamp the bound and treat anything past MAX_EVENT_SEQ as unbounded; no stored sequence exceeds it, so the results are equivalent. Co-Authored-By: Claude Fable 5 --- .../fabro-store/src/slate/run_store.rs | 58 ++++++++++++++++++- 1 file changed, 56 insertions(+), 2 deletions(-) diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index bf700f96c..36b28ef0b 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -389,10 +389,15 @@ impl RunDatabase { before_seq: Option, limit: usize, ) -> Result> { + // Event keys zero-pad seq to six digits (see `keys::run_event_key`), + // so an end bound past `MAX_EVENT_SEQ` would format as a seven-digit + // prefix that breaks lexicographic key order. No stored seq exceeds + // `MAX_EVENT_SEQ`, so clamp and treat that bound as unbounded. let end_seq = match before_seq { Some(seq) => u64::from(seq), None => u64::from(self.latest_event_seq().await?) + 1, - }; + } + .min(u64::from(keys::MAX_EVENT_SEQ) + 1); if end_seq <= 1 { return Ok(Vec::new()); } @@ -400,7 +405,9 @@ impl RunDatabase { let window_size = u64::try_from(limit.saturating_add(1)).unwrap_or(u64::MAX); let start_seq = u32::try_from(end_seq.saturating_sub(window_size).max(1)).unwrap_or(u32::MAX); - let end_seq = u32::try_from(end_seq).ok(); + let end_seq = u32::try_from(end_seq) + .ok() + .filter(|end| *end <= keys::MAX_EVENT_SEQ); let mut events = list_events_in_range_with_limit( &self.inner.db, &self.inner.run_id, @@ -1038,6 +1045,53 @@ mod tests { ); } + #[tokio::test] + async fn list_events_before_with_limit_reads_newest_page_at_max_event_seq() { + let run = fresh_run().await; + let run_id = run.run_id(); + run.inner + .event_seq + .as_ref() + .unwrap() + .store(keys::MAX_EVENT_SEQ - 1, Ordering::SeqCst); + run.append_event(&stage_prompt_payload(&run_id, 1, Some("alpha"))) + .await + .unwrap(); + run.append_event(&stage_prompt_payload(&run_id, 2, Some("beta"))) + .await + .unwrap(); + + let events = run.list_events_before_with_limit(None, 2).await.unwrap(); + + let seqs: Vec = events.iter().map(|event| event.seq).collect(); + assert_eq!(seqs, vec![keys::MAX_EVENT_SEQ, keys::MAX_EVENT_SEQ - 1]); + } + + #[tokio::test] + async fn list_events_before_with_limit_clamps_cursor_beyond_max_event_seq() { + let run = fresh_run().await; + let run_id = run.run_id(); + run.inner + .event_seq + .as_ref() + .unwrap() + .store(keys::MAX_EVENT_SEQ - 1, Ordering::SeqCst); + run.append_event(&stage_prompt_payload(&run_id, 1, Some("alpha"))) + .await + .unwrap(); + run.append_event(&stage_prompt_payload(&run_id, 2, Some("beta"))) + .await + .unwrap(); + + let events = run + .list_events_before_with_limit(Some(u32::MAX), 2) + .await + .unwrap(); + + let seqs: Vec = events.iter().map(|event| event.seq).collect(); + assert_eq!(seqs, vec![keys::MAX_EVENT_SEQ, keys::MAX_EVENT_SEQ - 1]); + } + #[tokio::test] async fn append_event_rejects_sequences_beyond_key_order_limit() { let run = fresh_run().await; From 7e7d7e445745eae3af84c7a42d8214336b53562c Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 10:05:01 -0400 Subject: [PATCH 4/5] Recover latest event seq with bounded probes instead of a full scan On a projection-cache miss, descending pagination recovered the latest sequence by scanning the run's entire event prefix, making a cold-cache order=desc request O(total_events). Binary-search the zero-padded sequence key space with single-entry probes instead, bounding recovery to O(log MAX_EVENT_SEQ) reads. The probe predicate (smallest stored sequence at or above a bound) stays monotone across gaps left by failed appends. Co-Authored-By: Claude Fable 5 --- .../fabro-store/src/slate/run_store.rs | 116 +++++++++++++++++- 1 file changed, 111 insertions(+), 5 deletions(-) diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index 0140d0f29..f6a701b79 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -445,8 +445,8 @@ impl RunDatabase { } /// 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. + /// Served from the projection cache when warm; otherwise recovered with + /// bounded probes of the event key space rather than a full history scan. async fn latest_event_seq(&self) -> Result { match self .inner @@ -455,9 +455,7 @@ impl RunDatabase { .await { Some((_, seq)) => Ok(seq), - None => Ok(recover_next_seq(&self.inner.db, &self.inner.run_id) - .await? - .saturating_sub(1)), + None => recover_latest_seq(&self.inner.db, &self.inner.run_id).await, } } @@ -625,6 +623,39 @@ where Ok(max_seq.saturating_add(1).max(1)) } +/// Smallest stored event sequence at or above `seq`, if any. +async fn first_event_seq_at_or_after(db: &R, run_id: &RunId, seq: u32) -> Result> +where + R: DbRead + Sync, +{ + let mut scan = EventScan::seek(db, run_id, seq).await?; + Ok(scan.next().await?.map(|(seq, _)| seq)) +} + +/// Largest stored event sequence for the run, or 0 when the run has no +/// events. Binary-searches the sequence space with single-entry probes so +/// recovery reads O(log `MAX_EVENT_SEQ`) entries instead of the full event +/// history. Probing for the smallest sequence at or above a bound is +/// monotone even when failed appends leave gaps in the sequence. +async fn recover_latest_seq(db: &R, run_id: &RunId) -> Result +where + R: DbRead + Sync, +{ + let Some(mut lo) = first_event_seq_at_or_after(db, run_id, 1).await? else { + return Ok(0); + }; + // Invariant: `lo` is a stored sequence and no stored sequence is >= `hi`. + let mut hi = keys::MAX_EVENT_SEQ + 1; + while lo + 1 < hi { + let mid = lo + (hi - lo) / 2; + match first_event_seq_at_or_after(db, run_id, mid).await? { + Some(seq) => lo = seq, + None => hi = mid, + } + } + Ok(lo) +} + /// Cursor over a run's stored events starting at `start_seq`, yielding raw /// `(seq, payload)` entries in ascending sequence order (event keys embed a /// zero-padded sequence, so key order matches sequence order). @@ -1117,6 +1148,81 @@ mod tests { assert_eq!(seqs, vec![keys::MAX_EVENT_SEQ, keys::MAX_EVENT_SEQ - 1]); } + #[tokio::test] + async fn recover_latest_seq_returns_zero_for_empty_history() { + let object_store = Arc::new(InMemory::new()); + let store = Database::new(object_store, "", Duration::from_millis(1), None); + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let run = store.create_run(&run_id).await.unwrap(); + + let latest = super::recover_latest_seq(&run.inner.db, &run_id) + .await + .unwrap(); + + assert_eq!(latest, 0); + } + + #[tokio::test] + async fn recover_latest_seq_finds_latest_across_sparse_gaps() { + let run = fresh_run().await; + let run_id = run.run_id(); + run.append_event(&stage_prompt_payload(&run_id, 1, Some("alpha"))) + .await + .unwrap(); + run.inner + .db + .put(keys::run_event_key(&run_id, 731_204, 0), b"{}") + .await + .unwrap(); + + let latest = super::recover_latest_seq(&run.inner.db, &run_id) + .await + .unwrap(); + + assert_eq!(latest, 731_204); + } + + #[tokio::test] + async fn recover_latest_seq_reads_max_event_seq() { + let run = fresh_run().await; + let run_id = run.run_id(); + run.inner + .db + .put(keys::run_event_key(&run_id, keys::MAX_EVENT_SEQ, 0), b"{}") + .await + .unwrap(); + + let latest = super::recover_latest_seq(&run.inner.db, &run_id) + .await + .unwrap(); + + assert_eq!(latest, keys::MAX_EVENT_SEQ); + } + + #[tokio::test] + async fn list_events_before_with_limit_serves_newest_page_from_cold_cache() { + let object_store = Arc::new(InMemory::new()); + let store = Database::new(object_store.clone(), "", Duration::from_millis(1), None); + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let run = store.create_run(&run_id).await.unwrap(); + run.append_event(&run_created_payload(&run_id)) + .await + .unwrap(); + for idx in 1..=4 { + run.append_event(&stage_prompt_payload(&run_id, idx, Some("alpha"))) + .await + .unwrap(); + } + + let reopened = Database::new(object_store, "", Duration::from_millis(1), None); + let reader = reopened.open_run_reader(&run_id).await.unwrap(); + + let events = reader.list_events_before_with_limit(None, 2).await.unwrap(); + + let seqs: Vec = events.iter().map(|event| event.seq).collect(); + assert_eq!(seqs, vec![5, 4, 3]); + } + #[tokio::test] async fn append_event_rejects_sequences_beyond_key_order_limit() { let run = fresh_run().await; From b77116994f1da15b70d48fead83b97b48d7de676 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 10:21:25 -0400 Subject: [PATCH 5/5] Harden event pagination bounds and scope desc params to the run route Guard EventScan seeks against sequences past MAX_EVENT_SEQ: a seven-digit start prefix sorts below six-digit event keys, so an unvalidated since_seq like 5000000 returned an incorrect slice of history instead of an empty page. An end bound past MAX_EVENT_SEQ now delegates to the unbounded scan, which is equivalent because no stored sequence exceeds it. Clamp the descending exclusive end to just past the newest stored event, so an oversized before_seq cursor pages from the newest event instead of probing empty key space and returning nothing. Split RunEventListParams out of EventListParams so before_seq and order are only accepted by /runs/{id}/events; the session, stage, pair transcript, and demo endpoints go back to ignoring them instead of accepting order=desc while returning ascending results. Co-Authored-By: Claude Fable 5 --- .../fabro-server/src/server/handler/events.rs | 29 ++++++- .../fabro-store/src/slate/run_store.rs | 83 +++++++++++++++++-- 2 files changed, 99 insertions(+), 13 deletions(-) diff --git a/lib/apps/fabro-server/src/server/handler/events.rs b/lib/apps/fabro-server/src/server/handler/events.rs index 5bca3c18b..3418c2cda 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -40,6 +40,27 @@ enum EventSequenceOrder { #[derive(serde::Deserialize)] pub(crate) struct EventListParams { + #[serde(default)] + since_seq: Option, + #[serde(default)] + limit: Option, +} + +impl EventListParams { + pub(crate) fn since_seq(&self) -> u32 { + self.since_seq.unwrap_or(1).max(1) + } + + pub(crate) fn limit(&self) -> usize { + self.limit.unwrap_or(100).clamp(1, 1000) + } +} + +/// Query parameters for `/runs/{id}/events` only. Descending pagination via +/// `before_seq` + `order` is not part of the shared `EventListParams` +/// contract used by the session, stage, and pair transcript endpoints. +#[derive(serde::Deserialize)] +struct RunEventListParams { #[serde(default)] since_seq: Option, #[serde(default)] @@ -50,12 +71,12 @@ pub(crate) struct EventListParams { limit: Option, } -impl EventListParams { - pub(crate) fn since_seq(&self) -> u32 { +impl RunEventListParams { + fn since_seq(&self) -> u32 { self.since_seq.unwrap_or(1).max(1) } - pub(crate) fn limit(&self) -> usize { + fn limit(&self) -> usize { self.limit.unwrap_or(100).clamp(1, 1000) } @@ -224,7 +245,7 @@ async fn append_run_event( async fn list_run_events( RequireRunManagementTarget(id, _actor): RequireRunManagementTarget, State(state): State>, - Query(params): Query, + Query(params): Query, ) -> Response { if let Some(detail) = params.cursor_error() { return ApiError::bad_request(detail).into_response(); diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index f6a701b79..9058627a6 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -407,20 +407,25 @@ impl RunDatabase { /// Returns up to `limit + 1` events immediately before `before_seq` in /// descending sequence order. Omitting `before_seq` starts at the newest + /// event, and a cursor beyond the newest event pages from the newest /// event. The extra item lets callers compute `has_more`. pub async fn list_events_before_with_limit( &self, before_seq: Option, limit: usize, ) -> Result> { - // Event keys zero-pad seq to six digits (see `keys::run_event_key`), - // so an end bound past `MAX_EVENT_SEQ` would format as a seven-digit - // prefix that breaks lexicographic key order. No stored seq exceeds - // `MAX_EVENT_SEQ`, so clamp and treat that bound as unbounded. + // Clamp the exclusive end to just past the newest stored event so an + // oversized cursor pages from the newest event instead of probing + // empty key space above it, and never past `MAX_EVENT_SEQ + 1`: event + // keys zero-pad seq to six digits (see `keys::run_event_key`), so a + // larger end bound would format as a seven-digit prefix that breaks + // lexicographic key order. + let newest = u64::from(self.latest_event_seq().await?); let end_seq = match before_seq { Some(seq) => u64::from(seq), - None => u64::from(self.latest_event_seq().await?) + 1, + None => u64::MAX, } + .min(newest + 1) .min(u64::from(keys::MAX_EVENT_SEQ) + 1); if end_seq <= 1 { return Ok(Vec::new()); @@ -660,7 +665,11 @@ where /// `(seq, payload)` entries in ascending sequence order (event keys embed a /// zero-padded sequence, so key order matches sequence order). struct EventScan { - iter: DbIterator, + // `None` when the requested start is beyond the storable sequence range: + // event keys zero-pad seq to six digits, so seeking past `MAX_EVENT_SEQ` + // would format a seven-digit prefix that breaks lexicographic order and + // returns an incorrect slice of history instead of an empty one. + iter: Option, } impl EventScan { @@ -668,8 +677,11 @@ impl EventScan { where R: DbRead + Sync, { + if start_seq > keys::MAX_EVENT_SEQ { + return Ok(Self { iter: None }); + } let iter = db.scan(keys::run_events_range(run_id, start_seq)).await?; - Ok(Self { iter }) + Ok(Self { iter: Some(iter) }) } /// Like `seek`, but stops before `end_seq` instead of scanning to the @@ -678,14 +690,25 @@ impl EventScan { where R: DbRead + Sync, { + if end_seq > keys::MAX_EVENT_SEQ { + // No stored sequence exceeds `MAX_EVENT_SEQ`, so a larger end + // bound is equivalent to an unbounded scan. + return Self::seek(db, run_id, start_seq).await; + } + if start_seq >= end_seq { + return Ok(Self { iter: None }); + } let range = keys::run_event_seq_prefix(run_id, start_seq) ..keys::run_event_seq_prefix(run_id, end_seq); let iter = db.scan(range).await?; - Ok(Self { iter }) + Ok(Self { iter: Some(iter) }) } async fn next(&mut self) -> Result> { - while let Some(entry) = self.iter.next().await? { + let Some(iter) = self.iter.as_mut() else { + return Ok(None); + }; + while let Some(entry) = iter.next().await? { let key = key_to_str(&entry.key)?; let Some(seq) = keys::parse_event_seq(key) else { continue; @@ -1148,6 +1171,48 @@ mod tests { assert_eq!(seqs, vec![keys::MAX_EVENT_SEQ, keys::MAX_EVENT_SEQ - 1]); } + #[tokio::test] + async fn list_events_from_with_limit_is_empty_beyond_key_order_limit() { + let run = fresh_run().await; + let run_id = run.run_id(); + run.inner + .event_seq + .as_ref() + .unwrap() + .store(keys::MAX_EVENT_SEQ - 1, Ordering::SeqCst); + run.append_event(&stage_prompt_payload(&run_id, 1, Some("alpha"))) + .await + .unwrap(); + run.append_event(&stage_prompt_payload(&run_id, 2, Some("beta"))) + .await + .unwrap(); + + let events = super::list_events_from_with_limit(&run.inner.db, &run_id, 5_000_000, 10) + .await + .unwrap(); + + assert!(events.is_empty()); + } + + #[tokio::test] + async fn list_events_before_with_limit_pages_from_newest_for_oversized_cursor() { + let run = fresh_run().await; + let run_id = run.run_id(); + for idx in 1..=5 { + run.append_event(&stage_prompt_payload(&run_id, idx, Some("alpha"))) + .await + .unwrap(); + } + + let events = run + .list_events_before_with_limit(Some(500_000), 2) + .await + .unwrap(); + + let seqs: Vec = events.iter().map(|event| event.seq).collect(); + assert_eq!(seqs, vec![6, 5, 4]); + } + #[tokio::test] async fn recover_latest_seq_returns_zero_for_empty_history() { let object_store = Arc::new(InMemory::new());