From f6ef6b375c9a6adc0a2d64f4e2160e166abdd5f5 Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Tue, 28 Jul 2026 19:54:26 -0700 Subject: [PATCH 1/5] feat(datagen): project fold trajectory started/completed sets to the API `DatagenTrajectory` already tracks `started` and `completed` position sets alongside `ordered`, but `folded_item_to_dto` projected only `ordered`, so neither set was reachable from the API, the client, or Python. Spec resume needs both: `completed` gates STEP_COMPLETED re-emission, `started` gates STEP_STARTED, and `started \ completed` is exactly the set of driver frames that were open when a writer died. Add `DatagenStreamPositionDto` and carry both sets on `FoldedDatagenItemDto`, projected in a deterministic order (the core sets are unordered). Both fields are `#[serde(default, skip_serializing_if = "Vec::is_empty")]`, so this is a backward-compatible wire addition: an older server simply yields empty sets rather than failing to decode. Co-Authored-By: Claude Opus 5 (1M context) --- crates/lance-context-api/src/lib.rs | 20 ++++++++++ crates/lance-context-core/src/api_impl.rs | 47 ++++++++++++++++++++--- crates/lance-context-core/src/datagen.rs | 20 ++++++++++ crates/lance-context/src/lib.rs | 10 ++--- python/src/lib.rs | 30 +++++++++++++-- 5 files changed, 112 insertions(+), 15 deletions(-) diff --git a/crates/lance-context-api/src/lib.rs b/crates/lance-context-api/src/lib.rs index e7df682..b9b3833 100644 --- a/crates/lance-context-api/src/lib.rs +++ b/crates/lance-context-api/src/lib.rs @@ -1155,6 +1155,19 @@ pub struct DatagenStepCursorDto { pub item_seq: i64, } +/// A position within one stream's step tree, without the `item_seq` a cursor carries. +/// Mirrors the Python `StepPosition` wire dict. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DatagenStreamPositionDto { + pub step_name: String, + pub step_kind: String, + pub step_index: i64, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub enclosing_step: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub selector_step: Option, +} + /// An item reconstructed by folding its events into latest state. /// Mirrors the Python `FoldedItem` wire dict. #[derive(Debug, Clone, Serialize, Deserialize)] @@ -1168,6 +1181,13 @@ pub struct FoldedDatagenItemDto { pub last_attempt: i32, pub fields: std::collections::BTreeMap, pub trajectory: Vec, + /// Positions with a STEP_STARTED — gates driver-frame (re-)opening on resume. `started` + /// minus `completed` = frames that were open when the process died. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub started: Vec, + /// Positions with a STEP_COMPLETED — gates STEP_COMPLETED re-emission on resume. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub completed: Vec, #[serde(default, skip_serializing_if = "Option::is_none")] pub query_tags: Option, /// `field_name -> event_id` for the folded blob fields, so a caller can resolve a blob by field diff --git a/crates/lance-context-core/src/api_impl.rs b/crates/lance-context-core/src/api_impl.rs index 8762dfc..b8f4a98 100644 --- a/crates/lance-context-core/src/api_impl.rs +++ b/crates/lance-context-core/src/api_impl.rs @@ -7,17 +7,18 @@ use lance_context_api::{ AddRolloutsResponse, AddRowsResponse, CompactRequest, CompactResponse, CompactStatsResponse, ContextError, ContextResult, ContextStoreApi, DatagenEventDto, DatagenFailureDto, DatagenFieldStateDto, DatagenRootItemStatusesResponse, DatagenStepCursorDto, DatagenStoreApi, - DatagenValueDto, DeleteRecordResponse, FoldedDatagenItemDto, GenericStoreApi, RecordDto, - RecordPatchDto, RelationshipDto, RetrieveRequest, RetrieveResultDto, RolloutRecordDto, - RolloutStoreApi, SchemaSpec, SearchRequest, SearchResultDto, StateMetadataDto, - UpdateRecordRequest, UpdateRecordResponse, UpsertRecordRequest, UpsertRecordResponse, - UpsertRecordsRequest, UpsertRecordsResponse, UpsertResultDto, + DatagenStreamPositionDto, DatagenValueDto, DeleteRecordResponse, FoldedDatagenItemDto, + GenericStoreApi, RecordDto, RecordPatchDto, RelationshipDto, RetrieveRequest, + RetrieveResultDto, RolloutRecordDto, RolloutStoreApi, SchemaSpec, SearchRequest, + SearchResultDto, StateMetadataDto, UpdateRecordRequest, UpdateRecordResponse, + UpsertRecordRequest, UpsertRecordResponse, UpsertRecordsRequest, UpsertRecordsResponse, + UpsertResultDto, }; use crate::datagen::{ DatagenBlobValue, DatagenEvent, DatagenEventType, DatagenFailure, DatagenFieldState, DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, DatagenStepCursor, - DatagenStepKind, DatagenValue, FoldedDatagenItem, + DatagenStepKind, DatagenStreamPosition, DatagenValue, FoldedDatagenItem, }; use crate::datagen_store::DatagenStore; use crate::generic_codec::Row; @@ -972,6 +973,8 @@ pub fn folded_item_to_dto(item: &FoldedDatagenItem) -> FoldedDatagenItemDto { .map(|(name, state)| (name.clone(), field_state_to_dto(state))) .collect(), trajectory: item.trajectory.ordered.iter().map(cursor_to_dto).collect(), + started: position_set_to_dto(&item.trajectory.started), + completed: position_set_to_dto(&item.trajectory.completed), query_tags: item.query_tags.clone(), blob_event_ids: item.blob_event_ids.clone(), } @@ -1003,6 +1006,38 @@ fn cursor_to_dto(cursor: &DatagenStepCursor) -> DatagenStepCursorDto { } } +fn position_to_dto(position: &DatagenStreamPosition) -> DatagenStreamPositionDto { + DatagenStreamPositionDto { + step_name: position.step.name.clone(), + step_kind: position.step.kind.as_str().to_string(), + step_index: position.index, + enclosing_step: position.enclosing.clone(), + selector_step: position.selector.clone(), + } +} + +/// Project a fold position set to DTOs in a deterministic order (the sets are unordered). +fn position_set_to_dto( + positions: &std::collections::HashSet, +) -> Vec { + let mut dtos: Vec = positions.iter().map(position_to_dto).collect(); + dtos.sort_by(|a, b| { + ( + &a.step_name, + a.step_index, + &a.enclosing_step, + &a.selector_step, + ) + .cmp(&( + &b.step_name, + b.step_index, + &b.enclosing_step, + &b.selector_step, + )) + }); + dtos +} + fn root_item_statuses_to_dto( statuses: &DatagenRootItemStatuses, ) -> DatagenRootItemStatusesResponse { diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 5183d71..856c44f 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -1243,6 +1243,26 @@ mod tests { assert_eq!(folded.last_item_seq, 2); } + #[test] + fn dto_projection_carries_started_and_completed_sets() { + let events = [ + created(0), + driver_started(1, "main", 0, None), + leaf_completed(2, "gen", 0, Some("main")), + ]; + let folded = fold_datagen_events(&events).unwrap().unwrap(); + let dto = crate::api_impl::folded_item_to_dto(&folded); + + // `trajectory` stays the ordered cursor list; the two sets ride alongside it so the + // caller can gate STEP_STARTED / STEP_COMPLETED re-emission on resume. + assert_eq!(dto.trajectory.len(), 1); + assert_eq!(dto.started.len(), 1); + assert_eq!(dto.started[0].step_name, "main"); + assert_eq!(dto.completed.len(), 1); + assert_eq!(dto.completed[0].step_name, "gen"); + assert_eq!(dto.completed[0].enclosing_step.as_deref(), Some("main")); + } + #[test] fn field_set_is_last_writer_wins_and_append_accumulates() { let mut set_v1 = leaf_completed(1, "gen", 0, Some("main")); diff --git a/crates/lance-context/src/lib.rs b/crates/lance-context/src/lib.rs index b0ee4e7..436ff2e 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -23,11 +23,11 @@ pub use lance_context_api::{ ColumnType, CompactRequest, CompactResponse, CompactStatsResponse, ContextError, ContextResult, ContextStoreApi, CreateDatagenStoreRequest, CreateGenericStoreRequest, CreateRolloutStoreRequest, DatagenEventDto, DatagenFailureDto, DatagenFieldStateDto, - DatagenRootItemStatusesResponse, DatagenStepCursorDto, DatagenStoreApi, DatagenValueDto, - DeleteRecordResponse, FoldedDatagenItemDto, GenericStoreApi, GenericStoreInfo, RecordDto, - RelationshipDto, RetrieveRequest, RetrieveResponse, RetrieveResultDto, RolloutRecordDto, - RolloutStoreApi, SchemaSpec, SearchResultDto, UpsertRecordRequest, UpsertRecordResponse, - ID_COLUMN, + DatagenRootItemStatusesResponse, DatagenStepCursorDto, DatagenStoreApi, + DatagenStreamPositionDto, DatagenValueDto, DeleteRecordResponse, FoldedDatagenItemDto, + GenericStoreApi, GenericStoreInfo, RecordDto, RelationshipDto, RetrieveRequest, + RetrieveResponse, RetrieveResultDto, RolloutRecordDto, RolloutStoreApi, SchemaSpec, + SearchResultDto, UpsertRecordRequest, UpsertRecordResponse, ID_COLUMN, }; #[cfg(feature = "remote")] diff --git a/python/src/lib.rs b/python/src/lib.rs index 7247c1d..f48fcf9 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -17,10 +17,10 @@ use lance_context::{ DatagenErrorInfo, DatagenEventDto, DatagenFailureDto, DatagenFieldChange, DatagenFieldStateDto, DatagenItemNode as UnifiedDatagenItemNode, DatagenItemTree as UnifiedDatagenItemTree, DatagenStepCursorDto, DatagenStore as UnifiedDatagenStore, DatagenStoreApi, - DatagenStreamPosition, DatagenStreamWriter as CoreDatagenStreamWriter, DatagenValueDto, - DatagenWriteContext, FoldedDatagenItemDto, GenericStore as UnifiedGenericStore, - GenericStoreApi, RolloutRecordDto, RolloutStore as UnifiedRolloutStore, RolloutStoreApi, - SchemaSpec, + DatagenStreamPosition, DatagenStreamPositionDto, + DatagenStreamWriter as CoreDatagenStreamWriter, DatagenValueDto, DatagenWriteContext, + FoldedDatagenItemDto, GenericStore as UnifiedGenericStore, GenericStoreApi, RolloutRecordDto, + RolloutStore as UnifiedRolloutStore, RolloutStoreApi, SchemaSpec, }; use lance_context_api::{ AddRecordRequest, CompactRequest, CompactResponse, CompactStatsResponse, ContextStoreApi, @@ -3173,6 +3173,18 @@ fn folded_item_to_py(py: Python<'_>, item: &FoldedDatagenItemDto) -> PyResult, cursor: &DatagenStepCursorDto) -> PyResult, position: &DatagenStreamPositionDto) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("step_name", &position.step_name)?; + dict.set_item("step_kind", &position.step_kind)?; + dict.set_item("step_index", position.step_index)?; + dict.set_item("enclosing_step", position.enclosing_step.clone())?; + dict.set_item("selector_step", position.selector_step.clone())?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + fn failure_to_py(py: Python<'_>, failure: &DatagenFailureDto) -> PyResult { let dict = PyDict::new(py); dict.set_item("at", cursor_to_py(py, &failure.at)?)?; From bd7021f8ef65ce2a5654979302b47ae0d7f430f0 Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Tue, 28 Jul 2026 20:37:25 -0700 Subject: [PATCH 2/5] feat(datagen): enforce field value-kind stability and export DatagenItemId Two follow-ups on the datagen fold/binding surface: - fold now rejects a field whose value kind drifts between writes, for both FIELD_SET (last-writer-wins) and FIELD_APPEND (within one list). Previously only the SET-vs-APPEND mix was caught, so a field could silently hold mixed-kind values. - DatagenItemId is registered as a pyclass and re-exported from the Python package, so callers compose and parse sub-item ids through the store's own path format instead of hand-splicing strings. Co-Authored-By: Claude Opus 5 (1M context) --- crates/lance-context-core/src/datagen.rs | 67 +++++++++++++++++- python/python/lance_context/__init__.py | 2 + python/python/lance_context/api.py | 3 + python/src/lib.rs | 86 ++++++++++++++++++++++++ 4 files changed, 156 insertions(+), 2 deletions(-) diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 856c44f..52b427a 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -1069,6 +1069,18 @@ fn apply_event(item: &mut FoldedDatagenItem, event: &DatagenEvent) -> Result<(), DatagenEventType::FieldSet => { let field_name = event.field_name.clone().unwrap(); let value = event.value.clone().unwrap(); + // Group D: a field's value kind is fixed by its first write. A later SET that + // drifts to another kind means the pipeline wrote the field inconsistently. + if let Some(DatagenFieldState::Set(existing)) = item.fields.get(&field_name) { + if existing.kind() != value.kind() { + return Err(format!( + "field '{}' changes value kind from {} to {}", + field_name, + existing.kind(), + value.kind() + )); + } + } record_blob_event_id(item, &field_name, &value, &event.event_id); item.fields .insert(field_name, DatagenFieldState::Set(value)); @@ -1077,12 +1089,25 @@ fn apply_event(item: &mut FoldedDatagenItem, event: &DatagenEvent) -> Result<(), let field_name = event.field_name.clone().unwrap(); let value = event.value.clone().unwrap(); record_blob_event_id(item, &field_name, &value, &event.event_id); - match item.fields.entry(field_name) { + match item.fields.entry(field_name.clone()) { std::collections::btree_map::Entry::Vacant(entry) => { entry.insert(DatagenFieldState::Appended(vec![value])); } std::collections::btree_map::Entry::Occupied(mut entry) => match entry.get_mut() { - DatagenFieldState::Appended(values) => values.push(value), + DatagenFieldState::Appended(values) => { + // Group D: the same kind rule holds within one appended list. + if let Some(existing) = values.first() { + if existing.kind() != value.kind() { + return Err(format!( + "field '{}' changes value kind from {} to {}", + field_name, + existing.kind(), + value.kind() + )); + } + } + values.push(value); + } DatagenFieldState::Set(_) => { return Err(format!( "field '{}' mixes FIELD_SET and FIELD_APPEND", @@ -1306,6 +1331,44 @@ mod tests { ); } + #[test] + fn field_value_kind_drift_is_rejected() { + let mut set_str = leaf_completed(1, "gen", 0, Some("main")); + set_str.event_type = DatagenEventType::FieldSet; + set_str.field_name = Some("draft".to_string()); + set_str.field_type = Some("str".to_string()); + set_str.codec_version = Some(1); + set_str.value = Some(DatagenValue::Str("v1".to_string())); + + let mut set_int = set_str.clone(); + set_int.item_seq = 2; + set_int.checkpoint_id = "c2".to_string(); + set_int.event_id = datagen_event_id("5", "c2", 0); + set_int.field_type = Some("int".to_string()); + set_int.value = Some(DatagenValue::Int(7)); + + let err = fold_datagen_events(&[created(0), set_str, set_int]).unwrap_err(); + assert!(err.contains("changes value kind"), "{err}"); + + // Same rule inside an appended list. + let mut append_str = leaf_completed(3, "b1", 0, Some("body")); + append_str.event_type = DatagenEventType::FieldAppend; + append_str.field_name = Some("revisions".to_string()); + append_str.field_type = Some("str".to_string()); + append_str.codec_version = Some(1); + append_str.value = Some(DatagenValue::Str("a".to_string())); + + let mut append_int = append_str.clone(); + append_int.item_seq = 4; + append_int.checkpoint_id = "c4".to_string(); + append_int.event_id = datagen_event_id("5", "c4", 0); + append_int.field_type = Some("int".to_string()); + append_int.value = Some(DatagenValue::Int(1)); + + let err = fold_datagen_events(&[created(0), append_str, append_int]).unwrap_err(); + assert!(err.contains("changes value kind"), "{err}"); + } + #[test] fn terminal_sets_lifecycle_status() { let mut terminal = event(2, DatagenEventType::Terminal); diff --git a/python/python/lance_context/__init__.py b/python/python/lance_context/__init__.py index 190ee97..ad4b054 100644 --- a/python/python/lance_context/__init__.py +++ b/python/python/lance_context/__init__.py @@ -5,6 +5,7 @@ AsyncRolloutStore, Context, ContextNamespace, + DatagenItemId, DatagenStore, DatagenStreamWriter, EmbeddingProvider, @@ -24,6 +25,7 @@ "AsyncRolloutStore", "Context", "ContextNamespace", + "DatagenItemId", "DatagenStore", "DatagenStreamWriter", "GenericStore", diff --git a/python/python/lance_context/api.py b/python/python/lance_context/api.py index 03273d0..35f1d9f 100644 --- a/python/python/lance_context/api.py +++ b/python/python/lance_context/api.py @@ -15,6 +15,9 @@ from ._internal import ( # pyright: ignore[reportMissingImports] ContextNamespace as _ContextNamespace, ) +from ._internal import ( # pyright: ignore[reportMissingImports] + DatagenItemId as DatagenItemId, +) from ._internal import ( # pyright: ignore[reportMissingImports] DatagenStore as _DatagenStore, ) diff --git a/python/src/lib.rs b/python/src/lib.rs index f48fcf9..33d6a2e 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -2790,6 +2790,91 @@ impl DatagenStore { } } +/// A structured datagen item id. Root ids come from the executor's source key; sub-item +/// ids extend a parent with one fan-out segment (`5/expand:0`). `str(id)` is the stored +/// path form. Pure and client-side — composing an id does no I/O. +#[pyclass] +#[derive(Clone)] +struct DatagenItemId { + inner: CoreDatagenItemId, +} + +#[pymethods] +impl DatagenItemId { + /// Build a root id from the executor's source key. + #[staticmethod] + fn from_source_key(key: &str) -> Self { + Self { + inner: CoreDatagenItemId::from_source_key(key), + } + } + + /// Parse a stored path string (`"5/expand:0"`) back into a structured id. + #[staticmethod] + fn parse(path: &str) -> PyResult { + Ok(Self { + inner: CoreDatagenItemId::parse(path).map_err(to_py_err)?, + }) + } + + /// Extend this id with one fan-out segment -> the sub-item's id. + fn child(&self, origin_step: &str, branch_idx: i64) -> Self { + Self { + inner: self.inner.child(origin_step, branch_idx), + } + } + + /// The parent stream's id (`None` on a root). + fn parent(&self) -> Option { + self.inner.parent().map(|inner| Self { inner }) + } + + /// The root of this id's tree (== self if root). + fn root(&self) -> Self { + Self { + inner: self.inner.root(), + } + } + + /// The fan-out step that created this sub-item (`None` on a root). + #[getter] + fn origin_step(&self) -> Option { + self.inner.origin_step().map(str::to_string) + } + + /// Which branch this sub-item is (`None` on a root). + #[getter] + fn branch_idx(&self) -> Option { + self.inner.branch_idx() + } + + /// Whether this id names a root stream. + #[getter] + fn is_root(&self) -> bool { + self.inner.is_root() + } + + fn __str__(&self) -> String { + self.inner.to_string() + } + + fn __repr__(&self) -> String { + format!("DatagenItemId({:?})", self.inner.to_string()) + } + + fn __eq__(&self, other: &Self) -> bool { + self.inner == other.inner + } + + fn __hash__(&self) -> u64 { + use std::collections::hash_map::DefaultHasher; + use std::hash::{Hash, Hasher}; + let mut hasher = DefaultHasher::new(); + self.inner.hash(&mut hasher); + hasher.finish() + } +} + /// A per-stream write handle. Owns the bookkeeping columns the caller should never /// touch (`item_seq`, `attempt`, `checkpoint_id`, `event_id`), stamping them onto /// every event it emits. Each method returns the event dict(s) to persist — the @@ -3581,6 +3666,7 @@ fn _internal(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add_class::()?; m.add_class::()?; Ok(()) } From 989a248ebd2da4d0376c44830c456184d331f05e Mon Sep 17 00:00:00 2001 From: yangjunz Date: Tue, 28 Jul 2026 22:01:55 -0700 Subject: [PATCH 3/5] feat(datagen): whole-store run overview and lazy/eager blob projection Two gaps in the datagen delta-log read surface: * `DatagenRunOverview` aggregates a whole store (one experiment) into root-granularity status counts, per-step completion counts, and a failure roll-up grouped by `run_id`. Each run bucket carries its own error-type counts plus a capped, deterministic sample of failing root item ids to drill into. Exposed as `DatagenStore::overview`, `GET /datagen/{name}/overview`, and `DatagenStore.overview()` in Python. * `DatagenBlobProjection` makes the blob projection explicit. `fold_item` stays lazy (bytes absent, resolved through `get_blob`); `fold_item_with`/`fold_item_with_blobs`/`fold_item(load_blobs=True)` materialize them inline, reading the payload column only when asked. Also fixes the lint CI failure by wrapping the PyO3 `DatagenItemId` in a thin Python class, matching the existing `DatagenStreamWriter` pattern. The wire additions are backward compatible: new DTO fields default, so no lockstep server upgrade is required. --- crates/lance-context-api/src/lib.rs | 46 +++ crates/lance-context-client/src/lib.rs | 43 +++ crates/lance-context-core/src/api_impl.rs | 73 ++++- crates/lance-context-core/src/datagen.rs | 269 +++++++++++++++++- .../lance-context-core/src/datagen_store.rs | 148 +++++++++- .../src/routes/datagen.rs | 25 +- crates/lance-context-server/src/routes/mod.rs | 4 + crates/lance-context/src/lib.rs | 13 +- crates/lance-context/src/unified_datagen.rs | 14 +- python/python/lance_context/api.py | 91 +++++- python/src/lib.rs | 73 ++++- 11 files changed, 766 insertions(+), 33 deletions(-) diff --git a/crates/lance-context-api/src/lib.rs b/crates/lance-context-api/src/lib.rs index b9b3833..1a17d92 100644 --- a/crates/lance-context-api/src/lib.rs +++ b/crates/lance-context-api/src/lib.rs @@ -327,6 +327,16 @@ pub trait DatagenStoreApi { item_id: &str, ) -> impl Future>> + Send; + /// Like [`DatagenStoreApi::fold_item`], but `load_blobs` selects the blob projection: + /// `false` (the `fold_item` default) leaves blob fields lazy — bytes absent, resolved later + /// through `get_blob` — while `true` materializes them inline, at the cost of reading the + /// payload column. + fn fold_item_with_blobs( + &self, + item_id: &str, + load_blobs: bool, + ) -> impl Future>> + Send; + fn root_item_statuses( &self, root_item_ids: &[String], @@ -346,6 +356,10 @@ pub trait DatagenStoreApi { root_item_id: &str, ) -> impl Future>> + Send; + /// Aggregate the whole log into a run overview: per-status root item counts, + /// failure counts by error type, and completed-step counts. + fn overview(&self) -> impl Future> + Send; + /// Materialize one FIELD_* event's offloaded blob bytes by event id. /// Returns `None` when the event or its payload is absent. fn get_blob( @@ -1221,6 +1235,38 @@ pub struct DatagenFailureDto { pub traceback: Option, } +/// Whole-run aggregation over a datagen log. `items` counts root items only +/// (`running + completed + filtered`); `failures` counts FAILED events, which are +/// non-terminal, so a failed-then-retried item still counts as `running`. +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct DatagenRunOverviewDto { + pub items: usize, + pub running: usize, + pub completed: usize, + pub filtered: usize, + pub failures: usize, + /// FAILED-event count per `error_type`. + #[serde(default)] + pub failures_by_error_type: std::collections::BTreeMap, + /// Failure roll-up grouped by the `run_id` that emitted the FAILED event. + #[serde(default)] + pub failures_by_run: std::collections::BTreeMap, + /// STEP_COMPLETED count per step name, across every item. + #[serde(default)] + pub completed_steps: std::collections::BTreeMap, +} + +/// One `run_id`'s slice of an overview's failure roll-up, with a capped sample of failing root +/// item ids to drill into. +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct DatagenFailureBucketDto { + pub failures: usize, + #[serde(default)] + pub failures_by_error_type: std::collections::BTreeMap, + #[serde(default)] + pub sample_root_item_ids: Vec, +} + #[derive(Debug, Serialize, Deserialize)] pub struct ListDatagenFailuresResponse { pub failures: Vec, diff --git a/crates/lance-context-client/src/lib.rs b/crates/lance-context-client/src/lib.rs index b024224..36db341 100644 --- a/crates/lance-context-client/src/lib.rs +++ b/crates/lance-context-client/src/lib.rs @@ -548,6 +548,19 @@ impl DatagenStoreApi for RemoteDatagenStore { Ok(resp.item) } + async fn fold_item_with_blobs( + &self, + item_id: &str, + load_blobs: bool, + ) -> ContextResult> { + let resp = self + .client + .fold_datagen_item_with_blobs(&self.store_name, item_id, load_blobs) + .await + .map_err(to_ctx_err)?; + Ok(resp.item) + } + async fn root_item_statuses( &self, root_item_ids: &[String], @@ -558,6 +571,13 @@ impl DatagenStoreApi for RemoteDatagenStore { .map_err(to_ctx_err) } + async fn overview(&self) -> ContextResult { + self.client + .datagen_overview(&self.store_name) + .await + .map_err(to_ctx_err) + } + async fn item_failures(&self, item_id: &str) -> ContextResult> { let resp = self .client @@ -1286,10 +1306,23 @@ impl ContextClient { &self, name: &str, item_id: &str, + ) -> Result { + self.fold_datagen_item_with_blobs(name, item_id, false) + .await + } + + /// Fold an item, choosing the blob projection: `load_blobs` materializes blob bytes inline + /// instead of leaving them to a later `get_blob`. + pub async fn fold_datagen_item_with_blobs( + &self, + name: &str, + item_id: &str, + load_blobs: bool, ) -> Result { let resp = self .http .get(self.url(&format!("/datagen/{}/items/{}", name, item_id))) + .query(&[("load_blobs", load_blobs)]) .send() .await?; Self::handle_response(resp).await @@ -1308,6 +1341,16 @@ impl ContextClient { Self::handle_response(resp).await } + /// Aggregate the whole datagen log into a run overview. + pub async fn datagen_overview(&self, name: &str) -> Result { + let resp = self + .http + .get(self.url(&format!("/datagen/{}/overview", name))) + .send() + .await?; + Self::handle_response(resp).await + } + pub async fn datagen_root_item_statuses( &self, name: &str, diff --git a/crates/lance-context-core/src/api_impl.rs b/crates/lance-context-core/src/api_impl.rs index b8f4a98..c570383 100644 --- a/crates/lance-context-core/src/api_impl.rs +++ b/crates/lance-context-core/src/api_impl.rs @@ -5,20 +5,21 @@ use uuid::Uuid; use lance_context_api::{ AddDatagenEventsResponse, AddRecordRequest, AddRecordsResponse, AddRolloutRequest, AddRolloutsResponse, AddRowsResponse, CompactRequest, CompactResponse, CompactStatsResponse, - ContextError, ContextResult, ContextStoreApi, DatagenEventDto, DatagenFailureDto, - DatagenFieldStateDto, DatagenRootItemStatusesResponse, DatagenStepCursorDto, DatagenStoreApi, - DatagenStreamPositionDto, DatagenValueDto, DeleteRecordResponse, FoldedDatagenItemDto, - GenericStoreApi, RecordDto, RecordPatchDto, RelationshipDto, RetrieveRequest, - RetrieveResultDto, RolloutRecordDto, RolloutStoreApi, SchemaSpec, SearchRequest, - SearchResultDto, StateMetadataDto, UpdateRecordRequest, UpdateRecordResponse, - UpsertRecordRequest, UpsertRecordResponse, UpsertRecordsRequest, UpsertRecordsResponse, - UpsertResultDto, + ContextError, ContextResult, ContextStoreApi, DatagenEventDto, DatagenFailureBucketDto, + DatagenFailureDto, DatagenFieldStateDto, DatagenRootItemStatusesResponse, + DatagenRunOverviewDto, DatagenStepCursorDto, DatagenStoreApi, DatagenStreamPositionDto, + DatagenValueDto, DeleteRecordResponse, FoldedDatagenItemDto, GenericStoreApi, RecordDto, + RecordPatchDto, RelationshipDto, RetrieveRequest, RetrieveResultDto, RolloutRecordDto, + RolloutStoreApi, SchemaSpec, SearchRequest, SearchResultDto, StateMetadataDto, + UpdateRecordRequest, UpdateRecordResponse, UpsertRecordRequest, UpsertRecordResponse, + UpsertRecordsRequest, UpsertRecordsResponse, UpsertResultDto, }; use crate::datagen::{ - DatagenBlobValue, DatagenEvent, DatagenEventType, DatagenFailure, DatagenFieldState, - DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, DatagenStepCursor, - DatagenStepKind, DatagenStreamPosition, DatagenValue, FoldedDatagenItem, + DatagenBlobProjection, DatagenBlobValue, DatagenEvent, DatagenEventType, DatagenFailure, + DatagenFieldState, DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, + DatagenRunOverview, DatagenStepCursor, DatagenStepKind, DatagenStreamPosition, DatagenValue, + FoldedDatagenItem, }; use crate::datagen_store::DatagenStore; use crate::generic_codec::Row; @@ -772,6 +773,25 @@ impl DatagenStoreApi for DatagenStore { }) } + async fn fold_item_with_blobs( + &self, + item_id: &str, + load_blobs: bool, + ) -> ContextResult> { + let blobs = if load_blobs { + DatagenBlobProjection::Eager + } else { + DatagenBlobProjection::Lazy + }; + let lookup = DatagenStore::fold_item_with(self, item_id, blobs) + .await + .map_err(to_ctx_err)?; + Ok(match lookup { + DatagenItemLookup::NeverStarted => None, + DatagenItemLookup::Found(item) => Some(folded_item_to_dto(&item)), + }) + } + async fn root_item_statuses( &self, root_item_ids: &[String], @@ -783,6 +803,11 @@ impl DatagenStoreApi for DatagenStore { Ok(root_item_statuses_to_dto(&statuses)) } + async fn overview(&self) -> ContextResult { + let overview = DatagenStore::overview(self).await.map_err(to_ctx_err)?; + Ok(run_overview_to_dto(&overview)) + } + async fn item_failures(&self, item_id: &str) -> ContextResult> { let failures = DatagenStore::item_failures(self, item_id) .await @@ -1049,6 +1074,32 @@ fn root_item_statuses_to_dto( } } +fn run_overview_to_dto(overview: &DatagenRunOverview) -> DatagenRunOverviewDto { + DatagenRunOverviewDto { + items: overview.items, + running: overview.running, + completed: overview.completed, + filtered: overview.filtered, + failures: overview.failures, + failures_by_error_type: overview.failures_by_error_type.clone(), + failures_by_run: overview + .failures_by_run + .iter() + .map(|(run_id, bucket)| { + ( + run_id.clone(), + DatagenFailureBucketDto { + failures: bucket.failures, + failures_by_error_type: bucket.failures_by_error_type.clone(), + sample_root_item_ids: bucket.sample_root_item_ids.clone(), + }, + ) + }) + .collect(), + completed_steps: overview.completed_steps.clone(), + } +} + fn failure_to_dto(failure: &DatagenFailure) -> DatagenFailureDto { DatagenFailureDto { at: cursor_to_dto(&failure.at), diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 52b427a..9211d49 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -764,6 +764,108 @@ impl DatagenItemTree { } } +/// Whole-run aggregation over a datagen log: how many items are in each lifecycle state, how far +/// they got, and where they failed. Built purely from events — no I/O. +/// +/// `items` counts *root* items only (fan-out sub-items roll up under their root), so +/// `running + completed + filtered` equals the number of roots the log has seen. `failures` counts +/// FAILED events, which are non-terminal: a failed-then-retried item still counts as `running`. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct DatagenRunOverview { + /// Root items seen in the log (`running + completed + filtered`). + pub items: usize, + pub running: usize, + pub completed: usize, + pub filtered: usize, + /// Total FAILED events across every item, including sub-items and retried attempts. + pub failures: usize, + /// FAILED-event count per `error_type`, so the dominant failure mode is one lookup away. + pub failures_by_error_type: BTreeMap, + /// Failure roll-up grouped by the `run_id` that emitted the FAILED event. One store spans every + /// attempt at an experiment, so this separates "this run's failures" from historical ones. + pub failures_by_run: BTreeMap, + /// STEP_COMPLETED count per step name, across every item — the run's step-level progress. + pub completed_steps: BTreeMap, +} + +/// One `run_id`'s slice of the failure roll-up: how many, of what kind, and a handful of root item +/// ids to open next. The sample is capped at [`FAILURE_SAMPLE_LIMIT`] and is deterministic (roots in +/// id order), so an overview stays small no matter how wide the run is. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct DatagenFailureBucket { + pub failures: usize, + pub failures_by_error_type: BTreeMap, + /// Up to [`FAILURE_SAMPLE_LIMIT`] distinct root item ids that failed under this run. + pub sample_root_item_ids: Vec, +} + +/// How many root item ids each [`DatagenFailureBucket`] samples. +pub const FAILURE_SAMPLE_LIMIT: usize = 5; + +impl DatagenRunOverview { + /// Aggregate a whole run's events. Events for many items may be interleaved; they are grouped by + /// `item_id` and folded independently, exactly like [`DatagenItemTree::build`]. + pub fn build(events: &[DatagenEvent]) -> Result { + let mut by_item: BTreeMap> = BTreeMap::new(); + for event in events { + by_item + .entry(event.item_id.clone()) + .or_default() + .push(event.clone()); + } + + let mut overview = Self::default(); + for item_events in by_item.values() { + let Some(item) = fold_datagen_events(item_events)? else { + continue; + }; + for cursor in &item.trajectory.ordered { + *overview + .completed_steps + .entry(cursor.position.step.name.clone()) + .or_default() += 1; + } + for failure in datagen_failures(item_events)? { + overview.failures += 1; + *overview + .failures_by_error_type + .entry(failure.error.error_type.clone()) + .or_default() += 1; + let bucket = overview + .failures_by_run + .entry(failure.run_id.clone()) + .or_default(); + bucket.failures += 1; + *bucket + .failures_by_error_type + .entry(failure.error.error_type.clone()) + .or_default() += 1; + // Items are visited in id order, so the sample is the same on every rebuild. + let root = item.root_item_id.to_string(); + if bucket.sample_root_item_ids.len() < FAILURE_SAMPLE_LIMIT + && !bucket.sample_root_item_ids.contains(&root) + { + bucket.sample_root_item_ids.push(root); + } + } + // Only roots are counted as items; a sub-item's outcome rolls up under its root. + if item.parent_item_id.is_some() { + continue; + } + overview.items += 1; + match item.status { + DatagenItemStatus::Running => overview.running += 1, + DatagenItemStatus::Completed => overview.completed += 1, + DatagenItemStatus::Filtered => overview.filtered += 1, + // A folded status is never Failed (failures live in the failure lens), but count it + // as still running rather than silently dropping the item from `items`. + DatagenItemStatus::Failed => overview.running += 1, + } + } + Ok(overview) + } +} + /// Result of a resumption fold. `NeverStarted` (no ITEM_CREATED) is the fresh-vs-restore fork the /// executor acts on; `Found` carries the folded item (whose `status` is the lifecycle status). #[derive(Debug, Clone, PartialEq)] @@ -943,9 +1045,35 @@ pub fn datagen_event_id(item_id: &str, checkpoint_id: &str, ordinal: u32) -> Str Uuid::new_v5(&Uuid::NAMESPACE_OID, input.as_bytes()).to_string() } +/// Whether a fold keeps blob field bytes it happens to have, or drops them to a pointer. +/// +/// The log's `value_blob` column is normally projected away on read, so a folded blob field is a +/// lazy pointer resolved through `blob_event_ids` + `get_blob`. `Eager` keeps whatever bytes the +/// events already carry, which lets a caller that scanned the blob column fold the payload in one +/// pass instead of one `get_blob` round trip per field. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub enum DatagenBlobProjection { + /// Blob fields fold to a pointer: `bytes` is dropped even when present. The default. + #[default] + Lazy, + /// Blob fields keep the bytes the events carry (`None` when the scan projected them away). + Eager, +} + /// Fold an item's events into its latest state. Returns `None` if there is no ITEM_CREATED (the item -/// was never started). +/// was never started). Blob fields fold lazily — see [`fold_datagen_events_with`] to keep bytes. pub fn fold_datagen_events(events: &[DatagenEvent]) -> Result, String> { + fold_datagen_events_with(events, DatagenBlobProjection::Lazy) +} + +/// Fold an item's events under an explicit blob projection. +/// +/// `DatagenBlobProjection::Eager` is only useful when the events were read with the blob column +/// projected in; with the default read path it folds identically to `Lazy`. +pub fn fold_datagen_events_with( + events: &[DatagenEvent], + blobs: DatagenBlobProjection, +) -> Result, String> { let ordered = normalize_events(events)?; let Some(first) = ordered.first() else { return Ok(None); @@ -957,9 +1085,28 @@ pub fn fold_datagen_events(events: &[DatagenEvent]) -> Result strip(value), + DatagenFieldState::Appended(values) => values.iter_mut().for_each(strip), + } + } +} + /// The ordered completed-step cursors of an item, in `item_seq` order. pub fn datagen_trajectory(events: &[DatagenEvent]) -> Result, String> { Ok(fold_datagen_events(events)? @@ -1755,4 +1902,124 @@ mod tests { assert_eq!(open.len(), 1); assert_eq!(open[0].step.name, "refine"); } + + /// A blob-valued FIELD_SET on item "5", at `seq`. + fn blob_set(seq: i64, field: &str, bytes: &[u8]) -> DatagenEvent { + let mut set = leaf_completed(seq, "gen", 0, Some("main")); + set.event_type = DatagenEventType::FieldSet; + set.field_name = Some(field.to_string()); + set.field_type = Some("blob".to_string()); + set.codec_version = Some(1); + set.value = Some(DatagenValue::Blob(DatagenBlobValue { + bytes: Some(bytes.to_vec()), + size: bytes.len() as i64, + checksum: None, + })); + set + } + + #[test] + fn blob_projection_selects_lazy_or_eager_bytes() { + let events = vec![created(0), blob_set(1, "image", b"png-bytes")]; + + let lazy = fold_datagen_events_with(&events, DatagenBlobProjection::Lazy) + .unwrap() + .unwrap(); + let DatagenFieldState::Set(DatagenValue::Blob(lazy_blob)) = + lazy.fields.get("image").unwrap() + else { + panic!("expected a blob field"); + }; + assert_eq!(lazy_blob.bytes, None); + assert_eq!(lazy_blob.size, 9); + // The default fold is the lazy one. + assert_eq!(fold_datagen_events(&events).unwrap().unwrap(), lazy); + + let eager = fold_datagen_events_with(&events, DatagenBlobProjection::Eager) + .unwrap() + .unwrap(); + let DatagenFieldState::Set(DatagenValue::Blob(eager_blob)) = + eager.fields.get("image").unwrap() + else { + panic!("expected a blob field"); + }; + assert_eq!(eager_blob.bytes.as_deref(), Some(&b"png-bytes"[..])); + } + + #[test] + fn overview_counts_roots_and_rolls_failures_up_by_run() { + // Two roots: "5" completes; "9" stays running with two failures, one of them under a + // second run_id. "5/expand:0" is a sub-item and must not inflate the root counts. + let mut events = vec![created(0), leaf_completed(1, "gen", 0, Some("main"))]; + let mut terminal = event(2, DatagenEventType::Terminal); + terminal.status = Some(DatagenItemStatus::Completed); + events.push(terminal); + + let mut sub_created = created(0); + sub_created.item_id = "5/expand:0".to_string(); + sub_created.parent_item_id = Some("5".to_string()); + sub_created.event_id = datagen_event_id("5/expand:0", "checkpoint-0", 0); + events.push(sub_created); + + let mut other_created = created(0); + other_created.item_id = "9".to_string(); + other_created.root_item_id = "9".to_string(); + other_created.event_id = datagen_event_id("9", "checkpoint-0", 0); + events.push(other_created); + for (seq, run_id, error_type) in [(1, "run-1", "ValueError"), (2, "run-2", "KeyError")] { + let mut failed = leaf_completed(seq, "score", 0, Some("main")); + failed.item_id = "9".to_string(); + failed.root_item_id = "9".to_string(); + failed.event_id = datagen_event_id("9", &format!("checkpoint-{seq}"), 0); + failed.event_type = DatagenEventType::Failed; + failed.run_id = run_id.to_string(); + failed.error_type = Some(error_type.to_string()); + events.push(failed); + } + + let overview = DatagenRunOverview::build(&events).unwrap(); + assert_eq!(overview.items, 2); + assert_eq!(overview.completed, 1); + assert_eq!(overview.running, 1); + assert_eq!(overview.filtered, 0); + assert_eq!(overview.failures, 2); + assert_eq!(overview.failures_by_error_type["ValueError"], 1); + assert_eq!(overview.completed_steps["gen"], 1); + + // Failures group by the run_id that emitted them, each carrying a root-id sample. + assert_eq!(overview.failures_by_run.len(), 2); + let run1 = &overview.failures_by_run["run-1"]; + assert_eq!(run1.failures, 1); + assert_eq!(run1.failures_by_error_type["ValueError"], 1); + assert_eq!(run1.sample_root_item_ids, vec!["9".to_string()]); + assert_eq!(overview.failures_by_run["run-2"].failures, 1); + } + + #[test] + fn overview_failure_sample_is_capped() { + let mut events = Vec::new(); + for idx in 0..(FAILURE_SAMPLE_LIMIT + 3) { + let item_id = format!("item-{idx:02}"); + let mut item_created = created(0); + item_created.item_id = item_id.clone(); + item_created.root_item_id = item_id.clone(); + item_created.event_id = datagen_event_id(&item_id, "checkpoint-0", 0); + events.push(item_created); + + let mut failed = leaf_completed(1, "score", 0, Some("main")); + failed.item_id = item_id.clone(); + failed.root_item_id = item_id.clone(); + failed.event_id = datagen_event_id(&item_id, "checkpoint-1", 0); + failed.event_type = DatagenEventType::Failed; + failed.error_type = Some("ValueError".to_string()); + events.push(failed); + } + + let overview = DatagenRunOverview::build(&events).unwrap(); + assert_eq!(overview.failures, FAILURE_SAMPLE_LIMIT + 3); + let bucket = &overview.failures_by_run["run-1"]; + assert_eq!(bucket.sample_root_item_ids.len(), FAILURE_SAMPLE_LIMIT); + // Items are folded in id order, so the sample is deterministic. + assert_eq!(bucket.sample_root_item_ids[0], "item-00"); + } } diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index 6c0df9b..dcc6064 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -29,11 +29,11 @@ use tokio::task::JoinHandle; use tracing::{info, warn}; use crate::datagen::{ - datagen_failures, datagen_trajectory, fold_datagen_events, open_stream_events, - DatagenBlobValue, DatagenEvent, DatagenEventType, DatagenFailure, DatagenItemLookup, - DatagenItemStatus, DatagenItemTree, DatagenNewStream, DatagenRootItemStatuses, - DatagenStepCursor, DatagenStepKind, DatagenStreamWriter, DatagenValue, DatagenWriteContext, - FoldedDatagenItem, + datagen_failures, datagen_trajectory, fold_datagen_events, fold_datagen_events_with, + open_stream_events, DatagenBlobProjection, DatagenBlobValue, DatagenEvent, DatagenEventType, + DatagenFailure, DatagenItemLookup, DatagenItemStatus, DatagenItemTree, DatagenNewStream, + DatagenRootItemStatuses, DatagenRunOverview, DatagenStepCursor, DatagenStepKind, + DatagenStreamWriter, DatagenValue, DatagenWriteContext, FoldedDatagenItem, }; use crate::store::{ column_as, column_as_optional, timestamp_from_micros, CompactionConfig, CompactionStats, @@ -217,6 +217,36 @@ impl DatagenStore { } } + /// Reconstruct one item's latest state under an explicit blob projection. + /// + /// `DatagenBlobProjection::Eager` scans the log's blob column too, so the folded blob fields + /// carry their bytes and no `get_blob` round trip is needed. That reads the whole payload — + /// prefer the lazy [`fold_item`](Self::fold_item) unless the caller wants every blob anyway. + pub async fn fold_item_with( + &self, + item_id: &str, + blobs: DatagenBlobProjection, + ) -> LanceResult { + let events = match blobs { + DatagenBlobProjection::Lazy => self.events_for_item(item_id).await?, + DatagenBlobProjection::Eager => { + self.events_with_blobs(&format!("item_id = '{}'", escape_sql_literal(item_id))) + .await? + } + }; + match fold_datagen_events_with(&events, blobs).map_err(invalid_input)? { + Some(item) => Ok(DatagenItemLookup::Found(item)), + None => Ok(DatagenItemLookup::NeverStarted), + } + } + + /// Aggregate the whole log into a run overview: per-status root item counts, failure counts by + /// error type, and completed-step counts. Reads every event once; loads no blob bytes. + pub async fn overview(&self) -> LanceResult { + let events = self.filtered_events("event_type IS NOT NULL").await?; + DatagenRunOverview::build(&events).map_err(invalid_input) + } + /// Classify every root item that shares `root_item_id` with the given roots by folded lifecycle /// status. A root not present in the log is simply absent from the result (never started). pub async fn root_item_statuses( @@ -406,10 +436,31 @@ impl DatagenStore { })) } + /// Like [`filtered_events`](Self::filtered_events) but projects the blob column in, so field + /// blob bytes arrive with the events instead of needing a per-field `get_blob`. + async fn events_with_blobs(&self, filter: &str) -> LanceResult> { + self.scan_events(filter, None).await + } + async fn filtered_events(&self, filter: &str) -> LanceResult> { let columns = self.non_blob_columns(); - let refs: Vec<&str> = columns.iter().map(String::as_str).collect(); - let scanner = self.lsm_scanner().await?.project(&refs).filter(filter)?; + self.scan_events(filter, Some(&columns)).await + } + + /// Scan events matching `filter`, projecting `columns` when given (all columns otherwise), and + /// order them by `(item_id, item_seq, event_id)`. + async fn scan_events( + &self, + filter: &str, + columns: Option<&[String]>, + ) -> LanceResult> { + let scanner = match columns { + Some(columns) => { + let refs: Vec<&str> = columns.iter().map(String::as_str).collect(); + self.lsm_scanner().await?.project(&refs).filter(filter)? + } + None => self.lsm_scanner().await?.filter(filter)?, + }; let mut stream = scanner.try_into_stream().await?; let mut events = Vec::new(); while let Some(batch) = stream.try_next().await? { @@ -1131,6 +1182,89 @@ mod tests { }); } + #[test] + fn eager_fold_materializes_blob_bytes_and_overview_aggregates_the_store() { + let directory = TempDir::new().unwrap(); + let uri = directory.path().to_string_lossy().to_string(); + let blob_bytes = b"eager-screenshot".to_vec(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + let mut store = DatagenStore::open(&uri).await.unwrap(); + store + .append(&[event( + "item-1", + 0, + "created", + 0, + DatagenEventType::ItemCreated, + )]) + .await + .unwrap(); + store + .append_checkpoint(&[ + field_event( + 1, + 0, + DatagenEventType::FieldSet, + "screenshot", + "image", + DatagenValue::Blob(DatagenBlobValue { + bytes: Some(blob_bytes.clone()), + size: blob_bytes.len() as i64, + checksum: None, + }), + ), + completed_step(2, 1), + ]) + .await + .unwrap(); + + // Lazy (the default) leaves the bytes behind; eager projects them into the fold. + let lazy = store.fold_item("item-1").await.unwrap(); + let DatagenFieldState::Set(DatagenValue::Blob(lazy_blob)) = + lazy.folded().unwrap().fields.get("screenshot").unwrap() + else { + panic!("screenshot should be a blob"); + }; + assert!(lazy_blob.bytes.is_none()); + + let eager = store + .fold_item_with("item-1", DatagenBlobProjection::Eager) + .await + .unwrap(); + let DatagenFieldState::Set(DatagenValue::Blob(eager_blob)) = + eager.folded().unwrap().fields.get("screenshot").unwrap() + else { + panic!("screenshot should be a blob"); + }; + assert_eq!(eager_blob.bytes.as_ref(), Some(&blob_bytes)); + + // A second root that fails, so the overview has something in every bucket. + let mut created = event("item-2", 0, "created-2", 0, DatagenEventType::ItemCreated); + created.item_id = "item-2".to_string(); + created.root_item_id = "item-2".to_string(); + store.append(&[created]).await.unwrap(); + let mut failed = completed_step(1, 0); + failed.item_id = "item-2".to_string(); + failed.root_item_id = "item-2".to_string(); + failed.checkpoint_id = "failed-2".to_string(); + failed.event_id = datagen_event_id("item-2", "failed-2", 0); + failed.event_type = DatagenEventType::Failed; + failed.error_type = Some("ValueError".to_string()); + store.append(&[failed]).await.unwrap(); + + let overview = store.overview().await.unwrap(); + assert_eq!(overview.items, 2); + assert_eq!(overview.running, 2); + assert_eq!(overview.failures, 1); + assert_eq!(overview.failures_by_error_type["ValueError"], 1); + assert_eq!(overview.completed_steps["grade"], 1); + let bucket = overview.failures_by_run.values().next().unwrap(); + assert_eq!(bucket.failures, 1); + assert_eq!(bucket.sample_root_item_ids, vec!["item-2".to_string()]); + }); + } + #[test] fn no_op_checkpoint_requires_and_persists_completion_marker() { let directory = TempDir::new().unwrap(); diff --git a/crates/lance-context-server/src/routes/datagen.rs b/crates/lance-context-server/src/routes/datagen.rs index afb00fe..b4ad1bc 100644 --- a/crates/lance-context-server/src/routes/datagen.rs +++ b/crates/lance-context-server/src/routes/datagen.rs @@ -149,18 +149,41 @@ pub async fn add_datagen_events( Ok((StatusCode::CREATED, Json(resp))) } +/// Blob projection for the fold endpoint. +#[derive(Debug, Default, serde::Deserialize)] +pub struct FoldParams { + /// When `true`, blob fields are materialized inline instead of left lazy for `get_blob`. + #[serde(default)] + pub load_blobs: bool, +} + pub async fn fold_datagen_item( State(state): State>, Path((name, item_id)): Path<(String, String)>, + Query(params): Query, ) -> Result, AppError> { let store_lock = state.get_or_open_datagen_store(&name).await?; let store = store_lock.read().await; - let item = DatagenStoreApi::fold_item(&*store, &item_id) + let item = DatagenStoreApi::fold_item_with_blobs(&*store, &item_id, params.load_blobs) .await .map_err(AppError::from_context)?; Ok(Json(GetFoldedDatagenItemResponse { item })) } +/// Aggregate the whole log into a run overview (per-status root counts, failure +/// counts by error type, completed-step counts). +pub async fn datagen_overview( + State(state): State>, + Path(name): Path, +) -> Result, AppError> { + let store_lock = state.get_or_open_datagen_store(&name).await?; + let store = store_lock.read().await; + let overview = DatagenStoreApi::overview(&*store) + .await + .map_err(AppError::from_context)?; + Ok(Json(overview)) +} + pub async fn datagen_item_failures( State(state): State>, Path((name, item_id)): Path<(String, String)>, diff --git a/crates/lance-context-server/src/routes/mod.rs b/crates/lance-context-server/src/routes/mod.rs index b163087..32de356 100644 --- a/crates/lance-context-server/src/routes/mod.rs +++ b/crates/lance-context-server/src/routes/mod.rs @@ -139,6 +139,10 @@ pub fn router() -> Router> { "/api/v1/datagen/{name}/items/{item_id}/failures", get(datagen::datagen_item_failures), ) + .route( + "/api/v1/datagen/{name}/overview", + get(datagen::datagen_overview), + ) .route( "/api/v1/datagen/{name}/root-status", get(datagen::datagen_root_item_statuses), diff --git a/crates/lance-context/src/lib.rs b/crates/lance-context/src/lib.rs index 436ff2e..7f0c39d 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -22,12 +22,13 @@ pub use lance_context_api::{ AddRolloutRequest, AddRolloutsResponse, AddRowsRequest, AddRowsResponse, ColumnSpec, ColumnType, CompactRequest, CompactResponse, CompactStatsResponse, ContextError, ContextResult, ContextStoreApi, CreateDatagenStoreRequest, CreateGenericStoreRequest, - CreateRolloutStoreRequest, DatagenEventDto, DatagenFailureDto, DatagenFieldStateDto, - DatagenRootItemStatusesResponse, DatagenStepCursorDto, DatagenStoreApi, - DatagenStreamPositionDto, DatagenValueDto, DeleteRecordResponse, FoldedDatagenItemDto, - GenericStoreApi, GenericStoreInfo, RecordDto, RelationshipDto, RetrieveRequest, - RetrieveResponse, RetrieveResultDto, RolloutRecordDto, RolloutStoreApi, SchemaSpec, - SearchResultDto, UpsertRecordRequest, UpsertRecordResponse, ID_COLUMN, + CreateRolloutStoreRequest, DatagenEventDto, DatagenFailureBucketDto, DatagenFailureDto, + DatagenFieldStateDto, DatagenRootItemStatusesResponse, DatagenRunOverviewDto, + DatagenStepCursorDto, DatagenStoreApi, DatagenStreamPositionDto, DatagenValueDto, + DeleteRecordResponse, FoldedDatagenItemDto, GenericStoreApi, GenericStoreInfo, RecordDto, + RelationshipDto, RetrieveRequest, RetrieveResponse, RetrieveResultDto, RolloutRecordDto, + RolloutStoreApi, SchemaSpec, SearchResultDto, UpsertRecordRequest, UpsertRecordResponse, + ID_COLUMN, }; #[cfg(feature = "remote")] diff --git a/crates/lance-context/src/unified_datagen.rs b/crates/lance-context/src/unified_datagen.rs index 7677b5f..19533b5 100644 --- a/crates/lance-context/src/unified_datagen.rs +++ b/crates/lance-context/src/unified_datagen.rs @@ -1,6 +1,6 @@ use lance_context_api::{ AddDatagenEventsResponse, ContextError, ContextResult, DatagenEventDto, DatagenFailureDto, - DatagenRootItemStatusesResponse, DatagenStoreApi, FoldedDatagenItemDto, + DatagenRootItemStatusesResponse, DatagenRunOverviewDto, DatagenStoreApi, FoldedDatagenItemDto, }; use lance_context_core::{ datagen_event_to_dto, datagen_events_from_dtos, fold_datagen_events, open_stream_events, @@ -163,6 +163,14 @@ impl DatagenStoreApi for DatagenStore { dispatch_ref!(self, fold_item, item_id) } + async fn fold_item_with_blobs( + &self, + item_id: &str, + load_blobs: bool, + ) -> ContextResult> { + dispatch_ref!(self, fold_item_with_blobs, item_id, load_blobs) + } + async fn root_item_statuses( &self, root_item_ids: &[String], @@ -170,6 +178,10 @@ impl DatagenStoreApi for DatagenStore { dispatch_ref!(self, root_item_statuses, root_item_ids) } + async fn overview(&self) -> ContextResult { + dispatch_ref!(self, overview) + } + async fn item_failures(&self, item_id: &str) -> ContextResult> { dispatch_ref!(self, item_failures, item_id) } diff --git a/python/python/lance_context/api.py b/python/python/lance_context/api.py index 35f1d9f..8ad3fdc 100644 --- a/python/python/lance_context/api.py +++ b/python/python/lance_context/api.py @@ -16,7 +16,7 @@ ContextNamespace as _ContextNamespace, ) from ._internal import ( # pyright: ignore[reportMissingImports] - DatagenItemId as DatagenItemId, + DatagenItemId as _DatagenItemId, ) from ._internal import ( # pyright: ignore[reportMissingImports] DatagenStore as _DatagenStore, @@ -48,7 +48,9 @@ "Context", "ContextNamespace", "GenericStore", + "DatagenItemId", "DatagenStore", + "DatagenStreamWriter", "EmbeddingProvider", "RemoteContext", "RolloutStore", @@ -2830,9 +2832,26 @@ def append(self, events: Iterable[Mapping[str, Any]]) -> int: """Append raw events as one MemWAL generation. Returns the new store version.""" return self._sync.append(list(events)) - def fold_item(self, item_id: str) -> dict[str, Any] | None: - """Fold an item's events into its latest state, or ``None`` if never started.""" - return self._sync.fold_item(item_id) + def fold_item( + self, item_id: str, *, load_blobs: bool = False + ) -> dict[str, Any] | None: + """Fold an item's events into its latest state, or ``None`` if never started. + + ``load_blobs`` materializes blob-field bytes inline; the default leaves them + lazy, to be resolved through :meth:`get_blob` / :meth:`load_blob`. + """ + return self._sync.fold_item(item_id, load_blobs) + + def overview(self) -> dict[str, Any]: + """Aggregate the whole store into a run overview. + + Returns root-item counts by status (``items`` / ``running`` / ``completed`` / + ``filtered``), ``completed_steps`` per step name, and a failure roll-up: total + ``failures``, ``failures_by_error_type``, and ``failures_by_run`` — one bucket + per ``run_id``, each with its own error-type counts and a small sample of + failing root item ids to drill into. + """ + return self._sync.overview() def root_item_statuses(self, root_item_ids: Sequence[str]) -> dict[str, str]: """Classify each root item id by folded lifecycle status. @@ -2908,6 +2927,70 @@ def __repr__(self) -> str: return f"DatagenStore(version={self._sync.version()})" +class DatagenItemId: + """A structured datagen item id. + + Root ids come from the executor's source key; sub-item ids extend a parent with one + fan-out segment (``5/expand:0``). ``str(id)`` is the stored path form. Pure and + client-side — composing an id does no I/O. + """ + + def __init__(self, inner: Any) -> None: + self._inner = inner + + @classmethod + def from_source_key(cls, key: str) -> "DatagenItemId": + """A root id for the executor's source key.""" + return cls(_DatagenItemId.from_source_key(key)) + + @classmethod + def parse(cls, path: str) -> "DatagenItemId": + """Parse a stored path form (``5`` or ``5/expand:0``).""" + return cls(_DatagenItemId.parse(path)) + + def child(self, origin_step: str, branch_idx: int) -> "DatagenItemId": + """The sub-item id forked at ``origin_step`` branch ``branch_idx``.""" + return DatagenItemId(self._inner.child(origin_step, branch_idx)) + + def parent(self) -> "DatagenItemId | None": + """The parent id, or ``None`` for a root.""" + inner = self._inner.parent() + return DatagenItemId(inner) if inner is not None else None + + def root(self) -> "DatagenItemId": + """The root this id descends from (itself when already a root).""" + return DatagenItemId(self._inner.root()) + + @property + def origin_step(self) -> str | None: + """The fan-out step this id was forked at, or ``None`` for a root.""" + return self._inner.origin_step + + @property + def branch_idx(self) -> int | None: + """The branch index within ``origin_step``, or ``None`` for a root.""" + return self._inner.branch_idx + + @property + def is_root(self) -> bool: + """Whether this id has no fan-out segment.""" + return self._inner.is_root + + def __str__(self) -> str: + return str(self._inner) + + def __repr__(self) -> str: + return f"DatagenItemId({str(self._inner)!r})" + + def __eq__(self, other: object) -> bool: + if not isinstance(other, DatagenItemId): + return NotImplemented + return bool(self._inner == other._inner) + + def __hash__(self) -> int: + return hash(str(self._inner)) + + class DatagenStreamWriter: """A per-stream write handle over the native writer. diff --git a/python/src/lib.rs b/python/src/lib.rs index 33d6a2e..8b66312 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -2676,9 +2676,20 @@ impl DatagenStore { } /// Fold an item's events into its latest state, or `None` if never started. - fn fold_item(&self, py: Python<'_>, item_id: &str) -> PyResult> { + /// `load_blobs` materializes blob-field bytes inline; the default leaves them + /// lazy, to be resolved through `get_blob`. + #[pyo3(signature = (item_id, load_blobs = false))] + fn fold_item( + &self, + py: Python<'_>, + item_id: &str, + load_blobs: bool, + ) -> PyResult> { let item = py - .allow_threads(|| self.runtime.block_on(self.store.fold_item(item_id))) + .allow_threads(|| { + self.runtime + .block_on(self.store.fold_item_with_blobs(item_id, load_blobs)) + }) .map_err(to_py_err)?; match item { None => Ok(None), @@ -2686,6 +2697,16 @@ impl DatagenStore { } } + /// Aggregate the whole store into a run overview dict: root-item counts by + /// status, completed-step counts, and a failure roll-up by `run_id` (each + /// bucket carrying a small sample of failing root item ids). + fn overview(&self, py: Python<'_>) -> PyResult { + let overview = py + .allow_threads(|| self.runtime.block_on(self.store.overview())) + .map_err(to_py_err)?; + run_overview_to_py(py, &overview) + } + /// Classify each root item id by folded lifecycle status. Missing ids (never /// started) are absent from the returned dict. fn root_item_statuses(&self, py: Python<'_>, root_item_ids: Vec) -> PyResult { @@ -3410,6 +3431,54 @@ fn item_tree_to_py(py: Python<'_>, tree: &UnifiedDatagenItemTree) -> PyResult, + overview: &lance_context::DatagenRunOverviewDto, +) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("items", overview.items)?; + dict.set_item("running", overview.running)?; + dict.set_item("completed", overview.completed)?; + dict.set_item("filtered", overview.filtered)?; + dict.set_item("failures", overview.failures)?; + dict.set_item( + "failures_by_error_type", + counts_to_py(py, &overview.failures_by_error_type)?, + )?; + dict.set_item( + "completed_steps", + counts_to_py(py, &overview.completed_steps)?, + )?; + let by_run = PyDict::new(py); + for (run_id, bucket) in &overview.failures_by_run { + let entry = PyDict::new(py); + entry.set_item("failures", bucket.failures)?; + entry.set_item( + "failures_by_error_type", + counts_to_py(py, &bucket.failures_by_error_type)?, + )?; + let samples = PyList::empty(py); + for root in &bucket.sample_root_item_ids { + samples.append(root)?; + } + entry.set_item("sample_root_item_ids", samples)?; + by_run.set_item(run_id, entry)?; + } + dict.set_item("failures_by_run", by_run)?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + +fn counts_to_py( + py: Python<'_>, + counts: &std::collections::BTreeMap, +) -> PyResult { + let dict = PyDict::new(py); + for (key, count) in counts { + dict.set_item(key, count)?; + } + Ok(dict.into_pyobject(py)?.unbind().into()) +} + /// A store over a user-declared schema. Rows are plain dicts; the schema is /// declared once at creation and persisted in the dataset. #[pyclass] From 8730e2cf01796f16af24c4edf40b38cb3a1d0699 Mon Sep 17 00:00:00 2001 From: yangjunz Date: Tue, 28 Jul 2026 23:10:09 -0700 Subject: [PATCH 4/5] feat(python): raise typed context-store errors instead of bare RuntimeError `to_py_err` flattened every `ContextError` to `PyRuntimeError`, discarding the distinction the enum already carries. Callers that retry appends could not tell a transport fault from a request the store rejected outright, so a deterministic rejection burned the whole retry budget before failing. Add a `ContextStoreError` hierarchy (rooted at `RuntimeError`, so `except RuntimeError` keeps working) with one leaf per `ContextError` variant, and route the `DatagenStore` methods through it. `InternalError` and `CompactionInProgressError` are the retryable pair; the rest are verdicts. --- python/python/lance_context/__init__.py | 12 +++ python/python/lance_context/api.py | 14 ++++ python/src/lib.rs | 107 ++++++++++++++++++++---- 3 files changed, 117 insertions(+), 16 deletions(-) diff --git a/python/python/lance_context/__init__.py b/python/python/lance_context/__init__.py index ad4b054..bdc1fae 100644 --- a/python/python/lance_context/__init__.py +++ b/python/python/lance_context/__init__.py @@ -1,15 +1,21 @@ from __future__ import annotations from .api import ( # pyright: ignore[reportMissingImports] + AlreadyExistsError, AsyncContext, AsyncRolloutStore, + CompactionInProgressError, Context, ContextNamespace, + ContextStoreError, DatagenItemId, DatagenStore, DatagenStreamWriter, EmbeddingProvider, GenericStore, + InternalError, + InvalidRequestError, + NotFoundError, RemoteContext, RolloutStore, __version__, @@ -21,15 +27,21 @@ ) __all__ = [ + "AlreadyExistsError", "AsyncContext", "AsyncRolloutStore", "Context", + "CompactionInProgressError", "ContextNamespace", + "ContextStoreError", "DatagenItemId", "DatagenStore", "DatagenStreamWriter", "GenericStore", "EmbeddingProvider", + "InternalError", + "InvalidRequestError", + "NotFoundError", "MultiModalEmbeddingProvider", "RemoteContext", "RolloutStore", diff --git a/python/python/lance_context/api.py b/python/python/lance_context/api.py index 8ad3fdc..7e193fe 100644 --- a/python/python/lance_context/api.py +++ b/python/python/lance_context/api.py @@ -11,6 +11,14 @@ if TYPE_CHECKING: from os import PathLike +from ._internal import ( # pyright: ignore[reportMissingImports] + AlreadyExistsError, + CompactionInProgressError, + ContextStoreError, + InternalError, + InvalidRequestError, + NotFoundError, +) from ._internal import Context as _Context # pyright: ignore[reportMissingImports] from ._internal import ( # pyright: ignore[reportMissingImports] ContextNamespace as _ContextNamespace, @@ -43,15 +51,21 @@ from .embeddings import EmbeddingProvider, _build_provider, supports_media __all__ = [ + "AlreadyExistsError", "AsyncContext", "AsyncRolloutStore", + "CompactionInProgressError", "Context", "ContextNamespace", + "ContextStoreError", "GenericStore", "DatagenItemId", "DatagenStore", "DatagenStreamWriter", "EmbeddingProvider", + "InternalError", + "InvalidRequestError", + "NotFoundError", "RemoteContext", "RolloutStore", "__version__", diff --git a/python/src/lib.rs b/python/src/lib.rs index 8b66312..63504d1 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -4,6 +4,7 @@ use std::collections::{BTreeMap, HashMap, HashSet}; use std::sync::Arc; use chrono::{DateTime, SecondsFormat, Utc}; +use pyo3::create_exception; use pyo3::exceptions::{PyRuntimeError, PyTypeError, PyValueError}; use pyo3::prelude::*; use pyo3::types::{PyBytes, PyDict, PyList, PyModule, PyType}; @@ -23,9 +24,10 @@ use lance_context::{ RolloutStore as UnifiedRolloutStore, RolloutStoreApi, SchemaSpec, }; use lance_context_api::{ - AddRecordRequest, CompactRequest, CompactResponse, CompactStatsResponse, ContextStoreApi, - RecordDto, RecordPatchDto, RelationshipDto, RetrieveRequest, RetrieveResultDto, SearchRequest, - SearchResultDto, StateMetadataDto, UpdateRecordRequest, UpsertRecordRequest, + AddRecordRequest, CompactRequest, CompactResponse, CompactStatsResponse, ContextError, + ContextStoreApi, RecordDto, RecordPatchDto, RelationshipDto, RetrieveRequest, + RetrieveResultDto, SearchRequest, SearchResultDto, StateMetadataDto, UpdateRecordRequest, + UpsertRecordRequest, }; use lance_context_client::RemoteContextStore; use lance_context_core::serde::CONTENT_TYPE_TEXT; @@ -2151,6 +2153,64 @@ fn to_py_err(err: E) -> PyErr { PyRuntimeError::new_err(err.to_string()) } +create_exception!( + _internal, + ContextStoreError, + PyRuntimeError, + "Base class for every error a context store raises." +); +create_exception!( + _internal, + NotFoundError, + ContextStoreError, + "The requested store, record, or id does not exist." +); +create_exception!( + _internal, + AlreadyExistsError, + ContextStoreError, + "The store or record being created already exists." +); +create_exception!( + _internal, + InvalidRequestError, + ContextStoreError, + "The request was rejected as malformed. Deterministic — retrying cannot help." +); +create_exception!( + _internal, + InternalError, + ContextStoreError, + "A transport or storage fault. The retryable case: the same call may yet succeed." +); +create_exception!( + _internal, + CompactionInProgressError, + ContextStoreError, + "A compaction holds the store. Retry once it finishes." +); + +/// Map a [`ContextError`] onto the exception class that says whether retrying can help. +/// +/// Every variant subclasses `RuntimeError`, which is what this crate raised before these +/// classes existed, so `except RuntimeError` keeps catching all of them. +/// +/// `ContextError` already separates a transport or storage fault from a request the store +/// refused outright; flattening both to `RuntimeError` forced callers to retry a +/// deterministic rejection until their budget ran out. `InternalError` is the retryable +/// one — the rest are verdicts about the request itself and will fail again identically. +fn ctx_to_py_err(err: ContextError) -> PyErr { + match err { + ContextError::NotFound(msg) => NotFoundError::new_err(msg), + ContextError::AlreadyExists(msg) => AlreadyExistsError::new_err(msg), + ContextError::InvalidRequest(msg) => InvalidRequestError::new_err(msg), + ContextError::Internal(msg) => InternalError::new_err(msg), + ContextError::CompactionInProgress => { + CompactionInProgressError::new_err(ContextError::CompactionInProgress.to_string()) + } + } +} + fn content_to_payloads( content: &Bound<'_, PyAny>, data_type: Option<&str>, @@ -2608,7 +2668,7 @@ impl DatagenStore { let store_res = py.allow_threads(|| { runtime.block_on(UnifiedDatagenStore::open_with_options(uri, storage_options)) }); - let store = store_res.map_err(to_py_err)?; + let store = store_res.map_err(ctx_to_py_err)?; Ok(Self::from_store(store, runtime)) } @@ -2623,7 +2683,7 @@ impl DatagenStore { let runtime = Arc::new(Runtime::new().map_err(to_py_err)?); let store_res = py.allow_threads(|| runtime.block_on(UnifiedDatagenStore::connect(base_url, name))); - let store = store_res.map_err(to_py_err)?; + let store = store_res.map_err(ctx_to_py_err)?; Ok(Self::from_store(store, runtime)) } @@ -2645,7 +2705,7 @@ impl DatagenStore { let store_res = py.allow_threads(|| { runtime.block_on(UnifiedDatagenStore::connect_or_create(base_url, &req)) }); - let store = store_res.map_err(to_py_err)?; + let store = store_res.map_err(ctx_to_py_err)?; Ok(Self::from_store(store, runtime)) } @@ -2661,7 +2721,7 @@ impl DatagenStore { let parsed = events_from_pylist(events)?; let resp = py .allow_threads(|| self.runtime.block_on(self.store.append_checkpoint(&parsed))) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; Ok(resp.version) } @@ -2671,7 +2731,7 @@ impl DatagenStore { let parsed = events_from_pylist(events)?; let resp = py .allow_threads(|| self.runtime.block_on(self.store.append(&parsed))) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; Ok(resp.version) } @@ -2690,7 +2750,7 @@ impl DatagenStore { self.runtime .block_on(self.store.fold_item_with_blobs(item_id, load_blobs)) }) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; match item { None => Ok(None), Some(item) => Ok(Some(folded_item_to_py(py, &item)?)), @@ -2703,7 +2763,7 @@ impl DatagenStore { fn overview(&self, py: Python<'_>) -> PyResult { let overview = py .allow_threads(|| self.runtime.block_on(self.store.overview())) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; run_overview_to_py(py, &overview) } @@ -2715,7 +2775,7 @@ impl DatagenStore { self.runtime .block_on(self.store.root_item_statuses(&root_item_ids)) }) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; let dict = PyDict::new(py); for (item_id, status) in statuses.statuses.iter() { dict.set_item(item_id, status)?; @@ -2727,7 +2787,7 @@ impl DatagenStore { fn item_failures(&self, py: Python<'_>, item_id: &str) -> PyResult { let failures = py .allow_threads(|| self.runtime.block_on(self.store.item_failures(item_id))) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; let list = PyList::empty(py); for failure in &failures { list.append(failure_to_py(py, failure)?)?; @@ -2739,7 +2799,7 @@ impl DatagenStore { fn get_blob(&self, py: Python<'_>, event_id: &str) -> PyResult>> { let bytes = py .allow_threads(|| self.runtime.block_on(self.store.get_blob(event_id))) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; Ok(bytes.map(|b| PyBytes::new(py, &b).unbind())) } @@ -2750,7 +2810,7 @@ impl DatagenStore { fn item_tree(&self, py: Python<'_>, root_item_id: &str) -> PyResult { let tree = py .allow_threads(|| self.runtime.block_on(self.store.item_tree(root_item_id))) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; item_tree_to_py(py, &tree) } @@ -2784,7 +2844,7 @@ impl DatagenStore { self.runtime .block_on(self.store.open_stream(&stream, &context)) }) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; Ok(DatagenStreamWriter { inner: writer }) } @@ -2806,7 +2866,7 @@ impl DatagenStore { self.runtime .block_on(self.store.resume_stream(item_id, &context)) }) - .map_err(to_py_err)?; + .map_err(ctx_to_py_err)?; Ok(writer.map(|inner| DatagenStreamWriter { inner })) } } @@ -3737,5 +3797,20 @@ fn _internal(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add("ContextStoreError", m.py().get_type::())?; + m.add("NotFoundError", m.py().get_type::())?; + m.add( + "AlreadyExistsError", + m.py().get_type::(), + )?; + m.add( + "InvalidRequestError", + m.py().get_type::(), + )?; + m.add("InternalError", m.py().get_type::())?; + m.add( + "CompactionInProgressError", + m.py().get_type::(), + )?; Ok(()) } From c93a2c2470f8c6cf48ccecf1e85029420fe994c2 Mon Sep 17 00:00:00 2001 From: yangjunz Date: Tue, 28 Jul 2026 23:17:59 -0700 Subject: [PATCH 5/5] fix(python): bind the error classes through aliases so pyright resolves them CI type-checks without a built `_internal`, so every name imported from it needs the per-name `pyright: ignore` the rest of this module already uses. Re-exporting the error classes directly leaked the unresolved symbols into `__init__`. --- python/python/lance_context/api.py | 34 ++++++++++++++++++++++++------ 1 file changed, 28 insertions(+), 6 deletions(-) diff --git a/python/python/lance_context/api.py b/python/python/lance_context/api.py index 7e193fe..744431a 100644 --- a/python/python/lance_context/api.py +++ b/python/python/lance_context/api.py @@ -12,17 +12,18 @@ from os import PathLike from ._internal import ( # pyright: ignore[reportMissingImports] - AlreadyExistsError, - CompactionInProgressError, - ContextStoreError, - InternalError, - InvalidRequestError, - NotFoundError, + AlreadyExistsError as _AlreadyExistsError, +) +from ._internal import ( # pyright: ignore[reportMissingImports] + CompactionInProgressError as _CompactionInProgressError, ) from ._internal import Context as _Context # pyright: ignore[reportMissingImports] from ._internal import ( # pyright: ignore[reportMissingImports] ContextNamespace as _ContextNamespace, ) +from ._internal import ( # pyright: ignore[reportMissingImports] + ContextStoreError as _ContextStoreError, +) from ._internal import ( # pyright: ignore[reportMissingImports] DatagenItemId as _DatagenItemId, ) @@ -35,6 +36,15 @@ from ._internal import ( # pyright: ignore[reportMissingImports] GenericStore as _GenericStore, ) +from ._internal import ( # pyright: ignore[reportMissingImports] + InternalError as _InternalError, +) +from ._internal import ( # pyright: ignore[reportMissingImports] + InvalidRequestError as _InvalidRequestError, +) +from ._internal import ( # pyright: ignore[reportMissingImports] + NotFoundError as _NotFoundError, +) from ._internal import ( # pyright: ignore[reportMissingImports] RemoteContext as _RemoteContext, ) @@ -75,6 +85,18 @@ __version__ = _version() +# The store error hierarchy, re-exported from the native module. `ContextStoreError` is +# rooted at `RuntimeError` — what this package raised before these classes existed — so +# `except RuntimeError` still catches every one of them. `InternalError` and +# `CompactionInProgressError` are the retryable pair; the rest are verdicts on the +# request itself, which an identical retry can only earn again. +ContextStoreError = _ContextStoreError +NotFoundError = _NotFoundError +AlreadyExistsError = _AlreadyExistsError +InvalidRequestError = _InvalidRequestError +InternalError = _InternalError +CompactionInProgressError = _CompactionInProgressError + def generate_id() -> str: """Generate a new time-ordered id (UUIDv7) as a string.