From 0e25ce005dc9b14bbd27fe88c2ee4ba74f1df7ad 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/3] 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 | 868 ++++++++++++------ .../lance-context-core/src/datagen_store.rs | 145 +-- crates/lance-context-core/src/lib.rs | 9 +- 3 files changed, 708 insertions(+), 314 deletions(-) diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 8014353..eb0350c 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -1,16 +1,18 @@ -use std::collections::{BTreeMap, BTreeSet, HashMap}; +use std::collections::{BTreeMap, HashMap, HashSet}; +use std::fmt; use chrono::{DateTime, Utc}; use serde_json::Value; use uuid::Uuid; /// Current schema version for the append-only datagen checkpoint log. -pub const DATAGEN_SCHEMA_VERSION: i32 = 1; +pub const DATAGEN_SCHEMA_VERSION: i32 = 2; /// One lifecycle or field-level event in a datagen item's checkpoint history. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum DatagenEventType { ItemCreated, + StepStarted, FieldSet, FieldAppend, StepCompleted, @@ -23,6 +25,7 @@ impl DatagenEventType { pub fn as_str(self) -> &'static str { match self { Self::ItemCreated => "ITEM_CREATED", + Self::StepStarted => "STEP_STARTED", Self::FieldSet => "FIELD_SET", Self::FieldAppend => "FIELD_APPEND", Self::StepCompleted => "STEP_COMPLETED", @@ -34,6 +37,7 @@ impl DatagenEventType { pub fn parse(value: &str) -> Result { match value { "ITEM_CREATED" => Ok(Self::ItemCreated), + "STEP_STARTED" => Ok(Self::StepStarted), "FIELD_SET" => Ok(Self::FieldSet), "FIELD_APPEND" => Ok(Self::FieldAppend), "STEP_COMPLETED" => Ok(Self::StepCompleted), @@ -44,7 +48,171 @@ impl DatagenEventType { } } -/// Terminal outcome of an item. +/// The composition kind of a step. Drives two behaviors: +/// - forks a stream: `MapReduce`/`Branch`/`SubPipeline` sub-items get their own `item_id`. +/// - drives a frame: `Sequence`/`Loop` are the `enclosing_step` frames; only these emit STEP_STARTED. +/// +/// `Conditional`/`Router` are selectors (wrap one chosen child); `Leaf`/`Root` are the endpoints. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub enum DatagenStepKind { + Root, + Leaf, + Sequence, + Loop, + MapReduce, + Branch, + SubPipeline, + Conditional, + Router, +} + +impl DatagenStepKind { + #[must_use] + pub fn as_str(self) -> &'static str { + match self { + Self::Root => "root", + Self::Leaf => "leaf", + Self::Sequence => "sequence", + Self::Loop => "loop", + Self::MapReduce => "map_reduce", + Self::Branch => "branch", + Self::SubPipeline => "sub_pipeline", + Self::Conditional => "conditional", + Self::Router => "router", + } + } + + pub fn parse(value: &str) -> Result { + match value { + "root" => Ok(Self::Root), + "leaf" => Ok(Self::Leaf), + "sequence" => Ok(Self::Sequence), + "loop" => Ok(Self::Loop), + "map_reduce" => Ok(Self::MapReduce), + "branch" => Ok(Self::Branch), + "sub_pipeline" => Ok(Self::SubPipeline), + "conditional" => Ok(Self::Conditional), + "router" => Ok(Self::Router), + other => Err(format!("unsupported datagen step kind '{other}'")), + } + } + + /// A driver frame (`Sequence`/`Loop`) is the only kind that emits STEP_STARTED. + #[must_use] + pub fn is_driver(self) -> bool { + matches!(self, Self::Sequence | Self::Loop) + } +} + +/// An item (stream) identity — a structured value stored as a materialized path string. Root ids are +/// built by the executor from a source key (e.g. `5`); sub-item ids extend a parent with one fan-out +/// segment (`5/expand:0`) via [`DatagenItemId::child`]. The store owns the parse<->format; the client +/// only ever holds the structured form. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct DatagenItemId { + root_key: String, + /// One `(origin_step, branch_idx)` per fan-out hop below the root. + segments: Vec<(String, i64)>, +} + +impl DatagenItemId { + /// Build a root id from the executor's source key (the base case). + #[must_use] + pub fn from_source_key(key: &str) -> Self { + Self { + root_key: key.to_string(), + segments: Vec::new(), + } + } + + /// Extend this id with one fan-out segment -> the sub-item's id. Pure, deterministic, no I/O. + /// The sole id-composition entry point. Valid to call for a branch that was never written. + #[must_use] + pub fn child(&self, origin_step: &str, branch_idx: i64) -> Self { + let mut segments = self.segments.clone(); + segments.push((origin_step.to_string(), branch_idx)); + Self { + root_key: self.root_key.clone(), + segments, + } + } + + /// The parent stream's id (`None` on a root). + #[must_use] + pub fn parent(&self) -> Option { + if self.segments.is_empty() { + return None; + } + let mut segments = self.segments.clone(); + segments.pop(); + Some(Self { + root_key: self.root_key.clone(), + segments, + }) + } + + /// The root of this id's tree (== self if root). + #[must_use] + pub fn root(&self) -> Self { + Self { + root_key: self.root_key.clone(), + segments: Vec::new(), + } + } + + /// The fan-out step that created this sub-item (`None` on a root). + #[must_use] + pub fn origin_step(&self) -> Option<&str> { + self.segments.last().map(|(step, _)| step.as_str()) + } + + /// Which branch this sub-item is (`None` on a root). + #[must_use] + pub fn branch_idx(&self) -> Option { + self.segments.last().map(|(_, idx)| *idx) + } + + #[must_use] + pub fn is_root(&self) -> bool { + self.segments.is_empty() + } + + /// Parse a stored path string back into a structured id. + pub fn parse(path: &str) -> Result { + let mut parts = path.split('/'); + let root_key = parts + .next() + .filter(|part| !part.is_empty()) + .ok_or_else(|| format!("empty datagen item id '{path}'"))? + .to_string(); + let mut segments = Vec::new(); + for part in parts { + let (step, idx) = part + .split_once(':') + .ok_or_else(|| format!("malformed item id segment '{part}' in '{path}'"))?; + let branch_idx = idx + .parse::() + .map_err(|_| format!("non-integer branch index '{idx}' in '{path}'"))?; + segments.push((step.to_string(), branch_idx)); + } + Ok(Self { + root_key, + segments, + }) + } +} + +impl fmt::Display for DatagenItemId { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(&self.root_key)?; + for (step, idx) in &self.segments { + write!(formatter, "/{step}:{idx}")?; + } + Ok(()) + } +} + +/// Terminal outcome of an item (the write-side input to [`DatagenStreamWriter::item_terminal`]). #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum DatagenTerminal { Completed, @@ -59,28 +227,43 @@ impl DatagenTerminal { Self::Filtered => "filtered", } } - - pub fn parse(value: &str) -> Result { - match value { - "completed" => Ok(Self::Completed), - "filtered" => Ok(Self::Filtered), - other => Err(format!("unsupported datagen terminal value '{other}'")), - } - } } -/// Current status derived exclusively by folding the event log. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +/// The value type of the `status` column. Two read lenses: lifecycle {Running, Completed, Filtered} +/// vs failure {Failed}. A folded item's status is only ever a lifecycle value; `Failed` surfaces only +/// through the failure lens (overview / failure history). +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)] pub enum DatagenItemStatus { - Pending, Running, Completed, Filtered, Failed, } -/// Lazy reference to an inline blob event. `bytes` is absent on normal fold and -/// trajectory reads; callers materialize it through `DatagenStore::get_blob`. +impl DatagenItemStatus { + #[must_use] + pub fn as_str(self) -> &'static str { + match self { + Self::Running => "running", + Self::Completed => "completed", + Self::Filtered => "filtered", + Self::Failed => "failed", + } + } + + pub fn parse(value: &str) -> Result { + match value { + "running" => Ok(Self::Running), + "completed" => Ok(Self::Completed), + "filtered" => Ok(Self::Filtered), + "failed" => Ok(Self::Failed), + other => Err(format!("unsupported datagen status '{other}'")), + } + } +} + +/// Lazy reference to an inline blob event. `bytes` is absent on normal (lazy) fold and trajectory +/// reads; callers materialize it through `DatagenStore::load_blob`. #[derive(Debug, Clone, PartialEq, Eq)] pub struct DatagenBlobValue { pub bytes: Option>, @@ -94,7 +277,7 @@ pub enum DatagenValue { Int(i64), Float(f64), Bool(bool), - String(String), + Str(String), Json(Value), Blob(DatagenBlobValue), } @@ -106,14 +289,159 @@ impl DatagenValue { Self::Int(_) => "int", Self::Float(_) => "float", Self::Bool(_) => "bool", - Self::String(_) => "str", + Self::Str(_) => "str", Self::Json(_) => "json", Self::Blob(_) => "blob", } } } -/// A single append-only row in `log.lance`. +/// A step's identity: its (globally-unique) name + its kind. Name and kind always travel together. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct DatagenStepId { + pub name: String, + pub kind: DatagenStepKind, +} + +/// A position within one stream's step tree — the coordinate a step write is attributed to. Maps 1:1 +/// to the Group C provenance columns. `enclosing`/`selector` are stored as bare step *names* (the log +/// has no enclosing/selector kind column); `None` means "directly under the stream root" / "no +/// selector". +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub struct DatagenStreamPosition { + pub step: DatagenStepId, + pub index: i64, + pub enclosing: Option, + pub selector: Option, +} + +/// How a field folds: FIELD_SET replaces (last-writer-wins); FIELD_APPEND accumulates in order. +#[derive(Debug, Clone, PartialEq)] +pub enum DatagenFieldState { + Set(DatagenValue), + Appended(Vec), +} + +/// A pointer to one completed step position (a single STEP_COMPLETED). +#[derive(Debug, Clone, PartialEq)] +pub struct DatagenStepCursor { + pub position: DatagenStreamPosition, + /// The STEP_COMPLETED's `item_seq` — the fold cutoff for "state as of this step". + pub item_seq: i64, +} + +/// The ordered list of cursors an item passed through, plus sets for O(1) skip lookup. `completed` +/// gates STEP_COMPLETED (re-)emission; `started` gates STEP_STARTED. `started \ completed` = frames +/// that were open when the process died. +#[derive(Debug, Clone, Default, PartialEq)] +pub struct DatagenTrajectory { + pub ordered: Vec, + pub completed: HashSet, + pub started: HashSet, +} + +/// Error payload, shared by the write side (input to `item_failed`) and the read side (composed into +/// [`DatagenFailure`]). +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DatagenErrorInfo { + pub error_type: String, + pub error_dump: Option, + pub traceback: Option, +} + +/// One failure record — a lightweight pointer to a FAILED event (no folded item). An item may have +/// 0..N of these across attempts. +#[derive(Debug, Clone, PartialEq)] +pub struct DatagenFailure { + pub at: DatagenStepCursor, + pub run_id: String, + pub attempt: i32, + pub error: DatagenErrorInfo, +} + +/// One item reconstructed by folding its events (latest state). Carries enough to continue processing +/// and to rebuild a `DatagenStreamWriter` on resume. Does not carry failures. +#[derive(Debug, Clone, PartialEq)] +pub struct FoldedDatagenItem { + pub item_id: DatagenItemId, + pub root_item_id: DatagenItemId, + pub parent_item_id: Option, + pub status: DatagenItemStatus, + /// Max `item_seq` -> resume continues at `last_item_seq + 1`. + pub last_item_seq: i64, + /// Max `attempt` seen -> resume runs at `last_attempt + 1`. + pub last_attempt: i32, + pub fields: BTreeMap, + pub trajectory: DatagenTrajectory, + pub query_tags: Option, + /// Internal `field_name -> event_id` map for the folded blob fields, so the store can resolve a + /// lazy blob without the caller handling an `event_id`. + pub blob_event_ids: BTreeMap, +} + +/// 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)] +pub enum DatagenItemLookup { + NeverStarted, + Found(FoldedDatagenItem), +} + +impl DatagenItemLookup { + #[must_use] + pub fn folded(&self) -> Option<&FoldedDatagenItem> { + match self { + Self::NeverStarted => None, + Self::Found(item) => Some(item), + } + } +} + +/// Bulk startup classification of root items. A missing id means "never started". +#[derive(Debug, Clone, Default, PartialEq)] +pub struct DatagenRootItemStatuses { + inner: HashMap, +} + +impl DatagenRootItemStatuses { + #[must_use] + pub fn from_map(inner: HashMap) -> Self { + Self { inner } + } + + /// The classified status of a root item, or `None` if it was never started. + #[must_use] + pub fn status(&self, item_id: &DatagenItemId) -> Option { + self.inner.get(&item_id.to_string()).copied() + } + + /// Whether this item reached a terminal lifecycle state (Completed | Filtered). + #[must_use] + pub fn is_terminated(&self, item_id: &DatagenItemId) -> bool { + matches!( + self.status(item_id), + Some(DatagenItemStatus::Completed | DatagenItemStatus::Filtered) + ) + } + + #[must_use] + pub fn len(&self) -> usize { + self.inner.len() + } + + #[must_use] + pub fn is_empty(&self) -> bool { + self.inner.is_empty() + } + + #[must_use] + pub fn iter(&self) -> impl Iterator { + self.inner.iter() + } +} + +/// One append-only row in `log.lance`. Item/root/parent ids are the stored path strings; fold parses +/// them into [`DatagenItemId`]. #[derive(Debug, Clone, PartialEq)] pub struct DatagenEvent { /// Deterministic idempotency key. The MemWAL read path de-duplicates by it. @@ -121,16 +449,16 @@ pub struct DatagenEvent { pub item_id: String, pub root_item_id: String, pub parent_item_id: Option, - /// Strictly increasing per item. A collision between different event ids is - /// treated as split-brain corruption during fold. + /// Strictly increasing per item. A collision between different event ids is split-brain corruption. pub item_seq: i64, /// Shared by every event emitted for one checkpoint boundary. pub checkpoint_id: String, pub event_type: DatagenEventType, pub step_name: Option, + pub step_kind: Option, pub step_index: Option, - pub step_instance_id: Option, - pub iteration: Option, + pub enclosing_step: Option, + pub selector_step: Option, pub attempt: i32, pub run_id: String, /// Fencing identity for the writer/lease that owned this item. @@ -140,9 +468,10 @@ pub struct DatagenEvent { pub field_type: Option, pub codec_version: Option, pub value: Option, - /// Query tags captured on ITEM_CREATED. They are not part of correctness. + /// Query tags captured on ITEM_CREATED. Not part of correctness. pub query_tags: Option, - pub terminal: Option, + /// The stored lifecycle/failure status (populated on ITEM_CREATED / TERMINAL / FAILED). + pub status: Option, pub error_type: Option, pub error_dump: Option, pub traceback: Option, @@ -188,41 +517,38 @@ impl DatagenEvent { if self.value.is_none() { return Err("field events require a value".to_string()); } - if self.step_name.as_deref().is_none_or(str::is_empty) - || self.step_index.is_none() - || self.step_instance_id.as_deref().is_none_or(str::is_empty) - { - return Err( - "field events require step_name, step_index, and step_instance_id" - .to_string(), - ); - } + self.require_step_provenance("field events")?; } - DatagenEventType::StepCompleted => { - if self.step_name.as_deref().is_none_or(str::is_empty) - || self.step_index.is_none() - || self.step_instance_id.as_deref().is_none_or(str::is_empty) - { - return Err( - "STEP_COMPLETED requires step_name, step_index, and step_instance_id" - .to_string(), - ); - } + DatagenEventType::StepStarted | DatagenEventType::StepCompleted => { + self.require_step_provenance(self.event_type.as_str())?; } DatagenEventType::Failed => { if self.error_type.as_deref().is_none_or(str::is_empty) { return Err("FAILED requires error_type".to_string()); } + self.require_step_provenance("FAILED")?; } - DatagenEventType::Terminal => { - if self.terminal.is_none() { - return Err("TERMINAL requires terminal".to_string()); - } - } + DatagenEventType::Terminal => match self.status { + Some(DatagenItemStatus::Completed | DatagenItemStatus::Filtered) => {} + _ => return Err("TERMINAL requires status completed or filtered".to_string()), + }, DatagenEventType::ItemCreated => {} } Ok(()) } + + fn require_step_provenance(&self, context: &str) -> Result<(), String> { + if self.step_name.as_deref().is_none_or(str::is_empty) { + return Err(format!("{context} require step_name")); + } + if self.step_kind.is_none() { + return Err(format!("{context} require step_kind")); + } + if self.step_index.is_none() { + return Err(format!("{context} require step_index")); + } + Ok(()) + } } /// Generate a deterministic event id for retry-safe checkpoint ingestion. @@ -232,89 +558,52 @@ pub fn datagen_event_id(item_id: &str, checkpoint_id: &str, ordinal: u32) -> Str Uuid::new_v5(&Uuid::NAMESPACE_OID, input.as_bytes()).to_string() } -#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] -pub struct DatagenStepCursor { - pub checkpoint_id: String, - pub step_name: String, - pub step_index: i64, - pub step_instance_id: String, - pub iteration: Option, - pub attempt: i32, -} - -#[derive(Debug, Clone, PartialEq)] -pub enum DatagenFieldState { - Set(DatagenValue), - Appended(Vec), -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct DatagenFailure { - pub event_id: String, - pub run_id: String, - pub checkpoint_id: String, - pub item_seq: i64, - pub step_name: Option, - pub error_type: String, - pub error_dump: Option, - pub traceback: Option, - pub failed_at: DateTime, -} - -/// Current item state reconstructed solely from the append-only log. -#[derive(Debug, Clone, PartialEq)] -pub struct FoldedDatagenItem { - pub item_id: String, - pub root_item_id: String, - pub parent_item_id: Option, - pub fields: BTreeMap, - pub completed_steps: BTreeSet, - pub status: DatagenItemStatus, - pub terminal: Option, - pub failure: Option, - pub query_tags: Option, - pub current_run_id: String, - pub last_item_seq: i64, - pub last_checkpoint_id: String, -} - -/// State captured immediately after a STEP_COMPLETED event. -#[derive(Debug, Clone, PartialEq)] -pub struct DatagenTrajectoryPoint { - pub cursor: DatagenStepCursor, - pub item: FoldedDatagenItem, -} - -pub fn fold_datagen_events(events: &[DatagenEvent]) -> Result { +/// Fold an item's events into its latest state. Returns `None` if there is no ITEM_CREATED (the item +/// was never started). +pub fn fold_datagen_events(events: &[DatagenEvent]) -> Result, String> { let ordered = normalize_events(events)?; - let first = ordered - .first() - .ok_or_else(|| "cannot fold an empty datagen event list".to_string())?; - let mut item = initial_item(first); + let Some(first) = ordered.first() else { + return Ok(None); + }; + if first.event_type != DatagenEventType::ItemCreated { + return Ok(None); + } + let mut item = initial_item(first)?; for event in ordered { apply_event(&mut item, event)?; } - Ok(item) + Ok(Some(item)) +} + +/// The ordered completed-step cursors of an item, in `item_seq` order. +pub fn datagen_trajectory(events: &[DatagenEvent]) -> Result, String> { + Ok(fold_datagen_events(events)? + .map(|item| item.trajectory.ordered) + .unwrap_or_default()) } -pub fn datagen_trajectory(events: &[DatagenEvent]) -> Result, String> { +/// The lightweight failure pointers of an item, in `item_seq` order. +pub fn datagen_failures(events: &[DatagenEvent]) -> Result, String> { let ordered = normalize_events(events)?; - let first = ordered - .first() - .ok_or_else(|| "cannot build a trajectory from an empty event list".to_string())?; - let mut item = initial_item(first); - let mut trajectory = Vec::new(); + let mut failures = Vec::new(); for event in ordered { - apply_event(&mut item, event)?; - if event.event_type == DatagenEventType::StepCompleted { - let cursor = step_cursor(event)?; - trajectory.push(DatagenTrajectoryPoint { - cursor, - item: item.clone(), + if event.event_type == DatagenEventType::Failed { + failures.push(DatagenFailure { + at: DatagenStepCursor { + position: stream_position(event)?, + item_seq: event.item_seq, + }, + run_id: event.run_id.clone(), + attempt: event.attempt, + error: DatagenErrorInfo { + error_type: event.error_type.clone().unwrap(), + error_dump: event.error_dump.clone(), + traceback: event.traceback.clone(), + }, }); } } - Ok(trajectory) + Ok(failures) } fn normalize_events(events: &[DatagenEvent]) -> Result, String> { @@ -352,69 +641,56 @@ fn normalize_events(events: &[DatagenEvent]) -> Result, Strin Ok(ordered) } -fn initial_item(first: &DatagenEvent) -> FoldedDatagenItem { - FoldedDatagenItem { - item_id: first.item_id.clone(), - root_item_id: first.root_item_id.clone(), - parent_item_id: first.parent_item_id.clone(), +fn initial_item(first: &DatagenEvent) -> Result { + let parent_item_id = match &first.parent_item_id { + Some(parent) => Some(DatagenItemId::parse(parent)?), + None => None, + }; + Ok(FoldedDatagenItem { + item_id: DatagenItemId::parse(&first.item_id)?, + root_item_id: DatagenItemId::parse(&first.root_item_id)?, + parent_item_id, + status: DatagenItemStatus::Running, + last_item_seq: first.item_seq, + last_attempt: first.attempt, fields: BTreeMap::new(), - completed_steps: BTreeSet::new(), - status: DatagenItemStatus::Pending, - terminal: None, - failure: None, + trajectory: DatagenTrajectory::default(), query_tags: None, - current_run_id: first.run_id.clone(), - last_item_seq: first.item_seq, - last_checkpoint_id: first.checkpoint_id.clone(), - } + blob_event_ids: BTreeMap::new(), + }) } fn apply_event(item: &mut FoldedDatagenItem, event: &DatagenEvent) -> Result<(), String> { - if event.item_id != item.item_id { + if DatagenItemId::parse(&event.item_id)? != item.item_id { return Err(format!( "event '{}' belongs to item '{}', expected '{}'", event.event_id, event.item_id, item.item_id )); } - if event.root_item_id != item.root_item_id { - return Err(format!( - "item '{}' changed root_item_id from '{}' to '{}'", - item.item_id, item.root_item_id, event.root_item_id - )); - } - if event.parent_item_id != item.parent_item_id { - return Err(format!( - "item '{}' changed parent_item_id during its trajectory", - item.item_id - )); - } - item.current_run_id = event.run_id.clone(); - item.last_item_seq = event.item_seq; - item.last_checkpoint_id = event.checkpoint_id.clone(); + item.last_item_seq = item.last_item_seq.max(event.item_seq); + item.last_attempt = item.last_attempt.max(event.attempt); match event.event_type { DatagenEventType::ItemCreated => { - item.status = DatagenItemStatus::Pending; - item.terminal = None; - item.failure = None; + item.status = DatagenItemStatus::Running; if event.query_tags.is_some() { item.query_tags = event.query_tags.clone(); } } + DatagenEventType::StepStarted => { + item.trajectory.started.insert(stream_position(event)?); + } DatagenEventType::FieldSet => { let field_name = event.field_name.clone().unwrap(); - item.fields.insert( - field_name, - DatagenFieldState::Set(event.value.clone().unwrap()), - ); - item.status = DatagenItemStatus::Running; - item.terminal = None; - item.failure = None; + let value = event.value.clone().unwrap(); + record_blob_event_id(item, &field_name, &value, &event.event_id); + item.fields.insert(field_name, DatagenFieldState::Set(value)); } DatagenEventType::FieldAppend => { let field_name = event.field_name.clone().unwrap(); let value = event.value.clone().unwrap(); + record_blob_event_id(item, &field_name, &value, &event.event_id); match item.fields.entry(field_name) { std::collections::btree_map::Entry::Vacant(entry) => { entry.insert(DatagenFieldState::Appended(vec![value])); @@ -429,60 +705,56 @@ fn apply_event(item: &mut FoldedDatagenItem, event: &DatagenEvent) -> Result<(), } }, } - item.status = DatagenItemStatus::Running; - item.terminal = None; - item.failure = None; } DatagenEventType::StepCompleted => { - item.completed_steps.insert(step_cursor(event)?); - item.status = DatagenItemStatus::Running; - item.terminal = None; - item.failure = None; - } - DatagenEventType::Failed => { - item.status = DatagenItemStatus::Failed; - item.terminal = None; - item.failure = Some(DatagenFailure { - event_id: event.event_id.clone(), - run_id: event.run_id.clone(), - checkpoint_id: event.checkpoint_id.clone(), + let position = stream_position(event)?; + item.trajectory.completed.insert(position.clone()); + item.trajectory.ordered.push(DatagenStepCursor { + position, item_seq: event.item_seq, - step_name: event.step_name.clone(), - error_type: event.error_type.clone().unwrap(), - error_dump: event.error_dump.clone(), - traceback: event.traceback.clone(), - failed_at: event.event_ts, }); } + DatagenEventType::Failed => { + // Failure lens only: a FAILED row leaves the item Running under the lifecycle lens. + } DatagenEventType::Terminal => { - let terminal = event.terminal.unwrap(); - item.status = match terminal { - DatagenTerminal::Completed => DatagenItemStatus::Completed, - DatagenTerminal::Filtered => DatagenItemStatus::Filtered, + item.status = match event.status { + Some(status @ (DatagenItemStatus::Completed | DatagenItemStatus::Filtered)) => status, + _ => return Err("TERMINAL event missing completed/filtered status".to_string()), }; - item.terminal = Some(terminal); - item.failure = None; } } Ok(()) } -fn step_cursor(event: &DatagenEvent) -> Result { - Ok(DatagenStepCursor { - checkpoint_id: event.checkpoint_id.clone(), - step_name: event - .step_name - .clone() - .ok_or_else(|| "STEP_COMPLETED missing step_name".to_string())?, - step_index: event +fn record_blob_event_id( + item: &mut FoldedDatagenItem, + field_name: &str, + value: &DatagenValue, + event_id: &str, +) { + if matches!(value, DatagenValue::Blob(_)) { + item.blob_event_ids + .insert(field_name.to_string(), event_id.to_string()); + } +} + +fn stream_position(event: &DatagenEvent) -> Result { + Ok(DatagenStreamPosition { + step: DatagenStepId { + name: event + .step_name + .clone() + .ok_or_else(|| "step event missing step_name".to_string())?, + kind: event + .step_kind + .ok_or_else(|| "step event missing step_kind".to_string())?, + }, + index: event .step_index - .ok_or_else(|| "STEP_COMPLETED missing step_index".to_string())?, - step_instance_id: event - .step_instance_id - .clone() - .ok_or_else(|| "STEP_COMPLETED missing step_instance_id".to_string())?, - iteration: event.iteration, - attempt: event.attempt, + .ok_or_else(|| "step event missing step_index".to_string())?, + enclosing: event.enclosing_step.clone(), + selector: event.selector_step.clone(), }) } @@ -495,17 +767,18 @@ mod tests { fn event(seq: i64, event_type: DatagenEventType) -> DatagenEvent { let checkpoint_id = format!("checkpoint-{seq}"); DatagenEvent { - event_id: datagen_event_id("item-1", &checkpoint_id, 0), - item_id: "item-1".to_string(), - root_item_id: "item-1".to_string(), + event_id: datagen_event_id("5", &checkpoint_id, 0), + item_id: "5".to_string(), + root_item_id: "5".to_string(), parent_item_id: None, item_seq: seq, checkpoint_id, event_type, step_name: None, + step_kind: None, step_index: None, - step_instance_id: None, - iteration: None, + enclosing_step: None, + selector_step: None, attempt: 0, run_id: "run-1".to_string(), writer_epoch: "writer-1".to_string(), @@ -514,7 +787,7 @@ mod tests { codec_version: None, value: None, query_tags: None, - terminal: None, + status: None, error_type: None, error_dump: None, traceback: None, @@ -523,77 +796,158 @@ mod tests { } } - fn completed_step(seq: i64) -> DatagenEvent { - let mut event = event(seq, DatagenEventType::StepCompleted); - event.step_name = Some("noop".to_string()); - event.step_index = Some(3); - event.step_instance_id = Some("loop/2/noop".to_string()); - event.iteration = Some(2); - event + fn created(seq: i64) -> DatagenEvent { + let mut created = event(seq, DatagenEventType::ItemCreated); + created.status = Some(DatagenItemStatus::Running); + created + } + + fn leaf_completed(seq: i64, name: &str, index: i64, enclosing: Option<&str>) -> DatagenEvent { + let mut completed = event(seq, DatagenEventType::StepCompleted); + completed.step_name = Some(name.to_string()); + completed.step_kind = Some(DatagenStepKind::Leaf); + completed.step_index = Some(index); + completed.enclosing_step = enclosing.map(str::to_string); + completed + } + + fn driver_started(seq: i64, name: &str, index: i64, enclosing: Option<&str>) -> DatagenEvent { + let mut started = event(seq, DatagenEventType::StepStarted); + started.step_name = Some(name.to_string()); + started.step_kind = Some(DatagenStepKind::Sequence); + started.step_index = Some(index); + started.enclosing_step = enclosing.map(str::to_string); + started } #[test] - fn no_op_step_is_present_in_fold_and_trajectory() { - let created = event(0, DatagenEventType::ItemCreated); - let completed = completed_step(1); + fn item_id_round_trips_and_navigates() { + let root = DatagenItemId::from_source_key("5"); + let child = root.child("expand", 0).child("enrich", 1); + assert_eq!(child.to_string(), "5/expand:0/enrich:1"); + assert_eq!(DatagenItemId::parse("5/expand:0/enrich:1").unwrap(), child); + assert_eq!(child.origin_step(), Some("enrich")); + assert_eq!(child.branch_idx(), Some(1)); + assert_eq!(child.parent().unwrap().to_string(), "5/expand:0"); + assert_eq!(child.root(), root); + assert!(root.is_root()); + assert_eq!(root.origin_step(), None); + } - let folded = fold_datagen_events(&[completed.clone(), created]).unwrap(); - assert_eq!(folded.status, DatagenItemStatus::Running); - assert_eq!(folded.completed_steps.len(), 1); - assert!(folded.fields.is_empty()); + #[test] + fn never_started_folds_to_none() { + let completed = leaf_completed(1, "gen", 0, Some("main")); + assert!(fold_datagen_events(&[completed]).unwrap().is_none()); + } - let trajectory = datagen_trajectory(&[completed]).unwrap(); - assert_eq!(trajectory.len(), 1); - assert_eq!(trajectory[0].cursor.step_instance_id, "loop/2/noop"); + #[test] + fn fold_tracks_started_and_completed_positions() { + let events = [ + created(0), + driver_started(1, "main", 0, None), + leaf_completed(2, "gen", 0, Some("main")), + ]; + let folded = fold_datagen_events(&events).unwrap().unwrap(); + assert_eq!(folded.status, DatagenItemStatus::Running); + assert_eq!(folded.trajectory.ordered.len(), 1); + assert_eq!(folded.trajectory.completed.len(), 1); + assert_eq!(folded.trajectory.started.len(), 1); + assert_eq!(folded.last_item_seq, 2); } #[test] - fn retry_duplicate_event_is_folded_once() { - let created = event(0, DatagenEventType::ItemCreated); - let mut append = event(1, DatagenEventType::FieldAppend); - append.field_name = Some("messages".to_string()); + fn field_set_is_last_writer_wins_and_append_accumulates() { + let mut set_v1 = leaf_completed(1, "gen", 0, Some("main")); + set_v1.event_type = DatagenEventType::FieldSet; + set_v1.field_name = Some("draft".to_string()); + set_v1.field_type = Some("str".to_string()); + set_v1.codec_version = Some(1); + set_v1.value = Some(DatagenValue::Str("v1".to_string())); + + let mut set_v2 = set_v1.clone(); + set_v2.item_seq = 2; + set_v2.checkpoint_id = "c2".to_string(); + set_v2.event_id = datagen_event_id("5", "c2", 0); + set_v2.value = Some(DatagenValue::Str("v2".to_string())); + + let mut append = leaf_completed(3, "b1", 0, Some("body")); + append.event_type = DatagenEventType::FieldAppend; + append.field_name = Some("revisions".to_string()); append.field_type = Some("json".to_string()); append.codec_version = Some(1); - append.value = Some(DatagenValue::Json(json!({"role": "assistant"}))); - append.step_name = Some("generate".to_string()); - append.step_index = Some(1); - append.step_instance_id = Some("generate/0".to_string()); - - let folded = fold_datagen_events(&[created, append.clone(), append.clone()]).unwrap(); + append.value = Some(DatagenValue::Json(json!({"n": "a"}))); + let mut append2 = append.clone(); + append2.item_seq = 4; + append2.checkpoint_id = "c4".to_string(); + append2.event_id = datagen_event_id("5", "c4", 0); + append2.value = Some(DatagenValue::Json(json!({"n": "b"}))); + + let folded = + fold_datagen_events(&[created(0), set_v1, set_v2, append, append2]).unwrap().unwrap(); assert_eq!( - folded.fields.get("messages"), - Some(&DatagenFieldState::Appended(vec![DatagenValue::Json( - json!({"role": "assistant"}) - )])) + folded.fields.get("draft"), + Some(&DatagenFieldState::Set(DatagenValue::Str("v2".to_string()))) + ); + assert_eq!( + folded.fields.get("revisions"), + Some(&DatagenFieldState::Appended(vec![ + DatagenValue::Json(json!({"n": "a"})), + DatagenValue::Json(json!({"n": "b"})), + ])) ); } + #[test] + fn terminal_sets_lifecycle_status() { + let mut terminal = event(2, DatagenEventType::Terminal); + terminal.status = Some(DatagenItemStatus::Completed); + let folded = fold_datagen_events(&[created(0), terminal]).unwrap().unwrap(); + assert_eq!(folded.status, DatagenItemStatus::Completed); + } + + #[test] + fn failed_leaves_item_running_but_surfaces_in_failures() { + let mut failed = leaf_completed(1, "check", 2, Some("solve")); + failed.event_type = DatagenEventType::Failed; + failed.status = Some(DatagenItemStatus::Failed); + failed.error_type = Some("ValueError".to_string()); + failed.attempt = 0; + + let folded = fold_datagen_events(&[created(0), failed.clone()]).unwrap().unwrap(); + assert_eq!(folded.status, DatagenItemStatus::Running); + + let failures = datagen_failures(&[created(0), failed]).unwrap(); + assert_eq!(failures.len(), 1); + assert_eq!(failures[0].error.error_type, "ValueError"); + assert_eq!(failures[0].at.position.step.name, "check"); + } + #[test] fn sequence_collision_is_rejected() { - let created = event(0, DatagenEventType::ItemCreated); - let first = completed_step(1); + let first = leaf_completed(1, "gen", 0, Some("main")); let mut second = first.clone(); second.event_id = "different-event".to_string(); second.checkpoint_id = "different-checkpoint".to_string(); - - let error = fold_datagen_events(&[created, first, second]).unwrap_err(); + let error = fold_datagen_events(&[created(0), first, second]).unwrap_err(); assert!(error.contains("conflicting events at item_seq 1")); } #[test] - fn later_run_can_supersede_a_failure() { - let created = event(0, DatagenEventType::ItemCreated); - let mut failed = event(1, DatagenEventType::Failed); - failed.error_type = Some("RuntimeError".to_string()); - - let mut retried = event(2, DatagenEventType::ItemCreated); - retried.run_id = "run-2".to_string(); - retried.checkpoint_id = "retry-created".to_string(); - retried.event_id = datagen_event_id("item-1", "retry-created", 0); - - let folded = fold_datagen_events(&[created, failed, retried]).unwrap(); - assert_eq!(folded.status, DatagenItemStatus::Pending); - assert_eq!(folded.current_run_id, "run-2"); - assert!(folded.failure.is_none()); + fn retry_duplicate_event_is_folded_once() { + let mut append = leaf_completed(1, "gen", 1, Some("main")); + append.event_type = DatagenEventType::FieldAppend; + append.field_name = Some("messages".to_string()); + append.field_type = Some("json".to_string()); + append.codec_version = Some(1); + append.value = Some(DatagenValue::Json(json!({"role": "assistant"}))); + + let folded = + fold_datagen_events(&[created(0), append.clone(), append]).unwrap().unwrap(); + assert_eq!( + folded.fields.get("messages"), + Some(&DatagenFieldState::Appended(vec![DatagenValue::Json( + json!({"role": "assistant"}) + )])) + ); } } diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index 228a3ef..f91000f 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -35,8 +35,9 @@ use tracing::{info, warn}; use uuid::Uuid; use crate::datagen::{ - datagen_trajectory, fold_datagen_events, DatagenBlobValue, DatagenEvent, DatagenEventType, - DatagenTrajectoryPoint, DatagenValue, FoldedDatagenItem, + datagen_failures, datagen_trajectory, fold_datagen_events, DatagenBlobValue, DatagenEvent, + DatagenEventType, DatagenFailure, DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, + DatagenStepCursor, DatagenStepKind, DatagenValue, }; use crate::rollout_store::derive_shard_id; use crate::store::{column_as, column_as_optional, timestamp_from_micros}; @@ -224,23 +225,45 @@ impl DatagenStore { self.filtered_events(&filter).await } - /// Reconstruct one item's latest state exclusively from its event log. - pub async fn fold_item(&self, item_id: &str) -> LanceResult> { + /// Reconstruct one item's latest state exclusively from its event log. Returns + /// [`DatagenItemLookup::NeverStarted`] when the item has no ITEM_CREATED (the fresh-vs-resume fork). + pub async fn fold_item(&self, item_id: &str) -> LanceResult { let events = self.events_for_item(item_id).await?; - if events.is_empty() { - return Ok(None); + match fold_datagen_events(&events).map_err(invalid_input)? { + Some(item) => Ok(DatagenItemLookup::Found(item)), + None => Ok(DatagenItemLookup::NeverStarted), } - fold_datagen_events(&events) - .map(Some) - .map_err(invalid_input) } - /// Reconstruct state after every completed step without loading blob bytes. - pub async fn trajectory(&self, item_id: &str) -> LanceResult> { - let events = self.events_for_item(item_id).await?; - if events.is_empty() { - return Ok(Vec::new()); + /// Classify every root item that shares `root_item_id` with the given roots by folded lifecycle + /// status. A root not present in the log is simply absent from the result (never started). + pub async fn root_item_statuses( + &self, + root_item_ids: &[&str], + ) -> LanceResult { + let mut statuses = HashMap::new(); + for root_item_id in root_item_ids { + let events = self.events_for_root(root_item_id).await?; + let root_events: Vec = events + .into_iter() + .filter(|event| &event.item_id == root_item_id) + .collect(); + if let Some(item) = fold_datagen_events(&root_events).map_err(invalid_input)? { + statuses.insert(root_item_id.to_string(), item.status); + } } + Ok(DatagenRootItemStatuses::from_map(statuses)) + } + + /// Read all failure records for an item directly from the failure lens. + pub async fn item_failures(&self, item_id: &str) -> LanceResult> { + let events = self.events_for_item(item_id).await?; + datagen_failures(&events).map_err(invalid_input) + } + + /// Reconstruct the ordered step cursors an item passed through, without loading blob bytes. + pub async fn trajectory(&self, item_id: &str) -> LanceResult> { + let events = self.events_for_item(item_id).await?; datagen_trajectory(&events).map_err(invalid_input) } @@ -653,9 +676,10 @@ pub fn datagen_log_schema() -> Schema { Field::new("checkpoint_id", DataType::Utf8, false), Field::new("event_type", DataType::Utf8, false), Field::new("step_name", DataType::Utf8, true), + Field::new("step_kind", DataType::Utf8, true), Field::new("step_index", DataType::Int64, true), - Field::new("step_instance_id", DataType::Utf8, true), - Field::new("iteration", DataType::Int64, true), + Field::new("enclosing_step", DataType::Utf8, true), + Field::new("selector_step", DataType::Utf8, true), Field::new("attempt", DataType::Int32, false), Field::new("run_id", DataType::Utf8, false), Field::new("writer_epoch", DataType::Utf8, false), @@ -674,7 +698,7 @@ pub fn datagen_log_schema() -> Schema { Field::new("payload_size", DataType::Int64, true), Field::new("payload_checksum", DataType::Utf8, true), Field::new("query_tags_json", DataType::LargeUtf8, true), - Field::new("terminal", DataType::Utf8, true), + Field::new("status", DataType::Utf8, true), Field::new("error_type", DataType::Utf8, true), Field::new("error_dump", DataType::LargeUtf8, true), Field::new("traceback", DataType::LargeUtf8, true), @@ -773,8 +797,9 @@ fn validate_checkpoint_batch(events: &[DatagenEvent]) -> LanceResult<()> { }) { if event.step_name != completion.step_name || event.step_index != completion.step_index - || event.step_instance_id != completion.step_instance_id - || event.iteration != completion.iteration + || event.step_kind != completion.step_kind + || event.enclosing_step != completion.enclosing_step + || event.selector_step != completion.selector_step { return Err(invalid_input( "all field events must share the STEP_COMPLETED step identity", @@ -793,9 +818,10 @@ fn events_to_batch(events: &[DatagenEvent]) -> LanceResult { let mut checkpoint_id = StringBuilder::new(); let mut event_type = StringBuilder::new(); let mut step_name = StringBuilder::new(); + let mut step_kind = StringBuilder::new(); let mut step_index = Int64Builder::new(); - let mut step_instance_id = StringBuilder::new(); - let mut iteration = Int64Builder::new(); + let mut enclosing_step = StringBuilder::new(); + let mut selector_step = StringBuilder::new(); let mut attempt = Int32Builder::new(); let mut run_id = StringBuilder::new(); let mut writer_epoch = StringBuilder::new(); @@ -812,7 +838,7 @@ fn events_to_batch(events: &[DatagenEvent]) -> LanceResult { let mut payload_size = Int64Builder::new(); let mut payload_checksum = StringBuilder::new(); let mut query_tags_json = LargeStringBuilder::new(); - let mut terminal = StringBuilder::new(); + let mut status = StringBuilder::new(); let mut error_type = StringBuilder::new(); let mut error_dump = LargeStringBuilder::new(); let mut traceback = LargeStringBuilder::new(); @@ -829,9 +855,10 @@ fn events_to_batch(events: &[DatagenEvent]) -> LanceResult { checkpoint_id.append_value(&event.checkpoint_id); event_type.append_value(event.event_type.as_str()); step_name.append_option(event.step_name.as_deref()); + step_kind.append_option(event.step_kind.map(DatagenStepKind::as_str)); step_index.append_option(event.step_index); - step_instance_id.append_option(event.step_instance_id.as_deref()); - iteration.append_option(event.iteration); + enclosing_step.append_option(event.enclosing_step.as_deref()); + selector_step.append_option(event.selector_step.as_deref()); attempt.append_value(event.attempt); run_id.append_value(&event.run_id); writer_epoch.append_value(&event.writer_epoch); @@ -853,7 +880,7 @@ fn events_to_batch(events: &[DatagenEvent]) -> LanceResult { _ => None, }); value_str.append_option(match &event.value { - Some(DatagenValue::String(value)) => Some(value.as_str()), + Some(DatagenValue::Str(value)) => Some(value.as_str()), _ => None, }); match &event.value { @@ -876,7 +903,7 @@ fn events_to_batch(events: &[DatagenEvent]) -> LanceResult { Some(tags) => query_tags_json.append_value(tags.to_string()), None => query_tags_json.append_null(), } - terminal.append_option(event.terminal.map(|value| value.as_str())); + status.append_option(event.status.map(DatagenItemStatus::as_str)); error_type.append_option(event.error_type.as_deref()); error_dump.append_option(event.error_dump.as_deref()); traceback.append_option(event.traceback.as_deref()); @@ -894,9 +921,10 @@ fn events_to_batch(events: &[DatagenEvent]) -> LanceResult { Arc::new(checkpoint_id.finish()), Arc::new(event_type.finish()), Arc::new(step_name.finish()), + Arc::new(step_kind.finish()), Arc::new(step_index.finish()), - Arc::new(step_instance_id.finish()), - Arc::new(iteration.finish()), + Arc::new(enclosing_step.finish()), + Arc::new(selector_step.finish()), Arc::new(attempt.finish()), Arc::new(run_id.finish()), Arc::new(writer_epoch.finish()), @@ -913,7 +941,7 @@ fn events_to_batch(events: &[DatagenEvent]) -> LanceResult { Arc::new(payload_size.finish()), Arc::new(payload_checksum.finish()), Arc::new(query_tags_json.finish()), - Arc::new(terminal.finish()), + Arc::new(status.finish()), Arc::new(error_type.finish()), Arc::new(error_dump.finish()), Arc::new(traceback.finish()), @@ -932,9 +960,10 @@ fn batch_to_events(batch: &RecordBatch) -> LanceResult> { let checkpoint_id = column_as::(batch, "checkpoint_id")?; let event_type = column_as::(batch, "event_type")?; let step_name = column_as_optional::(batch, "step_name"); + let step_kind = column_as_optional::(batch, "step_kind"); let step_index = column_as_optional::(batch, "step_index"); - let step_instance_id = column_as_optional::(batch, "step_instance_id"); - let iteration = column_as_optional::(batch, "iteration"); + let enclosing_step = column_as_optional::(batch, "enclosing_step"); + let selector_step = column_as_optional::(batch, "selector_step"); let attempt = column_as::(batch, "attempt")?; let run_id = column_as::(batch, "run_id")?; let writer_epoch = column_as::(batch, "writer_epoch")?; @@ -951,7 +980,7 @@ fn batch_to_events(batch: &RecordBatch) -> LanceResult> { let payload_size = column_as_optional::(batch, "payload_size"); let payload_checksum = column_as_optional::(batch, "payload_checksum"); let query_tags_json = column_as_optional::(batch, "query_tags_json"); - let terminal = column_as_optional::(batch, "terminal"); + let status = column_as_optional::(batch, "status"); let error_type = column_as_optional::(batch, "error_type"); let error_dump = column_as_optional::(batch, "error_dump"); let traceback = column_as_optional::(batch, "traceback"); @@ -978,7 +1007,7 @@ fn batch_to_events(batch: &RecordBatch) -> LanceResult> { row, "value_bool", )?)), - Some("str") => Some(DatagenValue::String( + Some("str") => Some(DatagenValue::Str( optional_large_string(value_str, row) .ok_or_else(|| invalid_input("value_kind=str requires value_str"))?, )), @@ -1025,9 +1054,13 @@ fn batch_to_events(batch: &RecordBatch) -> LanceResult> { checkpoint_id: checkpoint_id.value(row).to_string(), event_type: DatagenEventType::parse(event_type.value(row)).map_err(invalid_input)?, step_name: optional_string(step_name, row), + step_kind: match optional_string(step_kind, row) { + Some(value) => Some(DatagenStepKind::parse(&value).map_err(invalid_input)?), + None => None, + }, step_index: optional_i64(step_index, row), - step_instance_id: optional_string(step_instance_id, row), - iteration: optional_i64(iteration, row), + enclosing_step: optional_string(enclosing_step, row), + selector_step: optional_string(selector_step, row), attempt: attempt.value(row), run_id: run_id.value(row).to_string(), writer_epoch: writer_epoch.value(row).to_string(), @@ -1036,10 +1069,8 @@ fn batch_to_events(batch: &RecordBatch) -> LanceResult> { codec_version: optional_i32(codec_version, row), value, query_tags, - terminal: match optional_string(terminal, row) { - Some(value) => { - Some(crate::datagen::DatagenTerminal::parse(&value).map_err(invalid_input)?) - } + status: match optional_string(status, row) { + Some(value) => Some(DatagenItemStatus::parse(&value).map_err(invalid_input)?), None => None, }, error_type: optional_string(error_type, row), @@ -1131,7 +1162,7 @@ fn is_fenced_error(error: &LanceError) -> bool { mod tests { use super::*; use crate::datagen::{ - datagen_event_id, DatagenFieldState, DatagenItemStatus, DatagenTerminal, + datagen_event_id, DatagenFieldState, DatagenItemStatus, DatagenStepKind, DATAGEN_SCHEMA_VERSION, }; use chrono::{TimeZone, Utc}; @@ -1154,9 +1185,10 @@ mod tests { checkpoint_id: checkpoint_id.to_string(), event_type, step_name: None, + step_kind: None, step_index: None, - step_instance_id: None, - iteration: None, + enclosing_step: None, + selector_step: None, attempt: 0, run_id: "run-1".to_string(), writer_epoch: "writer-1".to_string(), @@ -1165,7 +1197,7 @@ mod tests { codec_version: None, value: None, query_tags: None, - terminal: None, + status: None, error_type: None, error_dump: None, traceback: None, @@ -1184,8 +1216,8 @@ mod tests { ) -> DatagenEvent { let mut event = event("item-1", seq, "grade-0", ordinal, event_type); event.step_name = Some("grade".to_string()); + event.step_kind = Some(DatagenStepKind::Leaf); event.step_index = Some(2); - event.step_instance_id = Some("root/grade/0".to_string()); event.field_name = Some(field_name.to_string()); event.field_type = Some(field_type.to_string()); event.codec_version = Some(1); @@ -1202,8 +1234,8 @@ mod tests { DatagenEventType::StepCompleted, ); event.step_name = Some("grade".to_string()); + event.step_kind = Some(DatagenStepKind::Leaf); event.step_index = Some(2); - event.step_instance_id = Some("root/grade/0".to_string()); event } @@ -1256,7 +1288,7 @@ mod tests { store.append_checkpoint(&checkpoint).await.unwrap(); let mut terminal = event("item-1", 5, "terminal", 0, DatagenEventType::Terminal); - terminal.terminal = Some(DatagenTerminal::Completed); + terminal.status = Some(DatagenItemStatus::Completed); store.append(&[terminal]).await.unwrap(); let events = store.events_for_item("item-1").await.unwrap(); @@ -1275,18 +1307,19 @@ mod tests { Some(blob_bytes) ); - let folded = store.fold_item("item-1").await.unwrap().unwrap(); + let folded = store.fold_item("item-1").await.unwrap(); + let folded = folded.folded().expect("item-1 was created"); assert_eq!(folded.status, DatagenItemStatus::Completed); assert_eq!( folded.fields.get("score"), Some(&DatagenFieldState::Set(DatagenValue::Int(i64::MAX))) ); - assert_eq!(folded.completed_steps.len(), 1); + assert_eq!(folded.trajectory.ordered.len(), 1); assert_eq!(folded.query_tags, Some(json!({"domain": "math"}))); let trajectory = store.trajectory("item-1").await.unwrap(); assert_eq!(trajectory.len(), 1); - assert_eq!(trajectory[0].cursor.step_name, "grade"); + assert_eq!(trajectory[0].position.step.name, "grade"); }); } @@ -1297,13 +1330,16 @@ mod tests { let runtime = tokio::runtime::Runtime::new().unwrap(); runtime.block_on(async { let mut store = DatagenStore::open(&uri).await.unwrap(); - let no_op = completed_step(0, 0); + let created = event("item-1", 0, "created", 0, DatagenEventType::ItemCreated); + store.append(&[created]).await.unwrap(); + let no_op = completed_step(1, 0); store.append_checkpoint(&[no_op]).await.unwrap(); - let folded = store.fold_item("item-1").await.unwrap().unwrap(); - assert_eq!(folded.completed_steps.len(), 1); + let folded = store.fold_item("item-1").await.unwrap(); + let folded = folded.folded().expect("item-1 was created"); + assert_eq!(folded.trajectory.ordered.len(), 1); let field_only = field_event( - 1, + 2, 0, DatagenEventType::FieldSet, "score", @@ -1360,6 +1396,9 @@ mod tests { let mut failed = event("root-b", 0, "failed-b", 0, DatagenEventType::Failed); failed.run_id = "run-failed".to_string(); failed.writer_epoch = "writer-b".to_string(); + failed.step_name = Some("expand".to_string()); + failed.step_kind = Some(DatagenStepKind::Leaf); + failed.step_index = Some(0); failed.error_type = Some("ValueError".to_string()); failed.error_dump = Some("bad source item".to_string()); writer_b.append(&[failed]).await.unwrap(); diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index 904e632..c4f0c60 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -20,10 +20,11 @@ mod store; pub use api_impl::rollout_record_to_dto; pub use context::{Context, ContextEntry, Snapshot}; pub use datagen::{ - datagen_event_id, datagen_trajectory, fold_datagen_events, DatagenBlobValue, DatagenEvent, - DatagenEventType, DatagenFailure, DatagenFieldState, DatagenItemStatus, DatagenStepCursor, - DatagenTerminal, DatagenTrajectoryPoint, DatagenValue, FoldedDatagenItem, - DATAGEN_SCHEMA_VERSION, + 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, }; pub use datagen_store::{datagen_log_schema, DatagenStore, DatagenStoreOptions}; pub use eval::{ From 91b4981b56e63891d71b60ce8b47b83718eb922e Mon Sep 17 00:00:00 2001 From: YangjunZ Date: Sun, 26 Jul 2026 17:27:54 -0700 Subject: [PATCH 2/3] 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.rs | 118 +++++++++ .../lance-context-core/src/datagen_store.rs | 59 ++++- docs/design/using-datagen-store.md | 235 ++++++++++++++++++ specs/datagen-checkpoint-schema.md | 73 ++++-- 4 files changed, 462 insertions(+), 23 deletions(-) create mode 100644 docs/design/using-datagen-store.md diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index eb0350c..8dd2061 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -950,4 +950,122 @@ mod tests { )])) ); } + + #[test] + fn step_kind_and_status_parse_round_trip() { + for kind in [ + DatagenStepKind::Root, + DatagenStepKind::Leaf, + DatagenStepKind::Sequence, + DatagenStepKind::Loop, + DatagenStepKind::MapReduce, + DatagenStepKind::Branch, + DatagenStepKind::SubPipeline, + DatagenStepKind::Conditional, + DatagenStepKind::Router, + ] { + assert_eq!(DatagenStepKind::parse(kind.as_str()).unwrap(), kind); + } + assert!(DatagenStepKind::Sequence.is_driver()); + assert!(DatagenStepKind::Loop.is_driver()); + assert!(!DatagenStepKind::MapReduce.is_driver()); + assert!(!DatagenStepKind::Leaf.is_driver()); + assert!(DatagenStepKind::parse("nope").is_err()); + + for status in [ + DatagenItemStatus::Running, + DatagenItemStatus::Completed, + DatagenItemStatus::Filtered, + DatagenItemStatus::Failed, + ] { + assert_eq!(DatagenItemStatus::parse(status.as_str()).unwrap(), status); + } + assert!(DatagenItemStatus::parse("nope").is_err()); + } + + #[test] + fn terminal_filtered_sets_filtered_status() { + let mut terminal = event(2, DatagenEventType::Terminal); + terminal.status = Some(DatagenItemStatus::Filtered); + let folded = fold_datagen_events(&[created(0), terminal]).unwrap().unwrap(); + assert_eq!(folded.status, DatagenItemStatus::Filtered); + } + + #[test] + fn selector_step_is_folded_onto_the_chosen_child_position() { + // A Conditional/Router writes no row; the chosen leaf records who selected it. + let mut chosen = leaf_completed(1, "stage2_qa", 0, Some("rubric_generation")); + chosen.selector_step = Some("if_stage2_qa".to_string()); + let folded = fold_datagen_events(&[created(0), chosen]).unwrap().unwrap(); + let cursor = &folded.trajectory.ordered[0]; + assert_eq!(cursor.position.selector.as_deref(), Some("if_stage2_qa")); + assert_eq!(cursor.position.step.name, "stage2_qa"); + } + + #[test] + fn fan_out_sub_item_folds_with_lineage() { + // A MapReduce/Branch sub-item is its own stream carrying denormalized parent/root ids. + let root = DatagenItemId::from_source_key("5"); + let child = root.child("solve_twice", 1); + let mut created = event(0, DatagenEventType::ItemCreated); + created.item_id = child.to_string(); + created.parent_item_id = Some(root.to_string()); + created.status = Some(DatagenItemStatus::Running); + created.event_id = datagen_event_id(&child.to_string(), "created", 0); + + let mut solved = leaf_completed(1, "solve", 0, Some("solve_attempt")); + solved.item_id = child.to_string(); + solved.root_item_id = root.to_string(); + solved.parent_item_id = Some(root.to_string()); + solved.event_id = datagen_event_id(&child.to_string(), "solve-0", 0); + + let folded = fold_datagen_events(&[created, solved]).unwrap().unwrap(); + assert_eq!(folded.item_id, child); + assert_eq!(folded.root_item_id, root); + assert_eq!(folded.parent_item_id, Some(root)); + assert_eq!(folded.item_id.origin_step(), Some("solve_twice")); + assert_eq!(folded.item_id.branch_idx(), Some(1)); + } + + #[test] + fn field_rejects_mixing_set_and_append() { + let mut set = leaf_completed(1, "gen", 0, Some("main")); + set.event_type = DatagenEventType::FieldSet; + set.field_name = Some("draft".to_string()); + set.field_type = Some("str".to_string()); + set.codec_version = Some(1); + set.value = Some(DatagenValue::Str("v1".to_string())); + + let mut append = set.clone(); + append.event_type = DatagenEventType::FieldAppend; + append.item_seq = 2; + append.checkpoint_id = "c2".to_string(); + append.event_id = datagen_event_id("5", "c2", 0); + + let error = fold_datagen_events(&[created(0), set, append]).unwrap_err(); + assert!(error.contains("mixes FIELD_SET and FIELD_APPEND")); + } + + #[test] + fn resume_open_frame_is_started_minus_completed() { + // A driver frame that opened but never completed = the frame live at crash time. + let mut main_completed = driver_started(3, "main", 0, None); + main_completed.event_type = DatagenEventType::StepCompleted; + main_completed.checkpoint_id = "main-done".to_string(); + main_completed.event_id = datagen_event_id("5", "main-done", 0); + let events = [ + created(0), + driver_started(1, "main", 0, None), + driver_started(2, "refine", 1, Some("main")), + main_completed, + ]; + let folded = fold_datagen_events(&events).unwrap().unwrap(); + let open: Vec<_> = folded + .trajectory + .started + .difference(&folded.trajectory.completed) + .collect(); + assert_eq!(open.len(), 1); + assert_eq!(open[0].step.name, "refine"); + } } diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index f91000f..6f36e4c 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -1162,8 +1162,8 @@ fn is_fenced_error(error: &LanceError) -> bool { mod tests { use super::*; use crate::datagen::{ - datagen_event_id, DatagenFieldState, DatagenItemStatus, DatagenStepKind, - DATAGEN_SCHEMA_VERSION, + datagen_event_id, DatagenFieldState, DatagenItemId, DatagenItemLookup, DatagenItemStatus, + DatagenStepKind, DATAGEN_SCHEMA_VERSION, }; use chrono::{TimeZone, Utc}; use serde_json::json; @@ -1455,4 +1455,59 @@ mod tests { ); }); } + + #[test] + fn fan_out_tree_reads_by_root_and_classifies_status() { + 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(); + + // Root item "7" fans out into one sub-item "7/solve_twice:0". + let mut root_created = event("7", 0, "created-root", 0, DatagenEventType::ItemCreated); + root_created.status = Some(DatagenItemStatus::Running); + store.append(&[root_created]).await.unwrap(); + + let mut child_created = event( + "7/solve_twice:0", + 0, + "created-child", + 0, + DatagenEventType::ItemCreated, + ); + child_created.parent_item_id = Some("7".to_string()); + child_created.status = Some(DatagenItemStatus::Running); + store.append(&[child_created]).await.unwrap(); + + let mut root_terminal = event("7", 1, "terminal", 0, DatagenEventType::Terminal); + root_terminal.status = Some(DatagenItemStatus::Completed); + store.append(&[root_terminal]).await.unwrap(); + + // The whole tree is one root filter, no join. + let tree = store.events_for_root("7").await.unwrap(); + assert_eq!(tree.len(), 3); + + // Bulk classification sees the root as terminated; the child is still running. + let statuses = store.root_item_statuses(&["7"]).await.unwrap(); + let root_id = DatagenItemId::from_source_key("7"); + assert!(statuses.is_terminated(&root_id)); + + let child = store.fold_item("7/solve_twice:0").await.unwrap(); + assert_eq!( + child.folded().unwrap().status, + DatagenItemStatus::Running + ); + assert_eq!( + child.folded().unwrap().parent_item_id, + Some(DatagenItemId::from_source_key("7")) + ); + + // A never-started sibling folds to NeverStarted. + assert_eq!( + store.fold_item("7/solve_twice:1").await.unwrap(), + DatagenItemLookup::NeverStarted + ); + }); + } } diff --git a/docs/design/using-datagen-store.md b/docs/design/using-datagen-store.md new file mode 100644 index 0000000..82506fd --- /dev/null +++ b/docs/design/using-datagen-store.md @@ -0,0 +1,235 @@ +# Using the Datagen Checkpoint Store + +`DatagenStore` is the durable checkpoint backend for a datagen experiment. It is +a single append-only Lance log: every write is one immutable event, and the +current state of any item is the *fold* of its events. Nothing is ever updated or +deleted. This document shows how a client (the `mai_datagen` executor, through +its pyo3 bindings) drives the store through a run — create, checkpoint, resume, +read — with runnable Rust snippets. + +For the schema and the design rationale, see +[`specs/datagen-checkpoint-schema.md`](../../specs/datagen-checkpoint-schema.md). + +## The model in one paragraph + +An experiment is one `log.lance`. Each row is a `DatagenEvent` with an +`event_type`. An **item** is one stream of events sharing an `item_id`; a source +item is a **root**, and fan-out steps project **sub-items**, each its own stream. +To learn an item's state you read its events and call `fold_datagen_events`, +which replays them into a `FoldedDatagenItem` (fields, status, trajectory). +Because reconstruction is a pure fold over an append-only log, a crash can never +leave a torn write: an unacknowledged batch either persisted whole or not at all. + +## Opening a store + +```rust +use lance_context_core::{DatagenStore, DatagenStoreOptions}; + +// One writer per shard. Concurrent writers pass distinct shard ids. +let mut store = DatagenStore::open("s3://bucket/exp/log.lance").await?; + +// Or with an explicit shard id (each writer owns its own MemWAL shard): +let mut store = DatagenStore::open_with_options( + "s3://bucket/exp/log.lance", + DatagenStoreOptions { shard_id: Some("worker-3".into()), ..Default::default() }, +).await?; +``` + +## Item identity — composed on the client, no round-trip + +Ids are structured values (`DatagenItemId`) stored as materialized path strings. +The client composes them purely; the store never allocates an id. + +```rust +use lance_context_core::DatagenItemId; + +let root = DatagenItemId::from_source_key("5"); // "5" +let child = root.child("solve_twice", 0); // "5/solve_twice:0" +let grandchild = child.child("judge", 1); // "5/solve_twice:0/judge:1" + +assert_eq!(grandchild.origin_step(), Some("judge")); +assert_eq!(grandchild.branch_idx(), Some(1)); +assert_eq!(grandchild.parent().unwrap(), child); +assert_eq!(grandchild.root(), root); +``` + +Because `root_item_id` and `parent_item_id` are denormalized onto every row, +reading a whole item tree is one filter (`events_for_root`), never a join. + +## Writing events + +An event is built with its provenance columns and appended. Field events plus +their `STEP_COMPLETED` marker for one step must go in a **single** call so a crash +cannot expose a half-checkpointed step — the batch is one durable generation. + +```rust +use lance_context_core::{ + datagen_event_id, DatagenEvent, DatagenEventType, DatagenItemStatus, + DatagenStepKind, DatagenValue, DATAGEN_SCHEMA_VERSION, +}; + +// 1. Announce the item once. +let created = DatagenEvent { + event_id: datagen_event_id("5", "created", 0), + item_id: "5".into(), + root_item_id: "5".into(), + parent_item_id: None, + item_seq: 0, + checkpoint_id: "created".into(), + event_type: DatagenEventType::ItemCreated, + status: Some(DatagenItemStatus::Running), + schema_version: DATAGEN_SCHEMA_VERSION, + ..blank_event() // your helper that zero-fills the optional columns +}; +store.append(&[created]).await?; + +// 2. Checkpoint one step: its field delta + exactly one STEP_COMPLETED, atomically. +let score = DatagenEvent { + event_id: datagen_event_id("5", "solve-0", 0), + item_id: "5".into(), + root_item_id: "5".into(), + item_seq: 1, + checkpoint_id: "solve-0".into(), + event_type: DatagenEventType::FieldSet, + step_name: Some("solve".into()), + step_kind: Some(DatagenStepKind::Leaf), + step_index: Some(0), + enclosing_step: Some("solve_attempt".into()), + field_name: Some("score".into()), + field_type: Some("int".into()), + codec_version: Some(1), + value: Some(DatagenValue::Int(9)), + schema_version: DATAGEN_SCHEMA_VERSION, + ..blank_event() +}; +let completed = DatagenEvent { + event_id: datagen_event_id("5", "solve-0", 1), + item_seq: 2, + checkpoint_id: "solve-0".into(), + event_type: DatagenEventType::StepCompleted, + // same step_name / step_kind / step_index / enclosing_step as above + ..score.clone() +}; +store.append_checkpoint(&[score, completed]).await?; +``` + +`append_checkpoint` enforces "exactly one `STEP_COMPLETED` per batch"; +`append` is the lower-level form used for lifecycle events (`ITEM_CREATED`, +`TERMINAL`, `FAILED`). + +### Retries are safe + +`event_id` is deterministic (`datagen_event_id(item_id, checkpoint_id, ordinal)`). +Replaying an ambiguously-acknowledged batch writes the same ids, and the fold +de-duplicates them — the second `append_checkpoint` of an identical batch is a +no-op in the folded result. + +### Kind determines what you write + +| Step kind | Emits | +|---|---| +| `Sequence`, `Loop` (drivers) | `STEP_STARTED` frame + `STEP_COMPLETED` | +| `MapReduce`, `Branch`, `SubPipeline` (fan-out) | only the reduce `STEP_COMPLETED`; sub-items are their own streams | +| `Conditional`, `Router` (selectors) | no row of their own — the chosen child sets `selector_step` | +| `Leaf` | `FIELD_SET` / `FIELD_APPEND` + `STEP_COMPLETED` | + +## Finishing an item + +```rust +// Success or filtered-out — a lifecycle terminal. +let terminal = DatagenEvent { + event_type: DatagenEventType::Terminal, + status: Some(DatagenItemStatus::Completed), // or Filtered + ..lifecycle_event("5", 3, "terminal") +}; +store.append(&[terminal]).await?; + +// A raised step — a failure-lens row. The item still folds to `running`. +let failed = DatagenEvent { + event_type: DatagenEventType::Failed, + status: Some(DatagenItemStatus::Failed), + step_name: Some("score".into()), + step_kind: Some(DatagenStepKind::Leaf), + step_index: Some(1), + error_type: Some("ValueError".into()), + ..lifecycle_event("5", 3, "failed") +}; +store.append(&[failed]).await?; +``` + +## Reading + +### Fold one item + +```rust +use lance_context_core::DatagenItemLookup; + +match store.fold_item("5").await? { + DatagenItemLookup::NeverStarted => { /* fresh: process from scratch */ } + DatagenItemLookup::Found(item) => { + // item.status, item.fields, item.trajectory, item.last_item_seq, item.last_attempt + // Resume: continue writing at last_item_seq + 1, attempt last_attempt + 1. + } +} +``` + +`NeverStarted` (no `ITEM_CREATED`) is the explicit fresh-vs-resume fork. A +`Found` item's `status` is always a *lifecycle* value (`running` / `completed` +/ `filtered`) — `failed` never appears here. + +### Resume: the open frame + +The trajectory records both started and completed positions. `started \ completed` +is the driver frame that was open when the process died — everything already +completed is skipped on re-run. + +```rust +let item = store.fold_item("5").await?.folded().unwrap(); +let open: Vec<_> = item.trajectory.started + .difference(&item.trajectory.completed) + .collect(); // the frame(s) to re-enter +``` + +### Whole item tree + +```rust +// Every event under root "5", including all fan-out sub-items — one filter. +let events = store.events_for_root("5").await?; +``` + +### Bulk startup classification + +```rust +// Classify many roots at once without folding fields (skip / resume / fresh). +let statuses = store.root_item_statuses(&["5", "6", "7"]).await?; +assert!(statuses.is_terminated(&DatagenItemId::from_source_key("5"))); +``` + +### Failures (the failure lens) + +```rust +let failures = store.item_failures("5").await?; // 0..N, across attempts +for failure in &failures { + println!("{} at {}", failure.error.error_type, failure.at.position.step.name); +} +let all = store.failures(Some("run-1")).await?; // run-wide forensics +``` + +### Blobs are lazy + +Field values that are blobs fold to a lazy `DatagenBlobValue { bytes: None, .. }`. +Materialize the bytes only when needed: + +```rust +let bytes = store.get_blob(&blob_event_id).await?; // O(single blob) take_rows +``` + +## Concurrency and maintenance + +- One live owner writes a given item at a time; a new owner takes over with a new + `writer_epoch` (fencing). +- Concurrent writers use distinct MemWAL shards. Reads union the base table and + all flushed shards, so every instance sees every writer's events. +- Each writer periodically merges only its own generations into the base table + (`cleanup_own_shard` / `spawn_periodic_cleanup`). Shared base-table compaction + is scheduled by one elected maintenance worker per experiment. diff --git a/specs/datagen-checkpoint-schema.md b/specs/datagen-checkpoint-schema.md index d999ee7..f045f8f 100644 --- a/specs/datagen-checkpoint-schema.md +++ b/specs/datagen-checkpoint-schema.md @@ -18,39 +18,62 @@ rebuildable, and excluded from checkpoint correctness. ## Event model -Every row is an immutable event: +Every row is an immutable event. The current schema version is +`DATAGEN_SCHEMA_VERSION = 2`. There are seven event types: -- `ITEM_CREATED` -- `FIELD_SET` -- `FIELD_APPEND` -- `STEP_COMPLETED` -- `FAILED` -- `TERMINAL` +- `ITEM_CREATED` — an item (stream) first appears, exactly once, `status = running`. +- `STEP_STARTED` — a driver frame (`Sequence`/`Loop`) opened. A structural marker. +- `FIELD_SET` — replaces a field's value (fold: last-writer-wins). +- `FIELD_APPEND` — accumulates onto a field (fold: append in order). +- `STEP_COMPLETED` — a checkpointed step boundary, carrying that step's field delta. +- `FAILED` — a step raised; a failure-lens row that leaves the item `running`. +- `TERMINAL` — the item finished, `status = completed | filtered`. Every completed step emits a `STEP_COMPLETED` event, even when no field changed. All field events and the completion marker for one step are written in the same -checkpoint batch. +checkpoint batch. Only driver kinds emit `STEP_STARTED`; fan-out kinds +(`MapReduce`/`Branch`/`SubPipeline`) emit only their reduce `STEP_COMPLETED`; +selector kinds (`Conditional`/`Router`) emit no row of their own — the chosen +child records the selector in `selector_step`. `event_id` is a deterministic idempotency key. Retrying an ambiguously acknowledged batch writes the same event ids, and the MemWAL LSM read path de-duplicates them. `item_seq` is strictly increasing per item; two different events at the same sequence are treated as a writer-fencing violation. +### Item identity + +An item id is a materialized path string owned by the store. A root id is the +executor's source key (`5`); a fan-out sub-item extends its parent with one +`step:idx` segment (`5/solve_twice:0`). Ids compose purely on the client, with +no store round-trip. `root_item_id` / `parent_item_id` are denormalized on every +row, so any subtree is a single filter with no join. + +### Read lenses + +State is read under two lenses: + +- *lifecycle* — fold the events; `status` is `running` until a `TERMINAL`. Drives + skip/resume/fresh classification and ignores `FAILED`. +- *failure* — read `FAILED` rows directly (forensics). An item may have 0..N + failures across attempts while still folding to `running` under lifecycle. + ## Schema | Column | Type | Purpose | |---|---|---| | `event_id` | string | Deterministic event identity and LSM primary key | -| `item_id` | string | Scoped item identity | -| `root_item_id` | string | Root of the projected item tree | -| `parent_item_id` | string? | Direct parent item | +| `item_id` | string | Scoped item identity (materialized path) | +| `root_item_id` | string | Root of the projected item tree (denormalized) | +| `parent_item_id` | string? | Direct parent item (denormalized) | | `item_seq` | int64 | Per-item event ordering | | `checkpoint_id` | string | Atomic step-boundary identity | -| `event_type` | string | Event kind | -| `step_name` | string? | Step provenance | -| `step_index` | int64? | Static step position | -| `step_instance_id` | string? | Runtime step identity | -| `iteration` | int64? | Loop/branch iteration | +| `event_type` | string | Event kind (one of the seven above) | +| `step_name` | string? | Step provenance (globally-unique step name) | +| `step_kind` | string? | Composition kind (`sequence`, `loop`, `map_reduce`, `branch`, `sub_pipeline`, `conditional`, `router`, `leaf`, `root`) | +| `step_index` | int64? | Static step position (loop iteration = driver frame's index) | +| `enclosing_step` | string? | Name of the `Sequence`/`Loop` driver frame this step ran under | +| `selector_step` | string? | Name of the `Conditional`/`Router` that chose this step | | `attempt` | int32 | Execution attempt | | `run_id` | string | Run attribution | | `writer_epoch` | string | Item ownership/fencing identity | @@ -67,13 +90,19 @@ events at the same sequence are treated as a writer-fencing violation. | `payload_size` | int64? | Blob size | | `payload_checksum` | string? | Blob integrity | | `query_tags_json` | large_string? | Non-authoritative query tags | -| `terminal` | string? | `completed` or `filtered` | +| `status` | string? | `running` / `completed` / `filtered` / `failed` (on ITEM_CREATED / TERMINAL / FAILED) | | `error_type` | string? | Failure type | | `error_dump` | large_string? | Serialized failure | | `traceback` | large_string? | Failure traceback | | `event_ts` | timestamp(us, UTC) | Event time | | `schema_version` | int32 | Log schema compatibility version | +The `status` column replaces the earlier `terminal` column: it carries a stored +`running` (from `ITEM_CREATED`) rather than deriving it, plus `failed` for the +failure lens. Structured `step_kind` / `enclosing_step` / `selector_step` replace +the earlier opaque `step_instance_id` + `iteration` provenance, so a resume can +rebuild step coordinates instead of trusting a stored cursor blob. + `value_blob` remains inline while MemWAL's LSM scanner cannot materialize blob-v2 columns. Normal fold and trajectory reads project it out. Blob access first locates `event_id` using lightweight columns and then calls `take_rows` @@ -85,9 +114,12 @@ for the exact `_rowid`. 2. Log writes are append-only; retries reuse deterministic event ids. 3. One live owner writes a given item at a time. Ownership changes require a new `writer_epoch`. -4. Resume reads `item_id = X`, orders by `item_seq`, and folds the events. -5. `FIELD_SET` replaces a field; `FIELD_APPEND` accumulates it. -6. `STEP_COMPLETED` reconstructs the resume cursor. +4. Resume reads `item_id = X`, orders by `item_seq`, and folds the events. An + item with no `ITEM_CREATED` folds to `NeverStarted` (fresh vs. restore fork). +5. `FIELD_SET` replaces a field; `FIELD_APPEND` accumulates it. Mixing the two + on one field is rejected. +6. `STEP_COMPLETED` reconstructs the resume cursor; `STEP_STARTED \ STEP_COMPLETED` + (started minus completed) is the driver frame open at crash time. 7. `FAILED` and `TERMINAL` are ordinary log events, not separate datasets. 8. Blob bytes are loaded only when the corresponding lazy reference is used. @@ -98,4 +130,3 @@ shards, so any instance sees every writer's events. Each writer periodically merges only its own generations into the base table. Shared base-table compaction and index refresh must be scheduled by one elected maintenance worker per experiment. - From 9a39772d1ad342ddc5a3faab225d37ec32f4a734 Mon Sep 17 00:00:00 2001 From: Yangjun Zhang Date: Sun, 26 Jul 2026 19:52:51 -0700 Subject: [PATCH 3/3] 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 | 36 +++++++++++-------- .../lance-context-core/src/datagen_store.rs | 9 ++--- crates/lance-context-core/src/lib.rs | 6 ++-- crates/lance-context/src/lib.rs | 2 +- 4 files changed, 29 insertions(+), 24 deletions(-) diff --git a/crates/lance-context-core/src/datagen.rs b/crates/lance-context-core/src/datagen.rs index 8dd2061..88f2fa6 100644 --- a/crates/lance-context-core/src/datagen.rs +++ b/crates/lance-context-core/src/datagen.rs @@ -195,10 +195,7 @@ impl DatagenItemId { .map_err(|_| format!("non-integer branch index '{idx}' in '{path}'"))?; segments.push((step.to_string(), branch_idx)); } - Ok(Self { - root_key, - segments, - }) + Ok(Self { root_key, segments }) } } @@ -382,6 +379,7 @@ pub struct FoldedDatagenItem { /// 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)] +#[allow(clippy::large_enum_variant)] // `Found` is the common path; boxing it adds indirection to the hot case. pub enum DatagenItemLookup { NeverStarted, Found(FoldedDatagenItem), @@ -434,7 +432,6 @@ impl DatagenRootItemStatuses { self.inner.is_empty() } - #[must_use] pub fn iter(&self) -> impl Iterator { self.inner.iter() } @@ -685,7 +682,8 @@ fn apply_event(item: &mut FoldedDatagenItem, event: &DatagenEvent) -> Result<(), let field_name = event.field_name.clone().unwrap(); let value = event.value.clone().unwrap(); record_blob_event_id(item, &field_name, &value, &event.event_id); - item.fields.insert(field_name, DatagenFieldState::Set(value)); + item.fields + .insert(field_name, DatagenFieldState::Set(value)); } DatagenEventType::FieldAppend => { let field_name = event.field_name.clone().unwrap(); @@ -719,7 +717,9 @@ fn apply_event(item: &mut FoldedDatagenItem, event: &DatagenEvent) -> Result<(), } DatagenEventType::Terminal => { item.status = match event.status { - Some(status @ (DatagenItemStatus::Completed | DatagenItemStatus::Filtered)) => status, + Some(status @ (DatagenItemStatus::Completed | DatagenItemStatus::Filtered)) => { + status + } _ => return Err("TERMINAL event missing completed/filtered status".to_string()), }; } @@ -882,8 +882,9 @@ mod tests { append2.event_id = datagen_event_id("5", "c4", 0); append2.value = Some(DatagenValue::Json(json!({"n": "b"}))); - let folded = - fold_datagen_events(&[created(0), set_v1, set_v2, append, append2]).unwrap().unwrap(); + let folded = fold_datagen_events(&[created(0), set_v1, set_v2, append, append2]) + .unwrap() + .unwrap(); assert_eq!( folded.fields.get("draft"), Some(&DatagenFieldState::Set(DatagenValue::Str("v2".to_string()))) @@ -901,7 +902,9 @@ mod tests { fn terminal_sets_lifecycle_status() { let mut terminal = event(2, DatagenEventType::Terminal); terminal.status = Some(DatagenItemStatus::Completed); - let folded = fold_datagen_events(&[created(0), terminal]).unwrap().unwrap(); + let folded = fold_datagen_events(&[created(0), terminal]) + .unwrap() + .unwrap(); assert_eq!(folded.status, DatagenItemStatus::Completed); } @@ -913,7 +916,9 @@ mod tests { failed.error_type = Some("ValueError".to_string()); failed.attempt = 0; - let folded = fold_datagen_events(&[created(0), failed.clone()]).unwrap().unwrap(); + let folded = fold_datagen_events(&[created(0), failed.clone()]) + .unwrap() + .unwrap(); assert_eq!(folded.status, DatagenItemStatus::Running); let failures = datagen_failures(&[created(0), failed]).unwrap(); @@ -941,8 +946,9 @@ mod tests { append.codec_version = Some(1); append.value = Some(DatagenValue::Json(json!({"role": "assistant"}))); - let folded = - fold_datagen_events(&[created(0), append.clone(), append]).unwrap().unwrap(); + let folded = fold_datagen_events(&[created(0), append.clone(), append]) + .unwrap() + .unwrap(); assert_eq!( folded.fields.get("messages"), Some(&DatagenFieldState::Appended(vec![DatagenValue::Json( @@ -987,7 +993,9 @@ mod tests { fn terminal_filtered_sets_filtered_status() { let mut terminal = event(2, DatagenEventType::Terminal); terminal.status = Some(DatagenItemStatus::Filtered); - let folded = fold_datagen_events(&[created(0), terminal]).unwrap().unwrap(); + let folded = fold_datagen_events(&[created(0), terminal]) + .unwrap() + .unwrap(); assert_eq!(folded.status, DatagenItemStatus::Filtered); } diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index 0ce83f3..7c66a57 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -36,8 +36,8 @@ use uuid::Uuid; use crate::datagen::{ datagen_failures, datagen_trajectory, fold_datagen_events, DatagenBlobValue, DatagenEvent, - DatagenEventType, DatagenFailure, DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, - DatagenStepCursor, DatagenStepKind, DatagenValue, + DatagenEventType, DatagenFailure, DatagenItemLookup, DatagenItemStatus, + DatagenRootItemStatuses, DatagenStepCursor, DatagenStepKind, DatagenValue, }; use crate::rollout_store::{align_batch_to_schema, derive_shard_id, is_not_found_error}; use crate::store::{column_as, column_as_optional, timestamp_from_micros}; @@ -1566,10 +1566,7 @@ mod tests { assert!(statuses.is_terminated(&root_id)); let child = store.fold_item("7/solve_twice:0").await.unwrap(); - assert_eq!( - child.folded().unwrap().status, - DatagenItemStatus::Running - ); + assert_eq!(child.folded().unwrap().status, DatagenItemStatus::Running); assert_eq!( child.folded().unwrap().parent_item_id, Some(DatagenItemId::from_source_key("7")) diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index 8c6cbee..44f041b 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -23,9 +23,9 @@ 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, + DatagenItemId, DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, + DatagenStepCursor, DatagenStepId, DatagenStepKind, DatagenStreamPosition, DatagenTerminal, + DatagenTrajectory, DatagenValue, FoldedDatagenItem, DATAGEN_SCHEMA_VERSION, }; pub use datagen_store::{datagen_log_schema, DatagenStore, DatagenStoreOptions}; pub use eval::{ diff --git a/crates/lance-context/src/lib.rs b/crates/lance-context/src/lib.rs index 2e4a4f4..93d6649 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -7,7 +7,7 @@ pub use lance_context_core::{ CompactionConfig, CompactionMetrics, CompactionStats, Context, ContextEntry, ContextNamespace, ContextRecord, ContextStoreOptions, DatagenBlobValue, DatagenEvent, DatagenEventType, DatagenFailure, DatagenFieldState, DatagenItemStatus, DatagenStepCursor, DatagenStore, - DatagenStoreOptions, DatagenTerminal, DatagenTrajectoryPoint, DatagenValue, FoldedDatagenItem, + DatagenStoreOptions, DatagenTerminal, DatagenTrajectory, DatagenValue, FoldedDatagenItem, IdIndexType, LifecycleQueryOptions, MetadataFilter, PartitionInfo, PartitionSelector, PartitionSpec, RecordFilters, Relationship, RetrieveResult, RolloutFilters, RolloutRecord, SearchResult, Snapshot, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE,