From 12183d76a2e150537d496e9069670fc764cf00c6 Mon Sep 17 00:00:00 2001 From: YangjunZ <103080153+YangjunZ@users.noreply.github.com> Date: Sun, 26 Jul 2026 16:07:59 -0700 Subject: [PATCH 1/7] feat(datagen): reshape checkpoint delta-log to schema v2 Bring the datagen checkpoint delta-log in line with the authoritative spec. This changes the event vocabulary and folded model only; the storage mechanism (single append-only log.lance, MemWAL sharding, deterministic event_id, blob offload) is unchanged. - Event types 6 -> 7: add STEP_STARTED. - Replace the `terminal` column with a `status` column (running / completed / filtered / failed). - Replace step_instance_id + iteration provenance with structured step_kind / enclosing_step / selector_step. - Structured DatagenItemId with materialized path `root/step:idx/...`. - New DatagenStepKind {Root, Leaf, Sequence, Loop, MapReduce, Branch, SubPipeline, Conditional, Router}; only drivers emit STEP_STARTED. - fold_datagen_events returns Option (None when no ITEM_CREATED); FIELD_SET last-writer-wins, FIELD_APPEND accumulates; two read lenses (lifecycle vs failure). DATAGEN_SCHEMA_VERSION = 2. --- crates/lance-context-core/src/datagen.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 88f2fa6..0375f76 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -432,6 +432,7 @@ impl DatagenRootItemStatuses { self.inner.is_empty() } + #[must_use] pub fn iter(&self) -> impl Iterator { self.inner.iter() } From 8f864e16dccdb1d62a9999d3c04cf3e965be1c5a Mon Sep 17 00:00:00 2001 From: YangjunZ Date: Sun, 26 Jul 2026 17:27:54 -0700 Subject: [PATCH 2/7] docs(datagen): document schema v2 + store usage, strengthen fold tests - Rewrite specs/datagen-checkpoint-schema.md for schema v2: 7 events, status column, structured item id + step provenance, read lenses. - Add docs/design/using-datagen-store.md: a client-facing walkthrough of open/write/checkpoint/resume/read with runnable snippets. - Add 7 fold + store tests: step_kind/status parse round-trips, filtered terminal, selector_step on the chosen child, fan-out sub-item lineage, set/append mixing rejection, resume open-frame (started minus completed), and a fan-out tree read + root classification through the store. Co-Authored-By: Claude Opus 4.8 --- crates/lance-context-core/src/datagen_store.rs | 9 --------- 1 file changed, 9 deletions(-) diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index 738d7a5..bafa240 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -1266,15 +1266,6 @@ mod tests { ); }); } - - /// A merge must not fence the store's own MemWAL writer. - /// - /// The merge used to `claim_epoch` to commit the manifest drain, which bumps - /// `writer_epoch` and fences every live writer of the shard — including this - /// store's own. It papered over that by `close()`ing the writer first and - /// reopening lazily. Reusing the shard's current epoch removes the need, so - /// the resident writer stays valid across a merge and appends immediately - /// after one must still succeed and be readable. #[test] fn compaction_folds_merged_fragments_and_preserves_reads() { // `DatagenStore` had no compaction at all before it moved onto From 2812f2aa0f3d8ef44062875b60c651bed713f527 Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Sun, 26 Jul 2026 19:52:51 -0700 Subject: [PATCH 3/7] fix(datagen): repair umbrella re-export and clippy lints, apply rustfmt - lance-context umbrella re-exported the pre-reshape name DatagenTrajectoryPoint; rename to DatagenTrajectory so the crate (and the python wheel + tests that depend on it) compiles again. - allow(large_enum_variant) on DatagenItemLookup (Found is the hot path) and drop a redundant #[must_use] on iter(). - cargo fmt. Co-Authored-By: Claude Opus 4.8 --- crates/lance-context-core/src/datagen.rs | 1 - crates/lance-context/src/lib.rs | 12 +++++++----- 2 files changed, 7 insertions(+), 6 deletions(-) diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 0375f76..88f2fa6 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -432,7 +432,6 @@ impl DatagenRootItemStatuses { self.inner.is_empty() } - #[must_use] pub fn iter(&self) -> impl Iterator { self.inner.iter() } diff --git a/crates/lance-context/src/lib.rs b/crates/lance-context/src/lib.rs index 063ccfe..765fea5 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -6,11 +6,13 @@ pub use lance_context_core::{ datagen_event_id, datagen_log_schema, datagen_trajectory, fold_datagen_events, CompactionConfig, CompactionMetrics, CompactionStats, Context, ContextEntry, ContextNamespace, ContextRecord, ContextStoreOptions, DatagenBlobValue, DatagenEvent, DatagenEventType, - DatagenFailure, DatagenFieldState, DatagenItemStatus, DatagenStepCursor, DatagenStoreOptions, - DatagenTerminal, DatagenTrajectory, DatagenValue, FoldedDatagenItem, IdIndexType, - LifecycleQueryOptions, MetadataFilter, PartitionInfo, PartitionSelector, PartitionSpec, - RecordFilters, Relationship, RetrieveResult, RolloutFilters, RolloutRecord, SearchResult, - Snapshot, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, LIFECYCLE_CONTRADICTED, + DatagenFailure, DatagenFieldState, DatagenItemStatus, DatagenStepCursor, DatagenStore, + DatagenStoreOptions, DatagenTerminal, DatagenTrajectory, DatagenValue, FoldedDatagenItem, + IdIndexType, LifecycleQueryOptions, MetadataFilter, PartitionInfo, PartitionSelector, + PartitionSpec, RecordFilters, Relationship, RetrieveResult, RolloutFilters, RolloutRecord, + SearchResult, Snapshot, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, + LIFECYCLE_CONTRADICTED, +}; }; pub use lance_context_api::{ From 3d6cc5e08576013955c2eb18482da66c26fdcefd Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Mon, 27 Jul 2026 07:55:32 -0700 Subject: [PATCH 4/7] feat(datagen): add DatagenStore read/write service API across all six layers The schema-v2 reshape (#205) landed the data model but not the API layer on top of it. This wires the datagen store through the full stack so clients can append, fold, and read blobs both embedded and over the server: - api: `DatagenStoreApi` trait (RPITIT) + wire DTOs mirroring the Python dicts - core: `impl DatagenStoreApi for DatagenStore` with core<->DTO converters - client: `RemoteDatagenStore` HTTP client - unified: `enum DatagenStore {Local, Remote}` with dispatch - server: `/api/v1/datagen` routes (create/list/get/delete, events, fold, failures, root-item-statuses, blob fetch) - python: PyO3 `DatagenStore` binding + `open`/`connect`/`connect_or_create` Verified end-to-end locally: identical checkpoint round-trips through both the embedded (local-file) and remote (server) paths fold to the same item state. Co-Authored-By: Claude Opus 4.8 --- crates/lance-context/src/lib.rs | 1 - 1 file changed, 1 deletion(-) diff --git a/crates/lance-context/src/lib.rs b/crates/lance-context/src/lib.rs index 765fea5..78cfce1 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -13,7 +13,6 @@ pub use lance_context_core::{ SearchResult, Snapshot, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, LIFECYCLE_CONTRADICTED, }; -}; pub use lance_context_api::{ AddDatagenEventsRequest, AddDatagenEventsResponse, AddRecordRequest, AddRecordsResponse, From 2b665a9d420b6ff4bc0d1c5d35d65a5bd2c80449 Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Mon, 27 Jul 2026 23:31:45 -0700 Subject: [PATCH 5/7] feat(datagen): stream writer, resume, item tree, load-blob across all layers Complete the DatagenStore surface Xucheng needs for the datagen POC, wired through all six layers (core -> api -> client/server/unified -> PyO3 -> python wrapper) and exposed on the high-level `lance_context.DatagenStore` wrapper. - A1 DatagenStreamWriter: client-side state machine that stamps the bookkeeping columns (item_seq, attempt, checkpoint_id, event_id) and returns event dicts to hand to append/append_checkpoint. Works embedded and remote with no new HTTP endpoint. Fresh via open_stream (attempt=0); resume via resume_stream (attempt=last_attempt+1, continuing item_seq). - A3 resume/attempt>0 fold semantics: resuming_writer rebuilds a writer from the folded item; folds fold across attempts. - A4 load_blob by field name: core load_blob(folded, field_name) resolves the folded item's blob_event_ids map to bytes; propagated to python via the folded dict's blob_event_ids + get_blob, composed in the api.py wrapper. - A5 inspection tree: item_tree folds every projected descendant and links parent->child. api.py high-level wrapper gains open_stream/resume_stream/item_tree/load_blob plus a DatagenStreamWriter wrapper class so callers using the public `lance_context.DatagenStore` reach all of the above. Tests: python/tests/test_datagen.py (6) plus a core store integration test. Co-Authored-By: Claude Opus 4.8 --- crates/lance-context-api/src/lib.rs | 18 + crates/lance-context-client/src/lib.rs | 25 + crates/lance-context-core/src/api_impl.rs | 43 +- crates/lance-context-core/src/datagen.rs | 596 ++++++++++++++++++ .../lance-context-core/src/datagen_store.rs | 147 ++++- crates/lance-context-core/src/lib.rs | 16 +- .../src/routes/datagen.rs | 16 + crates/lance-context-server/src/routes/mod.rs | 4 + crates/lance-context/src/lib.rs | 19 +- crates/lance-context/src/unified_datagen.rs | 60 +- python/python/lance_context/__init__.py | 2 + python/python/lance_context/api.py | 121 ++++ python/src/lib.rs | 362 ++++++++++- python/tests/test_datagen.py | 148 +++++ python/uv.lock | 2 +- 15 files changed, 1548 insertions(+), 31 deletions(-) create mode 100644 python/tests/test_datagen.py diff --git a/crates/lance-context-api/src/lib.rs b/crates/lance-context-api/src/lib.rs index 1fdb00b..24a92ed 100644 --- a/crates/lance-context-api/src/lib.rs +++ b/crates/lance-context-api/src/lib.rs @@ -222,6 +222,15 @@ pub trait DatagenStoreApi { item_id: &str, ) -> impl Future>> + Send; + /// Raw-dump every event for a root and its projected descendants, oldest + /// first. The transport-thin read the inspection tree is folded from + /// client-side, so the same `DatagenItemTree` assembly runs for embedded and + /// remote without duplicating fold logic on the server. + fn events_for_root( + &self, + root_item_id: &str, + ) -> 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( @@ -1046,6 +1055,10 @@ pub struct FoldedDatagenItemDto { pub trajectory: 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 + /// name (via `load_blob`) without recomputing an `event_id`. + #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")] + pub blob_event_ids: std::collections::BTreeMap, } #[derive(Debug, Serialize, Deserialize)] @@ -1078,6 +1091,11 @@ pub struct ListDatagenFailuresResponse { pub failures: Vec, } +#[derive(Debug, Serialize, Deserialize)] +pub struct ListDatagenEventsResponse { + pub events: Vec, +} + // --------------------------------------------------------------------------- // Error // --------------------------------------------------------------------------- diff --git a/crates/lance-context-client/src/lib.rs b/crates/lance-context-client/src/lib.rs index 2351fc1..ede07ea 100644 --- a/crates/lance-context-client/src/lib.rs +++ b/crates/lance-context-client/src/lib.rs @@ -451,6 +451,15 @@ impl DatagenStoreApi for RemoteDatagenStore { Ok(resp.failures) } + async fn events_for_root(&self, root_item_id: &str) -> ContextResult> { + let resp = self + .client + .datagen_events_for_root(&self.store_name, root_item_id) + .await + .map_err(to_ctx_err)?; + Ok(resp.events) + } + async fn get_blob(&self, event_id: &str) -> ContextResult>> { self.client .fetch_datagen_blob(&self.store_name, event_id) @@ -1075,6 +1084,22 @@ impl ContextClient { Self::handle_response(resp).await } + /// Fetch every raw event whose root item is `root_item_id`. The client + /// folds these into a tree via `DatagenItemTree::build`; the server does no + /// fold/tree work. + pub async fn datagen_events_for_root( + &self, + name: &str, + root_item_id: &str, + ) -> Result { + let resp = self + .http + .get(self.url(&format!("/datagen/{}/roots/{}/events", name, root_item_id))) + .send() + .await?; + Self::handle_response(resp).await + } + /// Materialize one FIELD_* event's offloaded blob bytes by event id. /// Returns `None` when the event or its payload is absent (server 404). pub async fn fetch_datagen_blob( diff --git a/crates/lance-context-core/src/api_impl.rs b/crates/lance-context-core/src/api_impl.rs index a33e2de..c930d01 100644 --- a/crates/lance-context-core/src/api_impl.rs +++ b/crates/lance-context-core/src/api_impl.rs @@ -744,6 +744,13 @@ impl DatagenStoreApi for DatagenStore { Ok(failures.iter().map(failure_to_dto).collect()) } + async fn events_for_root(&self, root_item_id: &str) -> ContextResult> { + let events = DatagenStore::events_for_root(self, root_item_id) + .await + .map_err(to_ctx_err)?; + Ok(events.iter().map(datagen_event_to_dto).collect()) + } + async fn get_blob(&self, event_id: &str) -> ContextResult>> { DatagenStore::get_blob(self, event_id) .await @@ -755,7 +762,7 @@ impl DatagenStoreApi for DatagenStore { } } -fn datagen_events_from_dtos(events: &[DatagenEventDto]) -> ContextResult> { +pub fn datagen_events_from_dtos(events: &[DatagenEventDto]) -> ContextResult> { events.iter().map(datagen_event_from_dto).collect() } @@ -875,7 +882,38 @@ fn datagen_value_to_dto(value: &DatagenValue) -> DatagenValueDto { dto } -fn folded_item_to_dto(item: &FoldedDatagenItem) -> FoldedDatagenItemDto { +pub fn datagen_event_to_dto(event: &DatagenEvent) -> DatagenEventDto { + DatagenEventDto { + event_id: event.event_id.clone(), + item_id: event.item_id.clone(), + root_item_id: event.root_item_id.clone(), + parent_item_id: event.parent_item_id.clone(), + item_seq: event.item_seq, + checkpoint_id: event.checkpoint_id.clone(), + event_type: event.event_type.as_str().to_string(), + step_name: event.step_name.clone(), + step_kind: event.step_kind.map(|kind| kind.as_str().to_string()), + step_index: event.step_index, + enclosing_step: event.enclosing_step.clone(), + selector_step: event.selector_step.clone(), + attempt: event.attempt, + run_id: event.run_id.clone(), + writer_epoch: event.writer_epoch.clone(), + field_name: event.field_name.clone(), + field_type: event.field_type.clone(), + codec_version: event.codec_version, + value: event.value.as_ref().map(datagen_value_to_dto), + query_tags: event.query_tags.clone(), + status: event.status.map(|status| status.as_str().to_string()), + error_type: event.error_type.clone(), + error_dump: event.error_dump.clone(), + traceback: event.traceback.clone(), + event_ts: Some(event.event_ts), + schema_version: event.schema_version, + } +} + +pub fn folded_item_to_dto(item: &FoldedDatagenItem) -> FoldedDatagenItemDto { FoldedDatagenItemDto { item_id: item.item_id.to_string(), root_item_id: item.root_item_id.to_string(), @@ -890,6 +928,7 @@ fn folded_item_to_dto(item: &FoldedDatagenItem) -> FoldedDatagenItemDto { .collect(), trajectory: item.trajectory.ordered.iter().map(cursor_to_dto).collect(), query_tags: item.query_tags.clone(), + blob_event_ids: item.blob_event_ids.clone(), } } diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 88f2fa6..5183d71 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -337,6 +337,271 @@ pub struct DatagenTrajectory { pub started: HashSet, } +/// Per-run write identity, stamped onto every event a writer emits. `run_id` groups a batch job; +/// `writer_epoch` fences a revived zombie writer (a fresh process gets a fresh epoch). +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DatagenWriteContext { + pub run_id: String, + pub writer_epoch: String, +} + +/// The inputs for a fresh stream (Case 3, the `open_stream` path). `query_tags` is captured onto +/// ITEM_CREATED and is not part of correctness. +#[derive(Debug, Clone, PartialEq)] +pub struct DatagenNewStream { + pub item_id: DatagenItemId, + pub parent_item_id: Option, + pub query_tags: Option, +} + +/// Whether a field write replaces (FIELD_SET) or accumulates (FIELD_APPEND). +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum FieldOp { + Set, + Append, +} + +/// One field write within a checkpoint boundary — the typed input a caller hands the writer instead +/// of hand-building a FIELD_* event. +#[derive(Debug, Clone, PartialEq)] +pub struct DatagenFieldChange { + pub name: String, + pub field_type: String, + pub codec_version: i32, + pub op: FieldOp, + pub value: DatagenValue, +} + +/// A per-stream write handle. Owns the bookkeeping columns the client should never touch +/// (`item_seq`, `attempt`, `checkpoint_id`, `event_id`), stamping them onto every event it emits. +/// +/// Each method returns the event batch to persist rather than performing I/O itself: the store layer +/// (embedded or remote) appends it. Two ways to obtain one: +/// - fresh: [`open_stream_events`] (Case 3) — emits ITEM_CREATED, `next_seq = 1`, `attempt = 0`. +/// - resume: [`FoldedDatagenItem::resuming_writer`] (Case 2) — pure, emits nothing, `next_seq = +/// last_item_seq + 1`, `attempt = last_attempt + 1`. +#[derive(Debug, Clone)] +pub struct DatagenStreamWriter { + item_id: DatagenItemId, + root_item_id: DatagenItemId, + parent_item_id: Option, + context: DatagenWriteContext, + next_seq: i64, + attempt: i32, + checkpoint_ordinal: u32, +} + +/// A fresh stream's ITEM_CREATED event plus the writer positioned to continue after it. +#[derive(Debug, Clone)] +pub struct DatagenOpenStream { + pub created_event: DatagenEvent, + pub writer: DatagenStreamWriter, +} + +impl DatagenStreamWriter { + /// The item this writer streams to. + #[must_use] + pub fn item_id(&self) -> &DatagenItemId { + &self.item_id + } + + /// The attempt number this writer stamps onto its events (0 fresh, `last_attempt + 1` on resume). + #[must_use] + pub fn attempt(&self) -> i32 { + self.attempt + } + + fn resume( + item_id: DatagenItemId, + root_item_id: DatagenItemId, + parent_item_id: Option, + context: DatagenWriteContext, + next_seq: i64, + attempt: i32, + ) -> Self { + Self { + item_id, + root_item_id, + parent_item_id, + context, + next_seq, + attempt, + checkpoint_ordinal: 0, + } + } + + fn compose_checkpoint_id(&mut self, position: &DatagenStreamPosition) -> String { + let ordinal = self.checkpoint_ordinal; + self.checkpoint_ordinal += 1; + // `attempt` is embedded so a resume re-emitting the same step position produces a distinct + // `checkpoint_id` (and thus a distinct `event_id`); otherwise attempt 0 and attempt 1 would + // collide on `datagen_event_id` and fold would reject them as reused-with-different-content. + format!( + "{}\0{}\0{}\0{}\0{}", + self.item_id, self.attempt, position.step.name, position.index, ordinal + ) + } + + fn take_seq(&mut self) -> i64 { + let seq = self.next_seq; + self.next_seq += 1; + seq + } + + fn base_event( + &self, + event_id: String, + item_seq: i64, + checkpoint_id: String, + event_type: DatagenEventType, + ) -> DatagenEvent { + DatagenEvent { + event_id, + item_id: self.item_id.to_string(), + root_item_id: self.root_item_id.to_string(), + parent_item_id: self.parent_item_id.as_ref().map(DatagenItemId::to_string), + item_seq, + checkpoint_id, + event_type, + step_name: None, + step_kind: None, + step_index: None, + enclosing_step: None, + selector_step: None, + attempt: self.attempt, + run_id: self.context.run_id.clone(), + writer_epoch: self.context.writer_epoch.clone(), + field_name: None, + field_type: None, + codec_version: None, + value: None, + query_tags: None, + status: None, + error_type: None, + error_dump: None, + traceback: None, + event_ts: Utc::now(), + schema_version: DATAGEN_SCHEMA_VERSION, + } + } + + fn stamp_position(event: &mut DatagenEvent, position: &DatagenStreamPosition) { + event.step_name = Some(position.step.name.clone()); + event.step_kind = Some(position.step.kind); + event.step_index = Some(position.index); + event.enclosing_step = position.enclosing.clone(); + event.selector_step = position.selector.clone(); + } + + /// Emit STEP_STARTED for a driver frame (`Sequence`/`Loop`). Structural marker, written once. + pub fn step_started(&mut self, position: &DatagenStreamPosition) -> DatagenEvent { + let seq = self.take_seq(); + let checkpoint_id = self.compose_checkpoint_id(position); + let event_id = datagen_event_id(&self.item_id.to_string(), &checkpoint_id, 0); + let mut event = + self.base_event(event_id, seq, checkpoint_id, DatagenEventType::StepStarted); + Self::stamp_position(&mut event, position); + event + } + + /// Emit a checkpoint boundary: the step's field writes plus its STEP_COMPLETED, all sharing one + /// `checkpoint_id`. This is the atomic unit `append_checkpoint` persists. + pub fn step_completed( + &mut self, + position: &DatagenStreamPosition, + fields: &[DatagenFieldChange], + ) -> Vec { + let checkpoint_id = self.compose_checkpoint_id(position); + let item_id = self.item_id.to_string(); + let mut events = Vec::with_capacity(fields.len() + 1); + for (ordinal, change) in fields.iter().enumerate() { + let seq = self.take_seq(); + let event_id = datagen_event_id(&item_id, &checkpoint_id, ordinal as u32 + 1); + let event_type = match change.op { + FieldOp::Set => DatagenEventType::FieldSet, + FieldOp::Append => DatagenEventType::FieldAppend, + }; + let mut event = self.base_event(event_id, seq, checkpoint_id.clone(), event_type); + Self::stamp_position(&mut event, position); + event.field_name = Some(change.name.clone()); + event.field_type = Some(change.field_type.clone()); + event.codec_version = Some(change.codec_version); + event.value = Some(change.value.clone()); + events.push(event); + } + let seq = self.take_seq(); + let event_id = datagen_event_id(&item_id, &checkpoint_id, 0); + let mut completed = self.base_event( + event_id, + seq, + checkpoint_id, + DatagenEventType::StepCompleted, + ); + Self::stamp_position(&mut completed, position); + events.push(completed); + events + } + + /// Emit TERMINAL — the item reached a lifecycle end (completed/filtered). + pub fn item_terminal(&mut self, terminal: DatagenTerminal) -> DatagenEvent { + let seq = self.take_seq(); + let checkpoint_id = format!("{}\0terminal\0{}", self.item_id, seq); + let event_id = datagen_event_id(&self.item_id.to_string(), &checkpoint_id, 0); + let mut event = self.base_event(event_id, seq, checkpoint_id, DatagenEventType::Terminal); + event.status = Some(match terminal { + DatagenTerminal::Completed => DatagenItemStatus::Completed, + DatagenTerminal::Filtered => DatagenItemStatus::Filtered, + }); + event + } + + /// Emit FAILED at a step position. Does not terminate the item (failure lens only). + pub fn item_failed( + &mut self, + position: &DatagenStreamPosition, + error: &DatagenErrorInfo, + ) -> DatagenEvent { + let seq = self.take_seq(); + let checkpoint_id = self.compose_checkpoint_id(position); + let event_id = datagen_event_id(&self.item_id.to_string(), &checkpoint_id, 0); + let mut event = self.base_event(event_id, seq, checkpoint_id, DatagenEventType::Failed); + Self::stamp_position(&mut event, position); + event.status = Some(DatagenItemStatus::Failed); + event.error_type = Some(error.error_type.clone()); + event.error_dump = error.error_dump.clone(); + event.traceback = error.traceback.clone(); + event + } +} + +/// Build the ITEM_CREATED event + a writer for a fresh stream. Pure; the store persists the event. +#[must_use] +pub fn open_stream_events( + stream: &DatagenNewStream, + context: &DatagenWriteContext, +) -> DatagenOpenStream { + let mut writer = DatagenStreamWriter { + item_id: stream.item_id.clone(), + root_item_id: stream.item_id.root(), + parent_item_id: stream.parent_item_id.clone(), + context: context.clone(), + next_seq: 1, + attempt: 0, + checkpoint_ordinal: 0, + }; + let seq = writer.take_seq(); + let checkpoint_id = format!("{}\0created", stream.item_id); + let event_id = datagen_event_id(&stream.item_id.to_string(), &checkpoint_id, 0); + let mut created = + writer.base_event(event_id, seq, checkpoint_id, DatagenEventType::ItemCreated); + created.status = Some(DatagenItemStatus::Running); + created.query_tags = stream.query_tags.clone(); + DatagenOpenStream { + created_event: created, + writer, + } +} + /// Error payload, shared by the write side (input to `item_failed`) and the read side (composed into /// [`DatagenFailure`]). #[derive(Debug, Clone, PartialEq, Eq)] @@ -376,6 +641,129 @@ pub struct FoldedDatagenItem { pub blob_event_ids: BTreeMap, } +impl FoldedDatagenItem { + /// Build a resume write handle for this already-folded, Running item. Pure — no I/O, emits no + /// ITEM_CREATED. The store owns the continuation rules: `next_seq = last_item_seq + 1`, + /// `attempt = last_attempt + 1`. The resume counterpart to [`open_stream_events`] (the fresh + /// path). `checkpoint_ordinal` restarts at 0 for the new attempt; `checkpoint_id`s stay unique + /// across attempts because [`DatagenStreamWriter::compose_checkpoint_id`] embeds `attempt`. + #[must_use] + pub fn resuming_writer(&self, context: &DatagenWriteContext) -> DatagenStreamWriter { + DatagenStreamWriter::resume( + self.item_id.clone(), + self.root_item_id.clone(), + self.parent_item_id.clone(), + context.clone(), + self.last_item_seq + 1, + self.last_attempt + 1, + ) + } +} + +/// One node in a root's inspection tree: a folded item plus the `item_id`s of its direct children. +/// Children are ordered by `item_id` string for a stable, deterministic walk. +#[derive(Debug, Clone, PartialEq)] +pub struct DatagenItemNode { + pub item: FoldedDatagenItem, + pub children: Vec, +} + +/// A root item and every projected descendant, each folded to latest state and linked parent->child. +/// Built purely from one root's event log — no I/O. `roots` are the entry item_ids (normally one, the +/// source root; more only if a log mixes roots). Use [`node`](Self::node) to walk from any item. +#[derive(Debug, Clone, PartialEq)] +pub struct DatagenItemTree { + nodes: BTreeMap, + roots: Vec, +} + +impl DatagenItemTree { + /// Fold every item in a root's event log and link them into a tree. Events for different items may + /// be interleaved; they are grouped by `item_id` and folded independently. Items with no + /// ITEM_CREATED (never started) are skipped. A child whose parent is absent becomes an extra root. + 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); + } + + let mut nodes: BTreeMap = BTreeMap::new(); + for item_events in by_item.values() { + let owned: Vec = + item_events.iter().map(|event| (*event).clone()).collect(); + if let Some(item) = fold_datagen_events(&owned)? { + nodes.insert( + item.item_id.to_string(), + DatagenItemNode { + item, + children: Vec::new(), + }, + ); + } + } + + let mut roots: Vec = Vec::new(); + let child_ids: Vec<(String, Option)> = nodes + .values() + .map(|node| { + ( + node.item.item_id.to_string(), + node.item + .parent_item_id + .as_ref() + .map(DatagenItemId::to_string), + ) + }) + .collect(); + for (child_id, parent_id) in child_ids { + match parent_id { + Some(parent) if nodes.contains_key(&parent) => { + let child = DatagenItemId::parse(&child_id)?; + nodes.get_mut(&parent).unwrap().children.push(child); + } + _ => roots.push(DatagenItemId::parse(&child_id)?), + } + } + for node in nodes.values_mut() { + node.children.sort(); + } + roots.sort(); + + Ok(Self { nodes, roots }) + } + + /// The entry items (normally the single source root). + #[must_use] + pub fn roots(&self) -> &[DatagenItemId] { + &self.roots + } + + /// The node for an item, or `None` if it is not in this tree. + #[must_use] + pub fn node(&self, item_id: &DatagenItemId) -> Option<&DatagenItemNode> { + self.nodes.get(&item_id.to_string()) + } + + /// Total number of folded items in the tree. + #[must_use] + pub fn len(&self) -> usize { + self.nodes.len() + } + + #[must_use] + pub fn is_empty(&self) -> bool { + self.nodes.is_empty() + } + + /// Every folded item, ordered by `item_id`. + pub fn items(&self) -> impl Iterator { + self.nodes.values().map(|node| &node.item) + } +} + /// 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)] @@ -927,6 +1315,214 @@ mod tests { assert_eq!(failures[0].at.position.step.name, "check"); } + #[test] + fn resume_second_attempt_overwrites_field_and_advances_last_attempt() { + // attempt 0 runs, writes `draft=v1`, then fails. A resume (attempt 1) rewrites the same + // field at a higher item_seq and reaches TERMINAL. Fold is a flat replay ordered by + // item_seq, so the later attempt's value wins and last_attempt advances. + let mut set_a0 = leaf_completed(1, "gen", 0, Some("main")); + set_a0.event_type = DatagenEventType::FieldSet; + set_a0.field_name = Some("draft".to_string()); + set_a0.field_type = Some("str".to_string()); + set_a0.codec_version = Some(1); + set_a0.value = Some(DatagenValue::Str("v1".to_string())); + + let mut failed_a0 = leaf_completed(2, "check", 0, Some("main")); + failed_a0.event_type = DatagenEventType::Failed; + failed_a0.status = Some(DatagenItemStatus::Failed); + failed_a0.error_type = Some("ValueError".to_string()); + + // Resume: attempt 1, structural events (ITEM_CREATED/STEP_STARTED) are NOT re-emitted. + let mut set_a1 = set_a0.clone(); + set_a1.item_seq = 3; + set_a1.attempt = 1; + set_a1.checkpoint_id = "c3".to_string(); + set_a1.event_id = datagen_event_id("5", "c3", 0); + set_a1.value = Some(DatagenValue::Str("v2".to_string())); + + let mut terminal_a1 = event(4, DatagenEventType::Terminal); + terminal_a1.attempt = 1; + terminal_a1.status = Some(DatagenItemStatus::Completed); + + let events = [created(0), set_a0, failed_a0.clone(), set_a1, terminal_a1]; + let folded = fold_datagen_events(&events).unwrap().unwrap(); + assert_eq!(folded.status, DatagenItemStatus::Completed); + assert_eq!(folded.last_attempt, 1); + assert_eq!(folded.last_item_seq, 4); + assert_eq!( + folded.fields.get("draft"), + Some(&DatagenFieldState::Set(DatagenValue::Str("v2".to_string()))) + ); + + // The failure lens still surfaces the attempt-0 failure, tagged with its attempt. + let failures = datagen_failures(&events).unwrap(); + assert_eq!(failures.len(), 1); + assert_eq!(failures[0].attempt, 0); + } + + fn write_context() -> DatagenWriteContext { + DatagenWriteContext { + run_id: "run-1".to_string(), + writer_epoch: "writer-1".to_string(), + } + } + + fn leaf_position(name: &str, index: i64, enclosing: Option<&str>) -> DatagenStreamPosition { + DatagenStreamPosition { + step: DatagenStepId { + name: name.to_string(), + kind: DatagenStepKind::Leaf, + }, + index, + enclosing: enclosing.map(str::to_string), + selector: None, + } + } + + fn set_field(name: &str, value: DatagenValue) -> DatagenFieldChange { + DatagenFieldChange { + name: name.to_string(), + field_type: "str".to_string(), + codec_version: 1, + op: FieldOp::Set, + value, + } + } + + #[test] + fn open_stream_writer_produces_a_foldable_lifecycle() { + let stream = DatagenNewStream { + item_id: DatagenItemId::from_source_key("5"), + parent_item_id: None, + query_tags: Some(json!({"lang": "en"})), + }; + let opened = open_stream_events(&stream, &write_context()); + let mut writer = opened.writer; + assert_eq!(writer.attempt(), 0); + + let position = leaf_position("gen", 0, Some("main")); + let checkpoint = writer.step_completed( + &position, + &[set_field("draft", DatagenValue::Str("v1".into()))], + ); + let terminal = writer.item_terminal(DatagenTerminal::Completed); + + let mut events = vec![opened.created_event]; + events.extend(checkpoint); + events.push(terminal); + + // Contiguous item_seq starting at 1, every event validates. + for (offset, ev) in events.iter().enumerate() { + ev.validate().unwrap(); + assert_eq!(ev.item_seq, offset as i64 + 1); + assert_eq!(ev.attempt, 0); + } + + let folded = fold_datagen_events(&events).unwrap().unwrap(); + assert_eq!(folded.status, DatagenItemStatus::Completed); + assert_eq!( + folded.fields.get("draft"), + Some(&DatagenFieldState::Set(DatagenValue::Str("v1".into()))) + ); + assert_eq!(folded.query_tags, Some(json!({"lang": "en"}))); + } + + #[test] + fn resuming_writer_continues_seq_bumps_attempt_and_avoids_event_id_collision() { + let stream = DatagenNewStream { + item_id: DatagenItemId::from_source_key("5"), + parent_item_id: None, + query_tags: None, + }; + let context = write_context(); + let opened = open_stream_events(&stream, &context); + let mut writer = opened.writer; + let position = leaf_position("gen", 0, Some("main")); + let attempt0 = writer.step_completed( + &position, + &[set_field("draft", DatagenValue::Str("v1".into()))], + ); + + let mut events = vec![opened.created_event]; + events.extend(attempt0.clone()); + + // Fold attempt 0, then resume from that folded state. + let folded0 = fold_datagen_events(&events).unwrap().unwrap(); + let mut resumed = folded0.resuming_writer(&context); + assert_eq!(resumed.attempt(), 1); + + let attempt1 = resumed.step_completed( + &position, + &[set_field("draft", DatagenValue::Str("v2".into()))], + ); + let terminal = resumed.item_terminal(DatagenTerminal::Completed); + + // Same step position across attempts must not collide on event_id (embeds attempt). + for a0 in &attempt0 { + for a1 in &attempt1 { + assert_ne!(a0.event_id, a1.event_id); + assert_ne!(a0.checkpoint_id, a1.checkpoint_id); + } + } + + events.extend(attempt1); + events.push(terminal); + assert!(events[events.len() - 2].item_seq > attempt0.last().unwrap().item_seq); + + let folded = fold_datagen_events(&events).unwrap().unwrap(); + assert_eq!(folded.status, DatagenItemStatus::Completed); + assert_eq!(folded.last_attempt, 1); + assert_eq!( + folded.fields.get("draft"), + Some(&DatagenFieldState::Set(DatagenValue::Str("v2".into()))) + ); + } + + #[test] + fn item_tree_links_parent_and_child_items() { + // Root "5" spawns child "5/expand:0". Build each via its own writer so item_ids/roots are set. + let context = write_context(); + let root_stream = DatagenNewStream { + item_id: DatagenItemId::from_source_key("5"), + parent_item_id: None, + query_tags: None, + }; + let root_open = open_stream_events(&root_stream, &context); + let mut root_writer = root_open.writer; + let root_terminal = root_writer.item_terminal(DatagenTerminal::Completed); + + let child_id = DatagenItemId::from_source_key("5").child("expand", 0); + let child_stream = DatagenNewStream { + item_id: child_id.clone(), + parent_item_id: Some(DatagenItemId::from_source_key("5")), + query_tags: None, + }; + let child_open = open_stream_events(&child_stream, &context); + let mut child_writer = child_open.writer; + let child_terminal = child_writer.item_terminal(DatagenTerminal::Completed); + + let events = vec![ + root_open.created_event, + root_terminal, + child_open.created_event, + child_terminal, + ]; + let tree = DatagenItemTree::build(&events).unwrap(); + assert_eq!(tree.len(), 2); + assert_eq!(tree.roots(), &[DatagenItemId::from_source_key("5")]); + + let root_node = tree.node(&DatagenItemId::from_source_key("5")).unwrap(); + assert_eq!(root_node.item.status, DatagenItemStatus::Completed); + assert_eq!(root_node.children, vec![child_id.clone()]); + + let child_node = tree.node(&child_id).unwrap(); + assert_eq!( + child_node.item.parent_item_id, + Some(DatagenItemId::from_source_key("5")) + ); + assert!(child_node.children.is_empty()); + } + #[test] fn sequence_collision_is_rejected() { let first = leaf_completed(1, "gen", 0, Some("main")); diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index bafa240..4a2633e 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -29,9 +29,11 @@ use tokio::task::JoinHandle; use tracing::{info, warn}; use crate::datagen::{ - datagen_failures, datagen_trajectory, fold_datagen_events, DatagenBlobValue, DatagenEvent, - DatagenEventType, DatagenFailure, DatagenItemLookup, DatagenItemStatus, - DatagenRootItemStatuses, DatagenStepCursor, DatagenStepKind, DatagenValue, + 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, }; use crate::store::{ column_as, column_as_optional, timestamp_from_micros, CompactionConfig, CompactionStats, @@ -152,6 +154,20 @@ impl DatagenStore { self.append(events).await } + /// Open a fresh stream: persist its ITEM_CREATED and return a [`DatagenStreamWriter`] positioned + /// at `item_seq = 1`, `attempt = 0`. The writer is client-side state; subsequent step batches it + /// produces are handed back to [`append`](Self::append)/[`append_checkpoint`](Self::append_checkpoint). + pub async fn open_stream( + &mut self, + stream: &DatagenNewStream, + context: &DatagenWriteContext, + ) -> LanceResult { + let opened = open_stream_events(stream, context); + self.append(std::slice::from_ref(&opened.created_event)) + .await?; + Ok(opened.writer) + } + /// Gracefully stop this store's resident MemWAL writer. pub async fn close(&mut self) -> LanceResult<()> { self.base.close().await @@ -172,6 +188,13 @@ impl DatagenStore { .await } + /// Fold a root and every projected descendant into an inspection tree (parent->child links, each + /// item at latest state). Pure over the root's event log; loads no blob bytes. + pub async fn item_tree(&self, root_item_id: &str) -> LanceResult { + let events = self.events_for_root(root_item_id).await?; + DatagenItemTree::build(&events).map_err(invalid_input) + } + /// Read failure events directly from the source-of-truth log. pub async fn failures(&self, run_id: Option<&str>) -> LanceResult> { let filter = match run_id { @@ -268,6 +291,23 @@ impl DatagenStore { .flatten()) } + /// Materialize a folded item's blob field by name, resolving the `event_id` for the caller. + /// + /// Returns `None` when `field_name` is not a blob field of `folded` (never written, or written + /// with a non-blob value); the blob bytes otherwise. This is the convenience wrapper over + /// [`FoldedDatagenItem::blob_event_ids`] + [`get_blob`](Self::get_blob) so callers never handle a + /// raw `event_id`. + pub async fn load_blob( + &self, + folded: &FoldedDatagenItem, + field_name: &str, + ) -> LanceResult>> { + match folded.blob_event_ids.get(field_name) { + Some(event_id) => self.get_blob(event_id).await, + None => Ok(None), + } + } + /// Number of flushed generations waiting across all writer shards. pub async fn pending_wal_generations(&self) -> LanceResult { self.base.pending_wal_generations().await @@ -921,8 +961,9 @@ fn invalid_input(message: impl Into) -> LanceError { mod tests { use super::*; use crate::datagen::{ - datagen_event_id, DatagenFieldState, DatagenItemId, DatagenItemLookup, DatagenItemStatus, - DatagenStepKind, DATAGEN_SCHEMA_VERSION, + datagen_event_id, DatagenFieldChange, DatagenFieldState, DatagenItemId, DatagenItemLookup, + DatagenItemStatus, DatagenStepId, DatagenStepKind, DatagenStreamPosition, DatagenTerminal, + FieldOp, DATAGEN_SCHEMA_VERSION, }; use chrono::{TimeZone, Utc}; use serde_json::json; @@ -1063,7 +1104,7 @@ mod tests { assert_eq!(blob.size, blob_bytes.len() as i64); assert_eq!( store.get_blob(&blob_event.event_id).await.unwrap(), - Some(blob_bytes) + Some(blob_bytes.clone()) ); let folded = store.fold_item("item-1").await.unwrap(); @@ -1076,6 +1117,14 @@ mod tests { assert_eq!(folded.trajectory.ordered.len(), 1); assert_eq!(folded.query_tags, Some(json!({"domain": "math"}))); + // load_blob resolves the blob field by name, and is None for a non-blob / absent field. + assert_eq!( + store.load_blob(folded, "screenshot").await.unwrap(), + Some(blob_bytes) + ); + assert_eq!(store.load_blob(folded, "score").await.unwrap(), None); + assert_eq!(store.load_blob(folded, "missing").await.unwrap(), None); + let trajectory = store.trajectory("item-1").await.unwrap(); assert_eq!(trajectory.len(), 1); assert_eq!(trajectory[0].position.step.name, "grade"); @@ -1367,4 +1416,90 @@ mod tests { assert_eq!(store.events_for_item("item-2").await.unwrap().len(), 1); }); } + + /// End-to-end through a real Lance store: `open_stream` persists ITEM_CREATED and hands back a + /// writer whose step/terminal batches append cleanly, and `item_tree` folds the persisted log + /// into a parent->child tree matching the pure-core result. + #[test] + fn open_stream_writer_appends_and_item_tree_folds_persisted_log() { + let directory = TempDir::new().unwrap(); + let uri = directory.path().to_string_lossy().to_string(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + let mut store = DatagenStore::open(&uri).await.unwrap(); + let context = DatagenWriteContext { + run_id: "run-1".to_string(), + writer_epoch: "writer-1".to_string(), + }; + + // Root "9" fans out into child "9/expand:0"; each streamed via its own writer. + let root_id = DatagenItemId::from_source_key("9"); + let mut root_writer = store + .open_stream( + &DatagenNewStream { + item_id: root_id.clone(), + parent_item_id: None, + query_tags: Some(json!({"lang": "en"})), + }, + &context, + ) + .await + .unwrap(); + assert_eq!(root_writer.attempt(), 0); + + let position = DatagenStreamPosition { + step: DatagenStepId { + name: "gen".to_string(), + kind: DatagenStepKind::Leaf, + }, + index: 0, + enclosing: None, + selector: None, + }; + let checkpoint = root_writer.step_completed( + &position, + &[DatagenFieldChange { + name: "draft".to_string(), + field_type: "str".to_string(), + codec_version: 1, + op: FieldOp::Set, + value: DatagenValue::Str("v1".into()), + }], + ); + store.append_checkpoint(&checkpoint).await.unwrap(); + let root_terminal = root_writer.item_terminal(DatagenTerminal::Completed); + store.append(&[root_terminal]).await.unwrap(); + + let child_id = root_id.child("expand", 0); + let mut child_writer = store + .open_stream( + &DatagenNewStream { + item_id: child_id.clone(), + parent_item_id: Some(root_id.clone()), + query_tags: None, + }, + &context, + ) + .await + .unwrap(); + let child_terminal = child_writer.item_terminal(DatagenTerminal::Completed); + store.append(&[child_terminal]).await.unwrap(); + + let tree = store.item_tree("9").await.unwrap(); + assert_eq!(tree.roots(), &[root_id.clone()]); + + let root_node = tree.node(&root_id).unwrap(); + assert_eq!(root_node.item.status, DatagenItemStatus::Completed); + assert_eq!(root_node.children, vec![child_id.clone()]); + assert_eq!(root_node.item.query_tags, Some(json!({"lang": "en"}))); + assert_eq!( + root_node.item.fields.get("draft"), + Some(&DatagenFieldState::Set(DatagenValue::Str("v1".into()))) + ); + + let child_node = tree.node(&child_id).unwrap(); + assert_eq!(child_node.item.parent_item_id, Some(root_id)); + assert!(child_node.children.is_empty()); + }); + } } diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index 40f2d6d..0de138c 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -22,14 +22,18 @@ mod storage; mod store; mod store_base; -pub use api_impl::rollout_record_to_dto; +pub use api_impl::{ + datagen_event_to_dto, datagen_events_from_dtos, folded_item_to_dto, rollout_record_to_dto, +}; pub use context::{Context, ContextEntry, Snapshot}; pub use datagen::{ - datagen_event_id, datagen_failures, datagen_trajectory, fold_datagen_events, DatagenBlobValue, - DatagenErrorInfo, DatagenEvent, DatagenEventType, DatagenFailure, DatagenFieldState, - DatagenItemId, DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, - DatagenStepCursor, DatagenStepId, DatagenStepKind, DatagenStreamPosition, DatagenTerminal, - DatagenTrajectory, DatagenValue, FoldedDatagenItem, DATAGEN_SCHEMA_VERSION, + datagen_event_id, datagen_failures, datagen_trajectory, fold_datagen_events, + open_stream_events, DatagenBlobValue, DatagenErrorInfo, DatagenEvent, DatagenEventType, + DatagenFailure, DatagenFieldChange, DatagenFieldState, DatagenItemId, DatagenItemLookup, + DatagenItemNode, DatagenItemStatus, DatagenItemTree, DatagenNewStream, DatagenOpenStream, + DatagenRootItemStatuses, DatagenStepCursor, DatagenStepId, DatagenStepKind, + DatagenStreamPosition, DatagenStreamWriter, DatagenTerminal, DatagenTrajectory, DatagenValue, + DatagenWriteContext, FieldOp, FoldedDatagenItem, DATAGEN_SCHEMA_VERSION, }; pub use datagen_store::{datagen_log_schema, DatagenStore, DatagenStoreOptions}; pub use eval::{ diff --git a/crates/lance-context-server/src/routes/datagen.rs b/crates/lance-context-server/src/routes/datagen.rs index bbdb35a..afb00fe 100644 --- a/crates/lance-context-server/src/routes/datagen.rs +++ b/crates/lance-context-server/src/routes/datagen.rs @@ -173,6 +173,22 @@ pub async fn datagen_item_failures( Ok(Json(ListDatagenFailuresResponse { failures })) } +/// Dump every raw event whose root item is `root_item_id`. The server does no +/// fold/tree assembly; the client builds the item tree from these events. +pub async fn datagen_events_for_root( + State(state): State>, + Path((name, root_item_id)): Path<(String, String)>, +) -> Result, AppError> { + let store_lock = state.get_or_open_datagen_store(&name).await?; + let store = store_lock.read().await; + let events = DatagenStoreApi::events_for_root(&*store, &root_item_id) + .await + .map_err(AppError::from_context)?; + Ok(Json(lance_context_api::ListDatagenEventsResponse { + events, + })) +} + #[derive(Debug, Default, serde::Deserialize)] pub struct RootStatusParams { /// Comma-separated list of root item ids to classify. diff --git a/crates/lance-context-server/src/routes/mod.rs b/crates/lance-context-server/src/routes/mod.rs index 17ab77d..18d4077 100644 --- a/crates/lance-context-server/src/routes/mod.rs +++ b/crates/lance-context-server/src/routes/mod.rs @@ -142,6 +142,10 @@ pub fn router() -> Router> { "/api/v1/datagen/{name}/root-status", get(datagen::datagen_root_item_statuses), ) + .route( + "/api/v1/datagen/{name}/roots/{root_item_id}/events", + get(datagen::datagen_events_for_root), + ) .route( "/api/v1/datagen/{name}/blobs/{event_id}", get(datagen::fetch_datagen_blob), diff --git a/crates/lance-context/src/lib.rs b/crates/lance-context/src/lib.rs index 78cfce1..aea0768 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -3,14 +3,17 @@ // Explicit re-exports from core (no glob to avoid recursion depth overflow) pub use lance_context_core::serde; pub use lance_context_core::{ - datagen_event_id, datagen_log_schema, datagen_trajectory, fold_datagen_events, - CompactionConfig, CompactionMetrics, CompactionStats, Context, ContextEntry, ContextNamespace, - ContextRecord, ContextStoreOptions, DatagenBlobValue, DatagenEvent, DatagenEventType, - DatagenFailure, DatagenFieldState, DatagenItemStatus, DatagenStepCursor, DatagenStore, - DatagenStoreOptions, DatagenTerminal, DatagenTrajectory, DatagenValue, FoldedDatagenItem, - IdIndexType, LifecycleQueryOptions, MetadataFilter, PartitionInfo, PartitionSelector, - PartitionSpec, RecordFilters, Relationship, RetrieveResult, RolloutFilters, RolloutRecord, - SearchResult, Snapshot, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, + datagen_event_id, datagen_event_to_dto, datagen_log_schema, datagen_trajectory, + fold_datagen_events, folded_item_to_dto, open_stream_events, CompactionConfig, + CompactionMetrics, CompactionStats, Context, ContextEntry, ContextNamespace, ContextRecord, + ContextStoreOptions, DatagenBlobValue, DatagenErrorInfo, DatagenEvent, DatagenEventType, + DatagenFailure, DatagenFieldChange, DatagenFieldState, DatagenItemId, DatagenItemNode, + DatagenItemStatus, DatagenItemTree, DatagenNewStream, DatagenOpenStream, DatagenStepCursor, + DatagenStepId, DatagenStepKind, DatagenStoreOptions, DatagenStreamPosition, + DatagenStreamWriter, DatagenTerminal, DatagenTrajectory, DatagenValue, DatagenWriteContext, + FieldOp, FoldedDatagenItem, IdIndexType, LifecycleQueryOptions, MetadataFilter, PartitionInfo, + PartitionSelector, PartitionSpec, RecordFilters, Relationship, RetrieveResult, RolloutFilters, + RolloutRecord, SearchResult, Snapshot, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, LIFECYCLE_CONTRADICTED, }; diff --git a/crates/lance-context/src/unified_datagen.rs b/crates/lance-context/src/unified_datagen.rs index b43861b..7677b5f 100644 --- a/crates/lance-context/src/unified_datagen.rs +++ b/crates/lance-context/src/unified_datagen.rs @@ -2,7 +2,11 @@ use lance_context_api::{ AddDatagenEventsResponse, ContextError, ContextResult, DatagenEventDto, DatagenFailureDto, DatagenRootItemStatusesResponse, DatagenStoreApi, FoldedDatagenItemDto, }; -use lance_context_core::{DatagenStore as LocalStore, DatagenStoreOptions}; +use lance_context_core::{ + datagen_event_to_dto, datagen_events_from_dtos, fold_datagen_events, open_stream_events, + DatagenEvent, DatagenItemId, DatagenItemTree, DatagenNewStream, DatagenStore as LocalStore, + DatagenStoreOptions, DatagenStreamWriter, DatagenWriteContext, +}; #[cfg(feature = "remote")] use lance_context_client::RemoteDatagenStore; @@ -58,6 +62,56 @@ impl DatagenStore { .map_err(|e| ContextError::Internal(e.to_string()))?; Ok(Self::Remote(store)) } + + /// Assemble the item tree rooted at `root_item_id`. Works for both `Local` + /// and `Remote`: raw events come from `events_for_root` and the fold/tree + /// assembly is the single-source [`DatagenItemTree::build`], so remote and + /// embedded produce identical trees. + pub async fn item_tree(&self, root_item_id: &str) -> Result { + let dtos = self.events_for_root(root_item_id).await?; + let events = datagen_events_from_dtos(&dtos)?; + DatagenItemTree::build(&events).map_err(ContextError::InvalidRequest) + } + + /// Open a fresh stream (Case 3): persist ITEM_CREATED and return a writer + /// positioned to continue after it. The writer is a pure client-side state + /// machine — its later events are appended by the caller — so this works for + /// both `Local` and `Remote` with no writer-specific endpoint. + pub async fn open_stream( + &mut self, + stream: &DatagenNewStream, + context: &DatagenWriteContext, + ) -> Result { + let opened = open_stream_events(stream, context); + let created = datagen_event_to_dto(&opened.created_event); + self.append(std::slice::from_ref(&created)).await?; + Ok(opened.writer) + } + + /// Rebuild a writer to resume an already-started item (Case 2). Pure — folds + /// the item to find `last_item_seq`/`last_attempt`, emits nothing. Returns + /// `None` if the item never started. + pub async fn resume_stream( + &self, + item_id: &str, + context: &DatagenWriteContext, + ) -> Result, ContextError> { + let dtos = self.events_for_root(&item_id_root(item_id)?).await?; + let events = datagen_events_from_dtos(&dtos)?; + let item_events: Vec = events + .into_iter() + .filter(|event| event.item_id == item_id) + .collect(); + let folded = fold_datagen_events(&item_events).map_err(ContextError::InvalidRequest)?; + Ok(folded.map(|item| item.resuming_writer(context))) + } +} + +fn item_id_root(item_id: &str) -> Result { + Ok(DatagenItemId::parse(item_id) + .map_err(ContextError::InvalidRequest)? + .root() + .to_string()) } macro_rules! dispatch_mut { @@ -120,6 +174,10 @@ impl DatagenStoreApi for DatagenStore { dispatch_ref!(self, item_failures, item_id) } + async fn events_for_root(&self, root_item_id: &str) -> ContextResult> { + dispatch_ref!(self, events_for_root, root_item_id) + } + async fn get_blob(&self, event_id: &str) -> ContextResult>> { dispatch_ref!(self, get_blob, event_id) } diff --git a/python/python/lance_context/__init__.py b/python/python/lance_context/__init__.py index 7072b2b..127e5ea 100644 --- a/python/python/lance_context/__init__.py +++ b/python/python/lance_context/__init__.py @@ -6,6 +6,7 @@ Context, ContextNamespace, DatagenStore, + DatagenStreamWriter, EmbeddingProvider, RemoteContext, RolloutStore, @@ -23,6 +24,7 @@ "Context", "ContextNamespace", "DatagenStore", + "DatagenStreamWriter", "EmbeddingProvider", "MultiModalEmbeddingProvider", "RemoteContext", diff --git a/python/python/lance_context/api.py b/python/python/lance_context/api.py index 90f3b75..d59d6aa 100644 --- a/python/python/lance_context/api.py +++ b/python/python/lance_context/api.py @@ -18,6 +18,9 @@ from ._internal import ( # pyright: ignore[reportMissingImports] DatagenStore as _DatagenStore, ) +from ._internal import ( # pyright: ignore[reportMissingImports] + DatagenStreamWriter as _DatagenStreamWriter, +) from ._internal import ( # pyright: ignore[reportMissingImports] RemoteContext as _RemoteContext, ) @@ -2839,5 +2842,123 @@ def get_blob(self, event_id: str) -> bytes | None: """Materialize one ``FIELD_*`` event's blob bytes by event id, or ``None``.""" return self._sync.get_blob(event_id) + def load_blob(self, folded: Mapping[str, Any], field_name: str) -> bytes | None: + """Resolve a folded item's blob field by name to its bytes, or ``None``. + + ``folded`` is a dict from :meth:`fold_item`; its ``blob_event_ids`` map points + each blob field to the event id that carries the bytes. Returns ``None`` when + ``field_name`` is not a blob field of ``folded`` (never written, or written with + a non-blob value). Works for embedded and remote stores. + """ + event_id = folded.get("blob_event_ids", {}).get(field_name) + if event_id is None: + return None + return self._sync.get_blob(event_id) + + def item_tree(self, root_item_id: str) -> dict[str, Any]: + """Assemble the inspection tree rooted at ``root_item_id``. + + Every projected descendant is folded to its latest state and linked + parent->child. Returns ``{"roots": [item_id, ...], "nodes": {item_id: {"item": + folded, "children": [item_id, ...]}}}``. Works for embedded and remote stores. + """ + return self._sync.item_tree(root_item_id) + + def open_stream( + self, + item_id: str, + *, + run_id: str, + writer_epoch: str, + parent_item_id: str | None = None, + query_tags: Any = None, + ) -> "DatagenStreamWriter": + """Open a fresh stream: persist ITEM_CREATED, return a writer to continue. + + ``run_id``/``writer_epoch`` stamp every event the writer emits; ``query_tags`` + is captured onto ITEM_CREATED. The writer is client-side state — hand its + emitted events back to :meth:`append`/:meth:`append_checkpoint`. Works for + embedded and remote stores. + """ + return DatagenStreamWriter( + self._sync.open_stream( + item_id, run_id, writer_epoch, parent_item_id, query_tags + ) + ) + + def resume_stream( + self, item_id: str, *, run_id: str, writer_epoch: str + ) -> "DatagenStreamWriter | None": + """Rebuild a writer to resume an already-started item, or ``None`` if absent. + + Pure — emits nothing; the returned writer continues the item's ``item_seq`` and + bumps ``attempt``. Works for embedded and remote stores. + """ + writer = self._sync.resume_stream(item_id, run_id, writer_epoch) + return DatagenStreamWriter(writer) if writer is not None else None + def __repr__(self) -> str: return f"DatagenStore(version={self._sync.version()})" + + +class DatagenStreamWriter: + """A per-stream write handle over the native writer. + + Owns the bookkeeping columns callers 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 — hand them to + :meth:`DatagenStore.append` or :meth:`DatagenStore.append_checkpoint`. Pure and + client-side, identical for embedded and remote stores. + """ + + def __init__(self, inner: _DatagenStreamWriter) -> None: + self._inner = inner + + @property + def item_id(self) -> str: + """The item this writer streams to.""" + return self._inner.item_id + + @property + def attempt(self) -> int: + """Attempt stamped onto events (0 fresh, ``last_attempt + 1`` on resume).""" + return self._inner.attempt + + def step_started(self, position: Mapping[str, Any]) -> dict[str, Any]: + """Emit STEP_STARTED for a driver frame (Sequence/Loop). Returns one event.""" + return self._inner.step_started(dict(position)) + + def step_completed( + self, position: Mapping[str, Any], fields: Sequence[Mapping[str, Any]] + ) -> list[dict[str, Any]]: + """Emit a checkpoint boundary: the step's field writes plus its STEP_COMPLETED. + + All share one ``checkpoint_id``; returns the list of event dicts (the atomic + unit :meth:`DatagenStore.append_checkpoint` persists). + """ + return self._inner.step_completed( + dict(position), [dict(field) for field in fields] + ) + + def item_terminal(self, terminal: str) -> dict[str, Any]: + """Emit TERMINAL; ``terminal`` is ``"completed"`` or ``"filtered"``.""" + return self._inner.item_terminal(terminal) + + def item_failed( + self, + position: Mapping[str, Any], + error_type: str, + *, + error_dump: str | None = None, + traceback: str | None = None, + ) -> dict[str, Any]: + """Emit FAILED at a step position (failure lens; does not terminate).""" + return self._inner.item_failed( + dict(position), error_type, error_dump, traceback + ) + + def __repr__(self) -> str: + return ( + f"DatagenStreamWriter(item_id={self._inner.item_id!r}, " + f"attempt={self._inner.attempt})" + ) diff --git a/python/src/lib.rs b/python/src/lib.rs index 79fff0a..182d9fd 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -4,7 +4,7 @@ use std::collections::{BTreeMap, HashMap, HashSet}; use std::sync::Arc; use chrono::{DateTime, SecondsFormat, Utc}; -use pyo3::exceptions::{PyRuntimeError, PyTypeError}; +use pyo3::exceptions::{PyRuntimeError, PyTypeError, PyValueError}; use pyo3::prelude::*; use pyo3::types::{PyBytes, PyDict, PyList, PyModule, PyType}; use pyo3::IntoPyObject; @@ -12,10 +12,13 @@ use serde_json::Value; use tokio::runtime::Runtime; use lance_context::{ - AddRolloutRequest, CreateDatagenStoreRequest, CreateRolloutStoreRequest, DatagenEventDto, - DatagenFailureDto, DatagenFieldStateDto, DatagenStepCursorDto, - DatagenStore as UnifiedDatagenStore, DatagenStoreApi, DatagenValueDto, FoldedDatagenItemDto, - RolloutRecordDto, RolloutStore as UnifiedRolloutStore, RolloutStoreApi, + datagen_event_to_dto, folded_item_to_dto, AddRolloutRequest, CreateDatagenStoreRequest, + CreateRolloutStoreRequest, DatagenErrorInfo, DatagenEventDto, DatagenFailureDto, + DatagenFieldChange, DatagenFieldStateDto, DatagenItemNode as UnifiedDatagenItemNode, + DatagenItemTree as UnifiedDatagenItemTree, DatagenStepCursorDto, + DatagenStore as UnifiedDatagenStore, DatagenStoreApi, DatagenStreamPosition, + DatagenStreamWriter as CoreDatagenStreamWriter, DatagenValueDto, DatagenWriteContext, + FoldedDatagenItemDto, RolloutRecordDto, RolloutStore as UnifiedRolloutStore, RolloutStoreApi, }; use lance_context_api::{ AddRecordRequest, CompactRequest, CompactResponse, CompactStatsResponse, ContextStoreApi, @@ -27,8 +30,10 @@ use lance_context_core::serde::CONTENT_TYPE_TEXT; use lance_context_core::{ datagen_event_id as core_datagen_event_id, CompactionConfig, CompactionMetrics, CompactionStats, Context as RustContext, ContextNamespace as RustContextNamespace, - ContextRecord, ContextStore, ContextStoreOptions, DistanceMetric, EvalConfig, EvalQuerySet, - ExportConfig, ExportTask, GroupBy, IdIndexType, LifecycleQueryOptions, PartitionInfo, + ContextRecord, ContextStore, ContextStoreOptions, DatagenBlobValue as CoreDatagenBlobValue, + DatagenItemId as CoreDatagenItemId, DatagenNewStream, DatagenStepId, DatagenStepKind, + DatagenTerminal, DatagenValue as CoreDatagenValue, DistanceMetric, EvalConfig, EvalQuerySet, + ExportConfig, ExportTask, FieldOp, GroupBy, IdIndexType, LifecycleQueryOptions, PartitionInfo, PartitionSelector, PartitionSpec, PreferenceForm, ReadProjection, RecordFilters, RecordPatch, Relationship, RetrievalMode, RetrieveResult, SearchResult, SplitConfig, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, @@ -2714,6 +2719,266 @@ impl DatagenStore { .map_err(to_py_err)?; Ok(bytes.map(|b| PyBytes::new(py, &b).unbind())) } + + /// Assemble the inspection tree rooted at `root_item_id`: every projected + /// descendant folded to latest state and linked parent->child. Returns a dict + /// `{"roots": [item_id, ...], "nodes": {item_id: {"item": folded, "children": + /// [item_id, ...]}}}`. + 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)?; + item_tree_to_py(py, &tree) + } + + /// Open a fresh stream: persist ITEM_CREATED and return a writer positioned to + /// continue after it. `run_id`/`writer_epoch` stamp every event the writer + /// emits; `query_tags` (JSON) is captured onto ITEM_CREATED. + #[pyo3(signature = (item_id, run_id, writer_epoch, parent_item_id = None, query_tags = None))] + fn open_stream( + &mut self, + py: Python<'_>, + item_id: &str, + run_id: &str, + writer_epoch: &str, + parent_item_id: Option<&str>, + query_tags: Option<&Bound<'_, PyAny>>, + ) -> PyResult { + let stream = DatagenNewStream { + item_id: CoreDatagenItemId::parse(item_id).map_err(to_py_err)?, + parent_item_id: parent_item_id + .map(CoreDatagenItemId::parse) + .transpose() + .map_err(to_py_err)?, + query_tags: query_tags.map(py_any_to_json).transpose()?, + }; + let context = DatagenWriteContext { + run_id: run_id.to_string(), + writer_epoch: writer_epoch.to_string(), + }; + let writer = py + .allow_threads(|| { + self.runtime + .block_on(self.store.open_stream(&stream, &context)) + }) + .map_err(to_py_err)?; + Ok(DatagenStreamWriter { inner: writer }) + } + + /// Rebuild a writer to resume an already-started item. Pure — emits nothing. + /// Returns `None` if the item never started. + fn resume_stream( + &self, + py: Python<'_>, + item_id: &str, + run_id: &str, + writer_epoch: &str, + ) -> PyResult> { + let context = DatagenWriteContext { + run_id: run_id.to_string(), + writer_epoch: writer_epoch.to_string(), + }; + let writer = py + .allow_threads(|| { + self.runtime + .block_on(self.store.resume_stream(item_id, &context)) + }) + .map_err(to_py_err)?; + Ok(writer.map(|inner| DatagenStreamWriter { inner })) + } +} + +/// 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 +/// caller hands them to `DatagenStore.append`/`append_checkpoint`. Pure and +/// client-side, so it is identical for embedded and remote stores. +#[pyclass] +struct DatagenStreamWriter { + inner: CoreDatagenStreamWriter, +} + +#[pymethods] +impl DatagenStreamWriter { + /// The item this writer streams to. + #[getter] + fn item_id(&self) -> String { + self.inner.item_id().to_string() + } + + /// The attempt number stamped onto emitted events (0 fresh, `last_attempt + 1` + /// on resume). + #[getter] + fn attempt(&self) -> i32 { + self.inner.attempt() + } + + /// Emit STEP_STARTED for a driver frame (Sequence/Loop). Returns one event dict. + fn step_started(&mut self, py: Python<'_>, position: &Bound<'_, PyDict>) -> PyResult { + let position = position_from_dict(position)?; + let event = self.inner.step_started(&position); + event_dto_to_py(py, &datagen_event_to_dto(&event)) + } + + /// Emit a checkpoint boundary: the step's field writes plus its STEP_COMPLETED, + /// all sharing one `checkpoint_id`. Returns a list of event dicts (the atomic + /// unit `append_checkpoint` persists). + fn step_completed( + &mut self, + py: Python<'_>, + position: &Bound<'_, PyDict>, + fields: &Bound<'_, PyList>, + ) -> PyResult { + let position = position_from_dict(position)?; + let changes = field_changes_from_pylist(fields)?; + let events = self.inner.step_completed(&position, &changes); + let list = PyList::empty(py); + for event in &events { + list.append(event_dto_to_py(py, &datagen_event_to_dto(event))?)?; + } + Ok(list.into_pyobject(py)?.unbind().into()) + } + + /// Emit TERMINAL — the item reached a lifecycle end. `terminal` is + /// `"completed"` or `"filtered"`. Returns one event dict. + fn item_terminal(&mut self, py: Python<'_>, terminal: &str) -> PyResult { + let terminal = match terminal { + "completed" => DatagenTerminal::Completed, + "filtered" => DatagenTerminal::Filtered, + other => { + return Err(PyValueError::new_err(format!( + "terminal must be 'completed' or 'filtered', got '{other}'" + ))) + } + }; + let event = self.inner.item_terminal(terminal); + event_dto_to_py(py, &datagen_event_to_dto(&event)) + } + + /// Emit FAILED at a step position (failure lens only; does not terminate the + /// item). Returns one event dict. + #[pyo3(signature = (position, error_type, error_dump = None, traceback = None))] + fn item_failed( + &mut self, + py: Python<'_>, + position: &Bound<'_, PyDict>, + error_type: &str, + error_dump: Option, + traceback: Option, + ) -> PyResult { + let position = position_from_dict(position)?; + let error = DatagenErrorInfo { + error_type: error_type.to_string(), + error_dump, + traceback, + }; + let event = self.inner.item_failed(&position, &error); + event_dto_to_py(py, &datagen_event_to_dto(&event)) + } +} + +fn position_from_dict(dict: &Bound<'_, PyDict>) -> PyResult { + let step_name = required_item(dict, "step_name", 0)?.extract::()?; + let step_kind_raw = required_item(dict, "step_kind", 0)?.extract::()?; + let step_kind = DatagenStepKind::parse(&step_kind_raw).map_err(to_py_err)?; + let index = required_item(dict, "index", 0)?.extract::()?; + let enclosing = optional_item(dict, "enclosing")? + .map(|value| value.extract::()) + .transpose()?; + let selector = optional_item(dict, "selector")? + .map(|value| value.extract::()) + .transpose()?; + Ok(DatagenStreamPosition { + step: DatagenStepId { + name: step_name, + kind: step_kind, + }, + index, + enclosing, + selector, + }) +} + +fn field_changes_from_pylist(fields: &Bound<'_, PyList>) -> PyResult> { + fields + .iter() + .enumerate() + .map(|(index, item)| { + let dict = item + .downcast::() + .map_err(|_| PyTypeError::new_err(format!("fields[{index}] must be a dict")))?; + field_change_from_dict(dict) + }) + .collect() +} + +fn field_change_from_dict(dict: &Bound<'_, PyDict>) -> PyResult { + let name = required_item(dict, "name", 0)?.extract::()?; + let field_type = required_item(dict, "field_type", 0)?.extract::()?; + let codec_version = optional_item(dict, "codec_version")? + .map(|value| value.extract::()) + .transpose()? + .unwrap_or(0); + let op = match optional_item(dict, "op")? + .map(|value| value.extract::()) + .transpose()? + .as_deref() + { + Some("append") => FieldOp::Append, + Some("set") | None => FieldOp::Set, + Some(other) => { + return Err(PyValueError::new_err(format!( + "field op must be 'set' or 'append', got '{other}'" + ))) + } + }; + let value_dict = required_item(dict, "value", 0)?; + let value_dict = value_dict + .downcast::() + .map_err(|_| PyTypeError::new_err("field 'value' must be a dict"))?; + let value = core_value_from_dict(value_dict)?; + Ok(DatagenFieldChange { + name, + field_type, + codec_version, + op, + value, + }) +} + +fn core_value_from_dict(dict: &Bound<'_, PyDict>) -> PyResult { + let kind = dict + .get_item("kind")? + .ok_or_else(|| PyRuntimeError::new_err("field value is missing 'kind'"))? + .extract::()?; + let inner = || { + dict.get_item("value")? + .ok_or_else(|| PyRuntimeError::new_err("field value is missing 'value'")) + }; + match kind.as_str() { + "int" => Ok(CoreDatagenValue::Int(inner()?.extract::()?)), + "float" => Ok(CoreDatagenValue::Float(inner()?.extract::()?)), + "bool" => Ok(CoreDatagenValue::Bool(inner()?.extract::()?)), + "str" => Ok(CoreDatagenValue::Str(inner()?.extract::()?)), + "json" => Ok(CoreDatagenValue::Json(py_any_to_json(&inner()?)?)), + "blob" => { + let bytes = dict + .get_item("bytes")? + .filter(|value| !value.is_none()) + .map(|value| value.extract::>()) + .transpose()? + .ok_or_else(|| PyRuntimeError::new_err("blob field value is missing 'bytes'"))?; + let size = bytes.len() as i64; + Ok(CoreDatagenValue::Blob(CoreDatagenBlobValue { + bytes: Some(bytes), + size, + checksum: None, + })) + } + other => Err(PyRuntimeError::new_err(format!( + "unsupported datagen value kind '{other}'" + ))), + } } /// Deterministic idempotency key for a checkpoint event, so retried batches @@ -2913,6 +3178,12 @@ fn folded_item_to_py(py: Python<'_>, item: &FoldedDatagenItemDto) -> PyResult None, }, )?; + + let blob_event_ids = PyDict::new(py); + for (field_name, event_id) in &item.blob_event_ids { + blob_event_ids.set_item(field_name, event_id)?; + } + dict.set_item("blob_event_ids", blob_event_ids)?; Ok(dict.into_pyobject(py)?.unbind().into()) } @@ -2954,6 +3225,82 @@ fn failure_to_py(py: Python<'_>, failure: &DatagenFailureDto) -> PyResult, event: &DatagenEventDto) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("event_id", &event.event_id)?; + dict.set_item("item_id", &event.item_id)?; + dict.set_item("root_item_id", &event.root_item_id)?; + dict.set_item("parent_item_id", event.parent_item_id.clone())?; + dict.set_item("item_seq", event.item_seq)?; + dict.set_item("checkpoint_id", &event.checkpoint_id)?; + dict.set_item("event_type", &event.event_type)?; + dict.set_item("step_name", event.step_name.clone())?; + dict.set_item("step_kind", event.step_kind.clone())?; + dict.set_item("step_index", event.step_index)?; + dict.set_item("enclosing_step", event.enclosing_step.clone())?; + dict.set_item("selector_step", event.selector_step.clone())?; + dict.set_item("attempt", event.attempt)?; + dict.set_item("run_id", &event.run_id)?; + dict.set_item("writer_epoch", &event.writer_epoch)?; + dict.set_item("field_name", event.field_name.clone())?; + dict.set_item("field_type", event.field_type.clone())?; + dict.set_item("codec_version", event.codec_version)?; + dict.set_item( + "value", + match &event.value { + Some(value) => Some(value_to_py(py, value)?), + None => None, + }, + )?; + dict.set_item( + "query_tags", + match &event.query_tags { + Some(tags) => Some(json_value_to_py(py, tags)?), + None => None, + }, + )?; + dict.set_item("status", event.status.clone())?; + dict.set_item("error_type", event.error_type.clone())?; + dict.set_item("error_dump", event.error_dump.clone())?; + dict.set_item("traceback", event.traceback.clone())?; + dict.set_item("event_ts", event.event_ts.map(|ts| ts.to_rfc3339()))?; + dict.set_item("schema_version", event.schema_version)?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + +/// One tree node: the folded item plus its direct children's item_ids. +fn item_node_to_py(py: Python<'_>, node: &UnifiedDatagenItemNode) -> PyResult { + let dict = PyDict::new(py); + dict.set_item( + "item", + folded_item_to_py(py, &folded_item_to_dto(&node.item))?, + )?; + let children = PyList::empty(py); + for child in &node.children { + children.append(child.to_string())?; + } + dict.set_item("children", children)?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + +fn item_tree_to_py(py: Python<'_>, tree: &UnifiedDatagenItemTree) -> PyResult { + let dict = PyDict::new(py); + let roots = PyList::empty(py); + for root in tree.roots() { + roots.append(root.to_string())?; + } + dict.set_item("roots", roots)?; + let nodes = PyDict::new(py); + for item in tree.items() { + let node = tree.node(&item.item_id).expect("item is in tree"); + nodes.set_item(item.item_id.to_string(), item_node_to_py(py, node)?)?; + } + dict.set_item("nodes", nodes)?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + #[pymodule] fn _internal(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_function(wrap_pyfunction!(version, m)?)?; @@ -2964,5 +3311,6 @@ fn _internal(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add_class::()?; Ok(()) } diff --git a/python/tests/test_datagen.py b/python/tests/test_datagen.py new file mode 100644 index 0000000..8cd3b6a --- /dev/null +++ b/python/tests/test_datagen.py @@ -0,0 +1,148 @@ +from __future__ import annotations + +import sys +from pathlib import Path + +PACKAGE_ROOT = Path(__file__).resolve().parents[2] / "python" / "python" +if str(PACKAGE_ROOT) not in sys.path: + sys.path.insert(0, str(PACKAGE_ROOT)) + +from lance_context.api import DatagenStore, DatagenStreamWriter # noqa: E402 + +_CONTEXT = {"run_id": "run-1", "writer_epoch": "writer-1"} + + +def _leaf_position(name: str, index: int) -> dict[str, object]: + return { + "step_name": name, + "step_kind": "leaf", + "index": index, + "enclosing": None, + "selector": None, + } + + +def _set_field(name: str, value: object, field_type: str = "str") -> dict[str, object]: + return { + "name": name, + "field_type": field_type, + "codec_version": 1, + "op": "set", + "value": {"kind": field_type, "value": value}, + } + + +def test_open_stream_writer_appends_and_folds(tmp_path: Path) -> None: + store = DatagenStore.open(str(tmp_path / "log")) + + writer = store.open_stream( + "5", run_id="run-1", writer_epoch="writer-1", query_tags={"lang": "en"} + ) + assert isinstance(writer, DatagenStreamWriter) + assert writer.item_id == "5" + assert writer.attempt == 0 + + checkpoint = writer.step_completed( + _leaf_position("gen", 0), [_set_field("draft", "v1")] + ) + store.append_checkpoint(checkpoint) + store.append([writer.item_terminal("completed")]) + + folded = store.fold_item("5") + assert folded is not None + assert folded["status"] == "completed" + assert folded["fields"]["draft"] == { + "mode": "set", + "value": {"kind": "str", "value": "v1"}, + } + assert folded["query_tags"] == {"lang": "en"} + + +def test_resume_stream_bumps_attempt(tmp_path: Path) -> None: + store = DatagenStore.open(str(tmp_path / "log")) + writer = store.open_stream("5", run_id="run-1", writer_epoch="writer-1") + store.append_checkpoint( + writer.step_completed(_leaf_position("gen", 0), [_set_field("draft", "v1")]) + ) + + resumed = store.resume_stream("5", run_id="run-1", writer_epoch="writer-2") + assert resumed is not None + assert resumed.attempt == 1 + + store.append_checkpoint( + resumed.step_completed(_leaf_position("gen", 0), [_set_field("draft", "v2")]) + ) + store.append([resumed.item_terminal("completed")]) + + folded = store.fold_item("5") + assert folded is not None + assert folded["last_attempt"] == 1 + assert folded["fields"]["draft"] == { + "mode": "set", + "value": {"kind": "str", "value": "v2"}, + } + + +def test_resume_stream_none_when_never_started(tmp_path: Path) -> None: + store = DatagenStore.open(str(tmp_path / "log")) + assert store.resume_stream("9", run_id="run-1", writer_epoch="writer-1") is None + + +def test_item_tree_links_parent_and_child(tmp_path: Path) -> None: + store = DatagenStore.open(str(tmp_path / "log")) + + root = store.open_stream("9", run_id="run-1", writer_epoch="writer-1") + store.append([root.item_terminal("completed")]) + + child = store.open_stream( + "9/expand:0", run_id="run-1", writer_epoch="writer-1", parent_item_id="9" + ) + store.append([child.item_terminal("completed")]) + + tree = store.item_tree("9") + assert tree["roots"] == ["9"] + root_node = tree["nodes"]["9"] + assert root_node["item"]["status"] == "completed" + assert root_node["children"] == ["9/expand:0"] + child_node = tree["nodes"]["9/expand:0"] + assert child_node["item"]["parent_item_id"] == "9" + assert child_node["children"] == [] + + +def test_load_blob_by_field_name(tmp_path: Path) -> None: + store = DatagenStore.open(str(tmp_path / "log")) + writer = store.open_stream("5", run_id="run-1", writer_epoch="writer-1") + + blob_field = { + "name": "screenshot", + "field_type": "blob", + "codec_version": 1, + "op": "set", + "value": {"kind": "blob", "bytes": b"payload", "size": 7}, + } + checkpoint = writer.step_completed(_leaf_position("shot", 0), [blob_field]) + store.append_checkpoint(checkpoint) + store.append([writer.item_terminal("completed")]) + + folded = store.fold_item("5") + assert folded is not None + assert store.load_blob(folded, "screenshot") == b"payload" + # A non-blob / absent field resolves to None. + assert store.load_blob(folded, "missing") is None + + +def test_item_failed_is_failure_lens_not_terminal(tmp_path: Path) -> None: + store = DatagenStore.open(str(tmp_path / "log")) + writer = store.open_stream("5", run_id="run-1", writer_epoch="writer-1") + store.append( + [writer.item_failed(_leaf_position("gen", 0), "ValueError", error_dump="boom")] + ) + + failures = store.item_failures("5") + assert len(failures) == 1 + assert failures[0]["error_type"] == "ValueError" + + # FAILED does not terminate the item. + folded = store.fold_item("5") + assert folded is not None + assert folded["status"] == "running" diff --git a/python/uv.lock b/python/uv.lock index d1c95f7..bd363d0 100644 --- a/python/uv.lock +++ b/python/uv.lock @@ -1185,7 +1185,7 @@ wheels = [ [[package]] name = "lance-context" -version = "0.6.3" +version = "0.6.4" source = { editable = "." } dependencies = [ { name = "pyarrow" }, From 7055af6e1c4e8675511d603670e6559a2b5f26ad Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Tue, 28 Jul 2026 14:13:38 -0700 Subject: [PATCH 6/7] style: rustfmt import block after merge --- python/src/lib.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/python/src/lib.rs b/python/src/lib.rs index bced3d0..7247c1d 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -18,8 +18,9 @@ use lance_context::{ 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, + DatagenWriteContext, FoldedDatagenItemDto, GenericStore as UnifiedGenericStore, + GenericStoreApi, RolloutRecordDto, RolloutStore as UnifiedRolloutStore, RolloutStoreApi, + SchemaSpec, }; use lance_context_api::{ AddRecordRequest, CompactRequest, CompactResponse, CompactStatsResponse, ContextStoreApi, From 3909251b9ab1b5c19c29ead74055a0a12a308168 Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Tue, 28 Jul 2026 14:28:24 -0700 Subject: [PATCH 7/7] fix(datagen): use slice::from_ref to satisfy clippy --- crates/lance-context-core/src/datagen_store.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index 4a2633e..6c0df9b 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -1486,7 +1486,7 @@ mod tests { store.append(&[child_terminal]).await.unwrap(); let tree = store.item_tree("9").await.unwrap(); - assert_eq!(tree.roots(), &[root_id.clone()]); + assert_eq!(tree.roots(), std::slice::from_ref(&root_id)); let root_node = tree.node(&root_id).unwrap(); assert_eq!(root_node.item.status, DatagenItemStatus::Completed);