diff --git a/crates/lance-context-api/src/lib.rs b/crates/lance-context-api/src/lib.rs index 9e30efb..e7df682 100644 --- a/crates/lance-context-api/src/lib.rs +++ b/crates/lance-context-api/src/lib.rs @@ -337,6 +337,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( @@ -1161,6 +1170,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)] @@ -1193,6 +1206,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 61200d5..b024224 100644 --- a/crates/lance-context-client/src/lib.rs +++ b/crates/lance-context-client/src/lib.rs @@ -567,6 +567,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) @@ -1313,6 +1322,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 fc44b0e..8762dfc 100644 --- a/crates/lance-context-core/src/api_impl.rs +++ b/crates/lance-context-core/src/api_impl.rs @@ -789,6 +789,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 @@ -800,7 +807,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() } @@ -920,7 +927,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(), @@ -935,6 +973,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 738d7a5..6c0df9b 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"); @@ -1266,15 +1315,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 @@ -1376,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(), std::slice::from_ref(&root_id)); + + 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 114a6c2..21bcf48 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -24,16 +24,19 @@ mod store_base; // Request/DTO conversions, exported so the server does not keep its own copies. // These were duplicated verbatim between here and `routes/`; see #214. pub use api_impl::{ - dto_to_relationship, patch_from_dto, record_from_add_request, record_to_dto, - relationship_to_dto, rollout_record_from_add_request, rollout_record_to_dto, + datagen_event_to_dto, datagen_events_from_dtos, dto_to_relationship, folded_item_to_dto, + patch_from_dto, record_from_add_request, record_to_dto, relationship_to_dto, + rollout_record_from_add_request, 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 ec7e10c..b163087 100644 --- a/crates/lance-context-server/src/routes/mod.rs +++ b/crates/lance-context-server/src/routes/mod.rs @@ -143,6 +143,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 c60616d..b0ee4e7 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -3,14 +3,18 @@ // 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, 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, + 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, }; pub use lance_context_api::{ 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 0929286..190ee97 100644 --- a/python/python/lance_context/__init__.py +++ b/python/python/lance_context/__init__.py @@ -6,6 +6,7 @@ Context, ContextNamespace, DatagenStore, + DatagenStreamWriter, EmbeddingProvider, GenericStore, RemoteContext, @@ -24,6 +25,7 @@ "Context", "ContextNamespace", "DatagenStore", + "DatagenStreamWriter", "GenericStore", "EmbeddingProvider", "MultiModalEmbeddingProvider", diff --git a/python/python/lance_context/api.py b/python/python/lance_context/api.py index 5320625..03273d0 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] GenericStore as _GenericStore, ) @@ -2843,10 +2846,128 @@ 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})" + ) + + class GenericStore: """Synchronous store over a schema you declare. diff --git a/python/src/lib.rs b/python/src/lib.rs index fe3eddc..7247c1d 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -12,11 +12,15 @@ use serde_json::{Map, Value}; use tokio::runtime::Runtime; use lance_context::{ - AddRolloutRequest, ColumnSpec, CreateDatagenStoreRequest, CreateGenericStoreRequest, - CreateRolloutStoreRequest, DatagenEventDto, DatagenFailureDto, DatagenFieldStateDto, - DatagenStepCursorDto, DatagenStore as UnifiedDatagenStore, DatagenStoreApi, DatagenValueDto, - FoldedDatagenItemDto, GenericStore as UnifiedGenericStore, GenericStoreApi, RolloutRecordDto, - RolloutStore as UnifiedRolloutStore, RolloutStoreApi, SchemaSpec, + datagen_event_to_dto, folded_item_to_dto, AddRolloutRequest, ColumnSpec, + CreateDatagenStoreRequest, CreateGenericStoreRequest, CreateRolloutStoreRequest, + 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, }; use lance_context_api::{ AddRecordRequest, CompactRequest, CompactResponse, CompactStatsResponse, ContextStoreApi, @@ -28,8 +32,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, @@ -2715,6 +2721,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 @@ -2914,6 +3180,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()) } @@ -2955,6 +3227,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()) +} + /// A store over a user-declared schema. Rows are plain dicts; the schema is /// declared once at creation and persisted in the dataset. #[pyclass] @@ -3210,6 +3558,7 @@ fn _internal(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; 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" },