diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 2e0dba95a..81794e101 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..247a22c6f 100644 --- a/lib/apps/fabro-cli/src/commands/run/events.rs +++ b/lib/apps/fabro-cli/src/commands/run/events.rs @@ -42,10 +42,18 @@ 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) => { + // 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 + } + _ => 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 131c3f956..3418c2cda 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -30,6 +30,14 @@ 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)] @@ -48,6 +56,43 @@ impl EventListParams { } } +/// 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)] + before_seq: Option, + #[serde(default)] + order: EventSequenceOrder, + #[serde(default)] + limit: Option, +} + +impl RunEventListParams { + fn since_seq(&self) -> u32 { + self.since_seq.unwrap_or(1).max(1) + } + + fn limit(&self) -> usize { + self.limit.unwrap_or(100).clamp(1, 1000) + } + + 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)] struct AttachParams { #[serde(default)] @@ -200,31 +245,44 @@ async fn append_run_event( async fn list_run_events( RequireRunManagementTarget(id, _actor): RequireRunManagementTarget, State(state): State>, - Query(params): Query, + Query(params): Query, ) -> Response { - let since_seq = params.since_seq(); + if let Some(detail) = params.cursor_error() { + return ApiError::bad_request(detail).into_response(); + } + 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(params.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 df85f5191..1a86ef59f 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -9997,6 +9997,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/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index 8627250f9..9058627a6 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -405,6 +405,65 @@ 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, 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> { + // 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::MAX, + } + .min(newest + 1) + .min(u64::from(keys::MAX_EVENT_SEQ) + 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() + .filter(|end| *end <= keys::MAX_EVENT_SEQ); + 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) + } + + /// Latest appended event sequence, or 0 when the run has no events. + /// 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 + .shared_projection_cache + .projection_snapshot(&self.inner.run_id) + .await + { + Some((_, seq)) => Ok(seq), + None => recover_latest_seq(&self.inner.db, &self.inner.run_id).await, + } + } + pub async fn get_event(&self, seq: u32) -> Result> { get_event(&self.inner.db, &self.inner.run_id, seq).await } @@ -569,11 +628,48 @@ 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). 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 { @@ -581,12 +677,38 @@ 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 + /// end of the run's event namespace. + async fn seek_before(db: &R, run_id: &RunId, start_seq: u32, end_seq: u32) -> Result + 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: 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; @@ -615,10 +737,30 @@ 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 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 scan = EventScan::seek(db, run_id, start_seq).await?; + let mut scan = match end_seq { + Some(end_seq) => EventScan::seek_before(db, run_id, start_seq, end_seq).await?, + None => EventScan::seek(db, run_id, start_seq).await?, + }; let mut events = Vec::new(); while events.len() < max_events { let Some((seq, value)) = scan.next().await? else { @@ -923,6 +1065,229 @@ 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_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 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()); + 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; diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index be21fe2dd..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; @@ -1564,6 +1550,56 @@ 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(); + loop { + let remaining = fetch_target - descending_events.len(); + let (page_events, has_more) = self + .fetch_run_events_page( + run_id, + EventPageCursor::Descending { before_seq }, + Some(remaining), + ) + .await?; + + let keeps_descending = descending_events + .last() + .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)); + } + + before_seq = page_events.last().map(|event| event.seq); + descending_events.extend(page_events); + if descending_events.len() >= fetch_target || !has_more || before_seq.is_none() { + break; + } + } + + 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, @@ -1578,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; @@ -1608,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, @@ -2085,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)) } @@ -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];