diff --git a/crates/lance-context-api/src/lib.rs b/crates/lance-context-api/src/lib.rs index b696982..1fdb00b 100644 --- a/crates/lance-context-api/src/lib.rs +++ b/crates/lance-context-api/src/lib.rs @@ -178,6 +178,60 @@ pub trait RolloutStoreApi { fn checkout(&mut self, version: u64) -> impl Future> + Send; } +// --------------------------------------------------------------------------- +// Datagen trait +// --------------------------------------------------------------------------- + +/// Remote-capable surface of a datagen checkpoint store — the append-only +/// delta-log of item lifecycle / field events plus the folded read lenses. +/// +/// A datagen store has no upsert/search/compaction: writers only append events +/// (`append` for lifecycle rows, `append_checkpoint` for an atomic step +/// boundary), and readers fold an item's events into latest state +/// (`fold_item`), classify root items in bulk (`root_item_statuses`), list an +/// item's failure pointers (`item_failures`), or materialize one field's +/// offloaded blob bytes (`get_blob`). +/// +/// Field blob bytes (`DatagenValueDto::bytes`) travel inline as base64 in JSON +/// or, for large payloads on the append endpoints, as raw multipart parts. By +/// the time a call reaches this trait the bytes are already materialized in +/// memory, so the signatures are transport-agnostic. +pub trait DatagenStoreApi { + fn append( + &mut self, + events: &[DatagenEventDto], + ) -> impl Future> + Send; + + fn append_checkpoint( + &mut self, + events: &[DatagenEventDto], + ) -> impl Future> + Send; + + fn fold_item( + &self, + item_id: &str, + ) -> impl Future>> + Send; + + fn root_item_statuses( + &self, + root_item_ids: &[String], + ) -> impl Future> + Send; + + fn item_failures( + &self, + 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( + &self, + event_id: &str, + ) -> impl Future>>> + Send; + + fn version(&self) -> u64; +} + // --------------------------------------------------------------------------- // Context lifecycle // --------------------------------------------------------------------------- @@ -826,6 +880,204 @@ pub struct GetRolloutResponse { pub record: Option, } +// --------------------------------------------------------------------------- +// Datagen lifecycle +// --------------------------------------------------------------------------- + +#[derive(Debug, Serialize, Deserialize)] +pub struct CreateDatagenStoreRequest { + /// Portable dataset name: 1-128 ASCII characters matching + /// `[A-Za-z0-9_][A-Za-z0-9._-]*`; `_registry` and `_stats` are reserved. + pub name: String, + #[serde(default)] + pub storage_options: Option>, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct DatagenStoreInfo { + pub name: String, + pub uri: String, + /// Dataset version. `None` in list responses, which are served from the + /// durable registry without opening each dataset. Single-store lookups + /// (`get`/`create`) always populate it. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub version: Option, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct ListDatagenStoresResponse { + pub stores: Vec, +} + +// --------------------------------------------------------------------------- +// Datagen values +// --------------------------------------------------------------------------- + +/// The wire form of a core `DatagenValue`. `kind` tags the payload: +/// `"int"`/`"float"`/`"bool"`/`"str"`/`"json"` carry a JSON scalar in `value`; +/// `"blob"` carries raw bytes in `bytes` (inline base64 in JSON, or a raw +/// multipart part on the append endpoints) plus `size`/`checksum`. Mirrors the +/// `DatagenValue.to_wire`/`from_wire` dict shape in the Python binding. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DatagenValueDto { + pub kind: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub value: Option, + #[serde( + default, + skip_serializing_if = "Option::is_none", + serialize_with = "serialize_base64_opt", + deserialize_with = "deserialize_base64_opt" + )] + pub bytes: Option>, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub size: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub checksum: Option, +} + +// --------------------------------------------------------------------------- +// Datagen events (append) +// --------------------------------------------------------------------------- + +/// One append-only datagen log row. Field names and semantics mirror the core +/// `DatagenEvent` and the Python `datagen_events` wire dict; enum-valued columns +/// (`event_type`, `step_kind`, `status`) travel as their canonical strings. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DatagenEventDto { + pub event_id: String, + pub item_id: String, + pub root_item_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub parent_item_id: Option, + pub item_seq: i64, + pub checkpoint_id: String, + pub event_type: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub step_name: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub step_kind: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub step_index: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub enclosing_step: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub selector_step: Option, + #[serde(default)] + pub attempt: i32, + pub run_id: String, + pub writer_epoch: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub field_name: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub field_type: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub codec_version: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub value: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub query_tags: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub status: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error_type: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error_dump: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub traceback: Option, + /// Defaults to the server's current time when omitted. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub event_ts: Option>, + pub schema_version: i32, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct AddDatagenEventsRequest { + pub events: Vec, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct AddDatagenEventsResponse { + pub version: u64, + pub count: usize, +} + +// --------------------------------------------------------------------------- +// Datagen folded read lenses +// --------------------------------------------------------------------------- + +/// A folded field: `mode = "set"` carries a single `value`; `mode = "append"` +/// carries an ordered `values` list. Mirrors the Python `FieldState` wire dict. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DatagenFieldStateDto { + pub mode: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub value: Option, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub values: Vec, +} + +/// One completed step position an item passed through (a single STEP_COMPLETED). +/// Mirrors the Python `StepCursor` wire dict. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DatagenStepCursorDto { + pub step_name: String, + pub step_kind: String, + pub step_index: i64, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub enclosing_step: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub selector_step: Option, + pub item_seq: i64, +} + +/// An item reconstructed by folding its events into latest state. +/// Mirrors the Python `FoldedItem` wire dict. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct FoldedDatagenItemDto { + pub item_id: String, + pub root_item_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub parent_item_id: Option, + pub status: String, + pub last_item_seq: i64, + pub last_attempt: i32, + pub fields: std::collections::BTreeMap, + pub trajectory: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub query_tags: Option, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetFoldedDatagenItemResponse { + pub item: Option, +} + +/// Bulk startup classification of root items. A missing id means "never started". +#[derive(Debug, Serialize, Deserialize)] +pub struct DatagenRootItemStatusesResponse { + pub statuses: std::collections::HashMap, +} + +/// One failure record for an item (the failure lens). Mirrors the Python +/// `Failure` wire dict, with the step position flattened under `at`. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DatagenFailureDto { + pub at: DatagenStepCursorDto, + pub run_id: String, + pub attempt: i32, + pub error_type: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error_dump: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub traceback: Option, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct ListDatagenFailuresResponse { + pub failures: Vec, +} + // --------------------------------------------------------------------------- // Error // --------------------------------------------------------------------------- diff --git a/crates/lance-context-client/src/lib.rs b/crates/lance-context-client/src/lib.rs index 71c4e39..2351fc1 100644 --- a/crates/lance-context-client/src/lib.rs +++ b/crates/lance-context-client/src/lib.rs @@ -361,6 +361,108 @@ impl RolloutStoreApi for RemoteRolloutStore { } } +pub struct RemoteDatagenStore { + client: ContextClient, + store_name: String, + cached_version: u64, +} + +impl RemoteDatagenStore { + pub async fn connect(base_url: &str, store_name: &str) -> Result { + let client = ContextClient::new(base_url); + let info = client.get_datagen_store(store_name).await?; + Ok(Self { + client, + store_name: store_name.to_string(), + cached_version: info.version.unwrap_or(0), + }) + } + + pub async fn connect_or_create( + base_url: &str, + req: &CreateDatagenStoreRequest, + ) -> Result { + let client = ContextClient::new(base_url); + let info = match client.get_datagen_store(&req.name).await { + Ok(info) => info, + Err(ClientError::Api { status: 404, .. }) => client.create_datagen_store(req).await?, + Err(e) => return Err(e), + }; + Ok(Self { + client, + store_name: req.name.clone(), + cached_version: info.version.unwrap_or(0), + }) + } +} + +impl DatagenStoreApi for RemoteDatagenStore { + async fn append( + &mut self, + events: &[DatagenEventDto], + ) -> ContextResult { + let resp = self + .client + .add_datagen_events(&self.store_name, events, false) + .await + .map_err(to_ctx_err)?; + self.cached_version = resp.version; + Ok(resp) + } + + async fn append_checkpoint( + &mut self, + events: &[DatagenEventDto], + ) -> ContextResult { + let resp = self + .client + .add_datagen_events(&self.store_name, events, true) + .await + .map_err(to_ctx_err)?; + self.cached_version = resp.version; + Ok(resp) + } + + async fn fold_item(&self, item_id: &str) -> ContextResult> { + let resp = self + .client + .fold_datagen_item(&self.store_name, item_id) + .await + .map_err(to_ctx_err)?; + Ok(resp.item) + } + + async fn root_item_statuses( + &self, + root_item_ids: &[String], + ) -> ContextResult { + self.client + .datagen_root_item_statuses(&self.store_name, root_item_ids) + .await + .map_err(to_ctx_err) + } + + async fn item_failures(&self, item_id: &str) -> ContextResult> { + let resp = self + .client + .datagen_item_failures(&self.store_name, item_id) + .await + .map_err(to_ctx_err)?; + Ok(resp.failures) + } + + async fn get_blob(&self, event_id: &str) -> ContextResult>> { + self.client + .fetch_datagen_blob(&self.store_name, event_id) + .await + .map_err(to_ctx_err) + } + + fn version(&self) -> u64 { + self.cached_version + } +} + fn to_ctx_err(err: ClientError) -> ContextError { match err { ClientError::Api { @@ -870,6 +972,130 @@ impl ContextClient { Self::handle_response(resp).await } + pub async fn create_datagen_store( + &self, + req: &CreateDatagenStoreRequest, + ) -> Result { + let resp = self + .http + .post(self.url("/datagen")) + .json(req) + .send() + .await?; + Self::handle_response(resp).await + } + + pub async fn list_datagen_stores(&self) -> Result { + let resp = self.http.get(self.url("/datagen")).send().await?; + Self::handle_response(resp).await + } + + pub async fn get_datagen_store(&self, name: &str) -> Result { + let resp = self + .http + .get(self.url(&format!("/datagen/{}", name))) + .send() + .await?; + Self::handle_response(resp).await + } + + pub async fn delete_datagen_store(&self, name: &str) -> Result<(), ClientError> { + let resp = self + .http + .delete(self.url(&format!("/datagen/{}", name))) + .send() + .await?; + if resp.status().is_success() { + Ok(()) + } else { + Err(Self::extract_error(resp).await) + } + } + + /// Append datagen events. `checkpoint = true` commits the batch as one atomic + /// step boundary (`append_checkpoint`); `false` appends a raw generation. + /// FIELD blobs are offloaded to a content-addressed artifact store before the + /// event reaches the log, so events are small JSON and travel inline. + pub async fn add_datagen_events( + &self, + name: &str, + events: &[DatagenEventDto], + checkpoint: bool, + ) -> Result { + let req = AddDatagenEventsRequest { + events: events.to_vec(), + }; + let resp = self + .http + .post(self.url(&format!("/datagen/{}/events", name))) + .query(&[("checkpoint", checkpoint)]) + .json(&req) + .send() + .await?; + Self::handle_response(resp).await + } + + pub async fn fold_datagen_item( + &self, + name: &str, + item_id: &str, + ) -> Result { + let resp = self + .http + .get(self.url(&format!("/datagen/{}/items/{}", name, item_id))) + .send() + .await?; + Self::handle_response(resp).await + } + + pub async fn datagen_item_failures( + &self, + name: &str, + item_id: &str, + ) -> Result { + let resp = self + .http + .get(self.url(&format!("/datagen/{}/items/{}/failures", name, item_id))) + .send() + .await?; + Self::handle_response(resp).await + } + + pub async fn datagen_root_item_statuses( + &self, + name: &str, + root_item_ids: &[String], + ) -> Result { + let resp = self + .http + .get(self.url(&format!("/datagen/{}/root-status", name))) + .query(&[("ids", root_item_ids.join(","))]) + .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( + &self, + name: &str, + event_id: &str, + ) -> Result>, ClientError> { + let resp = self + .http + .get(self.url(&format!("/datagen/{}/blobs/{}", name, event_id))) + .send() + .await?; + if resp.status().is_success() { + Ok(Some(resp.bytes().await?.to_vec())) + } else if resp.status().as_u16() == 404 { + Ok(None) + } else { + Err(Self::extract_error(resp).await) + } + } + async fn handle_response( resp: reqwest::Response, ) -> Result { diff --git a/crates/lance-context-core/src/api_impl.rs b/crates/lance-context-core/src/api_impl.rs index 09de0ec..4d945a6 100644 --- a/crates/lance-context-core/src/api_impl.rs +++ b/crates/lance-context-core/src/api_impl.rs @@ -3,14 +3,23 @@ use serde_json::Value; use uuid::Uuid; use lance_context_api::{ - AddRecordRequest, AddRecordsResponse, AddRolloutRequest, AddRolloutsResponse, CompactRequest, - CompactResponse, CompactStatsResponse, ContextError, ContextResult, ContextStoreApi, - DeleteRecordResponse, RecordDto, RecordPatchDto, RelationshipDto, RetrieveRequest, - RetrieveResultDto, RolloutRecordDto, RolloutStoreApi, SearchRequest, SearchResultDto, - StateMetadataDto, UpdateRecordRequest, UpdateRecordResponse, UpsertRecordRequest, - UpsertRecordResponse, UpsertRecordsRequest, UpsertRecordsResponse, UpsertResultDto, + AddDatagenEventsResponse, AddRecordRequest, AddRecordsResponse, AddRolloutRequest, + AddRolloutsResponse, CompactRequest, CompactResponse, CompactStatsResponse, ContextError, + ContextResult, ContextStoreApi, DatagenEventDto, DatagenFailureDto, DatagenFieldStateDto, + DatagenRootItemStatusesResponse, DatagenStepCursorDto, DatagenStoreApi, DatagenValueDto, + DeleteRecordResponse, FoldedDatagenItemDto, RecordDto, RecordPatchDto, RelationshipDto, + RetrieveRequest, RetrieveResultDto, RolloutRecordDto, RolloutStoreApi, SearchRequest, + SearchResultDto, StateMetadataDto, UpdateRecordRequest, UpdateRecordResponse, + UpsertRecordRequest, UpsertRecordResponse, UpsertRecordsRequest, UpsertRecordsResponse, + UpsertResultDto, }; +use crate::datagen::{ + DatagenBlobValue, DatagenEvent, DatagenEventType, DatagenFailure, DatagenFieldState, + DatagenItemLookup, DatagenItemStatus, DatagenRootItemStatuses, DatagenStepCursor, + DatagenStepKind, DatagenValue, FoldedDatagenItem, +}; +use crate::datagen_store::DatagenStore; use crate::record::{ ContextRecord, LifecycleQueryOptions, RecordFilters, RecordPatch, Relationship, StateMetadata, LIFECYCLE_ACTIVE, @@ -678,3 +687,253 @@ fn to_ctx_err(err: lance::Error) -> ContextError { ContextError::Internal(msg) } } + +impl DatagenStoreApi for DatagenStore { + async fn append( + &mut self, + events: &[DatagenEventDto], + ) -> ContextResult { + let core = datagen_events_from_dtos(events)?; + let count = core.len(); + let version = DatagenStore::append(self, &core) + .await + .map_err(to_ctx_err)?; + Ok(AddDatagenEventsResponse { version, count }) + } + + async fn append_checkpoint( + &mut self, + events: &[DatagenEventDto], + ) -> ContextResult { + let core = datagen_events_from_dtos(events)?; + let count = core.len(); + let version = DatagenStore::append_checkpoint(self, &core) + .await + .map_err(to_ctx_err)?; + Ok(AddDatagenEventsResponse { version, count }) + } + + async fn fold_item(&self, item_id: &str) -> ContextResult> { + let lookup = DatagenStore::fold_item(self, item_id) + .await + .map_err(to_ctx_err)?; + Ok(match lookup { + DatagenItemLookup::NeverStarted => None, + DatagenItemLookup::Found(item) => Some(folded_item_to_dto(&item)), + }) + } + + async fn root_item_statuses( + &self, + root_item_ids: &[String], + ) -> ContextResult { + let ids: Vec<&str> = root_item_ids.iter().map(String::as_str).collect(); + let statuses = DatagenStore::root_item_statuses(self, &ids) + .await + .map_err(to_ctx_err)?; + Ok(root_item_statuses_to_dto(&statuses)) + } + + async fn item_failures(&self, item_id: &str) -> ContextResult> { + let failures = DatagenStore::item_failures(self, item_id) + .await + .map_err(to_ctx_err)?; + Ok(failures.iter().map(failure_to_dto).collect()) + } + + async fn get_blob(&self, event_id: &str) -> ContextResult>> { + DatagenStore::get_blob(self, event_id) + .await + .map_err(to_ctx_err) + } + + fn version(&self) -> u64 { + DatagenStore::version(self) + } +} + +fn datagen_events_from_dtos(events: &[DatagenEventDto]) -> ContextResult> { + events.iter().map(datagen_event_from_dto).collect() +} + +fn datagen_event_from_dto(dto: &DatagenEventDto) -> ContextResult { + let event_type = + DatagenEventType::parse(&dto.event_type).map_err(ContextError::InvalidRequest)?; + let step_kind = dto + .step_kind + .as_deref() + .map(DatagenStepKind::parse) + .transpose() + .map_err(ContextError::InvalidRequest)?; + let status = dto + .status + .as_deref() + .map(DatagenItemStatus::parse) + .transpose() + .map_err(ContextError::InvalidRequest)?; + let value = dto.value.as_ref().map(datagen_value_from_dto).transpose()?; + Ok(DatagenEvent { + event_id: dto.event_id.clone(), + item_id: dto.item_id.clone(), + root_item_id: dto.root_item_id.clone(), + parent_item_id: dto.parent_item_id.clone(), + item_seq: dto.item_seq, + checkpoint_id: dto.checkpoint_id.clone(), + event_type, + step_name: dto.step_name.clone(), + step_kind, + step_index: dto.step_index, + enclosing_step: dto.enclosing_step.clone(), + selector_step: dto.selector_step.clone(), + attempt: dto.attempt, + run_id: dto.run_id.clone(), + writer_epoch: dto.writer_epoch.clone(), + field_name: dto.field_name.clone(), + field_type: dto.field_type.clone(), + codec_version: dto.codec_version, + value, + query_tags: dto.query_tags.clone(), + status, + error_type: dto.error_type.clone(), + error_dump: dto.error_dump.clone(), + traceback: dto.traceback.clone(), + event_ts: dto.event_ts.unwrap_or_else(Utc::now), + schema_version: dto.schema_version, + }) +} + +fn datagen_value_from_dto(dto: &DatagenValueDto) -> ContextResult { + let missing = || { + ContextError::InvalidRequest(format!( + "datagen value of kind '{}' is missing 'value'", + dto.kind + )) + }; + match dto.kind.as_str() { + "int" => Ok(DatagenValue::Int( + dto.value + .as_ref() + .and_then(Value::as_i64) + .ok_or_else(missing)?, + )), + "float" => Ok(DatagenValue::Float( + dto.value + .as_ref() + .and_then(Value::as_f64) + .ok_or_else(missing)?, + )), + "bool" => Ok(DatagenValue::Bool( + dto.value + .as_ref() + .and_then(Value::as_bool) + .ok_or_else(missing)?, + )), + "str" => Ok(DatagenValue::Str( + dto.value + .as_ref() + .and_then(Value::as_str) + .ok_or_else(missing)? + .to_string(), + )), + "json" => Ok(DatagenValue::Json(dto.value.clone().ok_or_else(missing)?)), + "blob" => Ok(DatagenValue::Blob(DatagenBlobValue { + bytes: dto.bytes.clone(), + size: dto + .size + .unwrap_or_else(|| dto.bytes.as_ref().map(|b| b.len() as i64).unwrap_or(0)), + checksum: dto.checksum.clone(), + })), + other => Err(ContextError::InvalidRequest(format!( + "unsupported datagen value kind '{other}'" + ))), + } +} + +fn datagen_value_to_dto(value: &DatagenValue) -> DatagenValueDto { + let mut dto = DatagenValueDto { + kind: value.kind().to_string(), + value: None, + bytes: None, + size: None, + checksum: None, + }; + match value { + DatagenValue::Int(inner) => dto.value = Some(Value::from(*inner)), + DatagenValue::Float(inner) => dto.value = Some(Value::from(*inner)), + DatagenValue::Bool(inner) => dto.value = Some(Value::from(*inner)), + DatagenValue::Str(inner) => dto.value = Some(Value::from(inner.clone())), + DatagenValue::Json(inner) => dto.value = Some(inner.clone()), + DatagenValue::Blob(blob) => { + dto.bytes = blob.bytes.clone(); + dto.size = Some(blob.size); + dto.checksum = blob.checksum.clone(); + } + } + dto +} + +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(), + parent_item_id: item.parent_item_id.as_ref().map(ToString::to_string), + status: item.status.as_str().to_string(), + last_item_seq: item.last_item_seq, + last_attempt: item.last_attempt, + fields: item + .fields + .iter() + .map(|(name, state)| (name.clone(), field_state_to_dto(state))) + .collect(), + trajectory: item.trajectory.ordered.iter().map(cursor_to_dto).collect(), + query_tags: item.query_tags.clone(), + } +} + +fn field_state_to_dto(state: &DatagenFieldState) -> DatagenFieldStateDto { + match state { + DatagenFieldState::Set(value) => DatagenFieldStateDto { + mode: "set".to_string(), + value: Some(datagen_value_to_dto(value)), + values: Vec::new(), + }, + DatagenFieldState::Appended(values) => DatagenFieldStateDto { + mode: "append".to_string(), + value: None, + values: values.iter().map(datagen_value_to_dto).collect(), + }, + } +} + +fn cursor_to_dto(cursor: &DatagenStepCursor) -> DatagenStepCursorDto { + DatagenStepCursorDto { + step_name: cursor.position.step.name.clone(), + step_kind: cursor.position.step.kind.as_str().to_string(), + step_index: cursor.position.index, + enclosing_step: cursor.position.enclosing.clone(), + selector_step: cursor.position.selector.clone(), + item_seq: cursor.item_seq, + } +} + +fn root_item_statuses_to_dto( + statuses: &DatagenRootItemStatuses, +) -> DatagenRootItemStatusesResponse { + DatagenRootItemStatusesResponse { + statuses: statuses + .iter() + .map(|(id, status)| (id.clone(), status.as_str().to_string())) + .collect(), + } +} + +fn failure_to_dto(failure: &DatagenFailure) -> DatagenFailureDto { + DatagenFailureDto { + at: cursor_to_dto(&failure.at), + run_id: failure.run_id.clone(), + attempt: failure.attempt, + error_type: failure.error.error_type.clone(), + error_dump: failure.error.error_dump.clone(), + traceback: failure.error.traceback.clone(), + } +} diff --git a/crates/lance-context-server/src/error.rs b/crates/lance-context-server/src/error.rs index c5bb116..6e48cbb 100644 --- a/crates/lance-context-server/src/error.rs +++ b/crates/lance-context-server/src/error.rs @@ -1,7 +1,7 @@ use axum::http::StatusCode; use axum::response::{IntoResponse, Response}; use axum::Json; -use lance_context_api::{ErrorBody, ErrorResponse}; +use lance_context_api::{ContextError, ErrorBody, ErrorResponse}; use lance_context_core::LanceError; #[derive(Debug)] @@ -44,6 +44,19 @@ impl AppError { other => AppError::Internal(other.to_string()), } } + + /// Map an API-layer [`ContextError`] onto the server's error taxonomy. + /// Datagen routes call the `DatagenStoreApi` trait (which owns the DTO↔core + /// conversion) and so surface `ContextError` rather than raw `LanceError`. + pub fn from_context(err: ContextError) -> Self { + match err { + ContextError::NotFound(msg) => AppError::NotFound(msg), + ContextError::AlreadyExists(msg) => AppError::AlreadyExists(msg), + ContextError::InvalidRequest(msg) => AppError::InvalidRequest(msg), + ContextError::Internal(msg) => AppError::Internal(msg), + ContextError::CompactionInProgress => AppError::CompactionInProgress, + } + } } impl IntoResponse for AppError { diff --git a/crates/lance-context-server/src/routes/datagen.rs b/crates/lance-context-server/src/routes/datagen.rs new file mode 100644 index 0000000..bbdb35a --- /dev/null +++ b/crates/lance-context-server/src/routes/datagen.rs @@ -0,0 +1,225 @@ +use std::sync::Arc; + +use axum::extract::{Path, Query, State}; +use axum::http::StatusCode; +use axum::Json; +use lance_context_api::{ + AddDatagenEventsRequest, AddDatagenEventsResponse, CreateDatagenStoreRequest, DatagenStoreApi, + DatagenStoreInfo, GetFoldedDatagenItemResponse, ListDatagenFailuresResponse, + ListDatagenStoresResponse, +}; +use lance_context_core::{DatagenStore, DatagenStoreOptions}; +use tokio::sync::RwLock; + +use crate::error::AppError; +use crate::state::AppState; + +/// Upper bound on a single datagen append request body. FIELD blobs are +/// offloaded to a content-addressed artifact store before the event reaches the +/// log, so events themselves are small; the ceiling still bounds a pathological +/// batch. +pub const MAX_DATAGEN_UPLOAD_BYTES: usize = 256 * 1024 * 1024; + +pub async fn create_datagen_store( + State(state): State>, + Json(req): Json, +) -> Result<(StatusCode, Json), AppError> { + AppState::validate_name(&req.name)?; + if state + .datagen_registry + .write() + .await + .contains(&req.name) + .await + .map_err(AppError::from_lance)? + { + return Err(AppError::AlreadyExists(format!( + "Datagen store '{}' already exists", + req.name + ))); + } + + let uri = state.datagen_uri(&req.name); + let options = DatagenStoreOptions { + storage_options: req.storage_options, + shard_id: state.instance_id.clone(), + merge_after_generations: None, + cleanup_interval_secs: None, + }; + let store = DatagenStore::open_with_options(&uri, options) + .await + .map_err(AppError::from_lance)?; + let version = store.version(); + + let store = Arc::new(RwLock::new(store)); + state.register_datagen(&req.name, &uri, store).await?; + + Ok(( + StatusCode::CREATED, + Json(DatagenStoreInfo { + name: req.name, + uri, + version: Some(version), + }), + )) +} + +pub async fn list_datagen_stores( + State(state): State>, +) -> Result, AppError> { + let entries = state + .datagen_registry + .write() + .await + .list() + .await + .map_err(AppError::from_lance)?; + let stores = entries + .into_iter() + .map(|entry| DatagenStoreInfo { + name: entry.name, + uri: entry.uri, + version: None, + }) + .collect(); + Ok(Json(ListDatagenStoresResponse { stores })) +} + +pub async fn get_datagen_store( + State(state): State>, + Path(name): Path, +) -> Result, AppError> { + let store_lock = state.get_or_open_datagen_store(&name).await?; + let store = store_lock.read().await; + Ok(Json(DatagenStoreInfo { + name: name.clone(), + uri: state.datagen_uri(&name), + version: Some(store.version()), + })) +} + +pub async fn delete_datagen_store( + State(state): State>, + Path(name): Path, +) -> Result { + if !state.unregister_datagen(&name).await? { + return Err(AppError::NotFound(format!( + "Datagen store '{}' does not exist", + name + ))); + } + let uri = state.datagen_uri(&name); + if let Err(e) = tokio::fs::remove_dir_all(&uri).await { + tracing::warn!("Failed to remove datagen data at {}: {}", uri, e); + } + Ok(StatusCode::NO_CONTENT) +} + +/// Which append semantics the `/events` endpoint applies. +#[derive(Debug, Default, serde::Deserialize)] +pub struct AppendParams { + /// When `true`, the batch is committed as one atomic checkpoint (FIELD_* + /// events plus exactly one STEP_COMPLETED) via `append_checkpoint`. Default + /// appends the events as one raw MemWAL generation. + #[serde(default)] + pub checkpoint: bool, +} + +pub async fn add_datagen_events( + State(state): State>, + Path(name): Path, + Query(params): Query, + Json(req): Json, +) -> Result<(StatusCode, Json), AppError> { + if req.events.is_empty() { + return Err(AppError::InvalidRequest( + "events array must not be empty".to_string(), + )); + } + + let store_lock = state.get_or_open_datagen_store(&name).await?; + let mut store = store_lock.write().await; + let resp = if params.checkpoint { + DatagenStoreApi::append_checkpoint(&mut *store, &req.events).await + } else { + DatagenStoreApi::append(&mut *store, &req.events).await + } + .map_err(AppError::from_context)?; + + Ok((StatusCode::CREATED, Json(resp))) +} + +pub async fn fold_datagen_item( + State(state): State>, + Path((name, 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 item = DatagenStoreApi::fold_item(&*store, &item_id) + .await + .map_err(AppError::from_context)?; + Ok(Json(GetFoldedDatagenItemResponse { item })) +} + +pub async fn datagen_item_failures( + State(state): State>, + Path((name, 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 failures = DatagenStoreApi::item_failures(&*store, &item_id) + .await + .map_err(AppError::from_context)?; + Ok(Json(ListDatagenFailuresResponse { failures })) +} + +#[derive(Debug, Default, serde::Deserialize)] +pub struct RootStatusParams { + /// Comma-separated list of root item ids to classify. + pub ids: Option, +} + +pub async fn datagen_root_item_statuses( + State(state): State>, + Path(name): Path, + Query(params): Query, +) -> Result, AppError> { + let ids: Vec = params + .ids + .as_deref() + .map(|raw| { + raw.split(',') + .map(str::trim) + .filter(|s| !s.is_empty()) + .map(str::to_string) + .collect() + }) + .unwrap_or_default(); + + let store_lock = state.get_or_open_datagen_store(&name).await?; + let store = store_lock.read().await; + let resp = DatagenStoreApi::root_item_statuses(&*store, &ids) + .await + .map_err(AppError::from_context)?; + Ok(Json(resp)) +} + +/// Materialize one FIELD_* event's offloaded blob bytes by event id. The bytes +/// are opaque, so they return as `application/octet-stream`; `404` when the +/// event or its payload is absent. +pub async fn fetch_datagen_blob( + State(state): State>, + Path((name, event_id)): Path<(String, String)>, +) -> Result { + use axum::http::header; + use axum::response::IntoResponse; + + let store_lock = state.get_or_open_datagen_store(&name).await?; + let store = store_lock.read().await; + let bytes = DatagenStoreApi::get_blob(&*store, &event_id) + .await + .map_err(AppError::from_context)? + .ok_or_else(|| AppError::NotFound(format!("Datagen event '{}' has no blob", event_id)))?; + + Ok(([(header::CONTENT_TYPE, "application/octet-stream")], bytes).into_response()) +} diff --git a/crates/lance-context-server/src/routes/mod.rs b/crates/lance-context-server/src/routes/mod.rs index c42eb15..17ab77d 100644 --- a/crates/lance-context-server/src/routes/mod.rs +++ b/crates/lance-context-server/src/routes/mod.rs @@ -1,5 +1,6 @@ pub mod compact; pub mod contexts; +pub mod datagen; pub mod health; pub mod records; pub mod rollouts; @@ -117,6 +118,34 @@ pub fn router() -> Router> { "/api/v1/internal/merge-wal/{name}", post(rollouts::merge_wal), ) + .route("/api/v1/datagen", post(datagen::create_datagen_store)) + .route("/api/v1/datagen", get(datagen::list_datagen_stores)) + .route("/api/v1/datagen/{name}", get(datagen::get_datagen_store)) + .route( + "/api/v1/datagen/{name}", + delete(datagen::delete_datagen_store), + ) + .route( + "/api/v1/datagen/{name}/events", + post(datagen::add_datagen_events) + .layer(DefaultBodyLimit::max(datagen::MAX_DATAGEN_UPLOAD_BYTES)), + ) + .route( + "/api/v1/datagen/{name}/items/{item_id}", + get(datagen::fold_datagen_item), + ) + .route( + "/api/v1/datagen/{name}/items/{item_id}/failures", + get(datagen::datagen_item_failures), + ) + .route( + "/api/v1/datagen/{name}/root-status", + get(datagen::datagen_root_item_statuses), + ) + .route( + "/api/v1/datagen/{name}/blobs/{event_id}", + get(datagen::fetch_datagen_blob), + ) } #[cfg(test)] diff --git a/crates/lance-context-server/src/state.rs b/crates/lance-context-server/src/state.rs index 1bea836..d0dd394 100644 --- a/crates/lance-context-server/src/state.rs +++ b/crates/lance-context-server/src/state.rs @@ -5,8 +5,8 @@ use std::sync::Arc; use std::time::Duration; use lance_context_core::{ - join_uri, validate_store_name, ContextStore, ContextStoreOptions, RolloutRegistry, - RolloutStore, RolloutStoreOptions, Session, + join_uri, validate_store_name, ContextStore, ContextStoreOptions, DatagenStore, + DatagenStoreOptions, RolloutRegistry, RolloutStore, RolloutStoreOptions, Session, }; use lru::LruCache; use tokio::sync::{Mutex, RwLock}; @@ -74,6 +74,15 @@ pub struct AppState { /// than Lance's default 6 GiB *per store*. `None` restores the per-store /// default session (leak-prone; only when the budget is configured to `0`). pub rollout_session: Option>, + /// Bounded LRU of resident datagen-store handles, mirroring + /// [`Self::rollout_stores`]. Datagen delta-log datasets are also one per + /// experiment, so the same residency bound and durable-registry existence + /// model applies. + pub datagen_stores: Mutex>>>, + /// Durable directory of which datagen stores exist. A separate registry + /// dataset from [`Self::rollout_registry`] so the two store kinds never + /// collide on a shared name. + pub datagen_registry: RwLock, } /// Process-wide admission control for the total artifact-blob payload held in @@ -185,6 +194,10 @@ impl AppState { let blob_budget = (config.rollout_max_inflight_blob_bytes > 0) .then(|| BlobBudget::new(config.rollout_max_inflight_blob_bytes)); let rollout_session = build_rollout_session(config.rollout_cache_bytes); + let datagen_registry_uri = join_uri(&base_uri, "_registry.datagen.lance"); + let datagen_registry = RolloutRegistry::open_or_create(&datagen_registry_uri, None) + .await + .map_err(AppError::from_lance)?; Ok(Self { stores: RwLock::new(std::collections::HashMap::new()), rollout_stores: Mutex::new(LruCache::new(capacity)), @@ -196,6 +209,8 @@ impl AppState { rollout_flush_interval_secs: config.rollout_flush_interval_secs, blob_budget, rollout_session, + datagen_stores: Mutex::new(LruCache::new(capacity)), + datagen_registry: RwLock::new(datagen_registry), }) } @@ -222,6 +237,10 @@ impl AppState { let registry = RolloutRegistry::open_or_create(®istry_uri, None) .await .expect("open test registry"); + let datagen_registry_uri = join_uri(&base_uri, "_registry.datagen.lance"); + let datagen_registry = RolloutRegistry::open_or_create(&datagen_registry_uri, None) + .await + .expect("open test datagen registry"); Self { stores: RwLock::new(std::collections::HashMap::new()), rollout_stores: Mutex::new(LruCache::new( @@ -235,6 +254,10 @@ impl AppState { rollout_flush_interval_secs: 0, blob_budget: None, rollout_session: build_rollout_session(2 * 1024 * 1024 * 1024), + datagen_stores: Mutex::new(LruCache::new( + NonZeroUsize::new(DEFAULT_ROLLOUT_CACHE_CAPACITY).unwrap(), + )), + datagen_registry: RwLock::new(datagen_registry), } } @@ -392,6 +415,109 @@ impl AppState { validate_store_name(name).map_err(AppError::InvalidRequest) } + /// Datagen delta-log datasets live under a distinct `.datagen.lance` suffix + /// so a datagen store shares neither a rollout nor a context store's path. + pub fn datagen_uri(&self, name: &str) -> String { + join_uri(&self.base_uri, &format!("{}.datagen.lance", name)) + } + + fn datagen_store_options(&self) -> DatagenStoreOptions { + DatagenStoreOptions { + storage_options: None, + shard_id: self.instance_id.clone(), + merge_after_generations: (self.rollout_merge_after_generations > 0) + .then_some(self.rollout_merge_after_generations), + cleanup_interval_secs: None, + } + } + + /// Record that a datagen store exists, in both the durable registry and the + /// in-memory LRU. Called by the create route after the dataset is written. + pub async fn register_datagen( + &self, + name: &str, + uri: &str, + store: Arc>, + ) -> Result<(), AppError> { + Self::validate_name(name)?; + self.datagen_registry + .write() + .await + .upsert(name, uri) + .await + .map_err(AppError::from_lance)?; + self.datagen_stores + .lock() + .await + .put(name.to_string(), store); + Ok(()) + } + + /// Remove a datagen store from the durable registry and evict any resident + /// handle. Returns whether the store existed. + pub async fn unregister_datagen(&self, name: &str) -> Result { + Self::validate_name(name)?; + let existed = self + .datagen_registry + .write() + .await + .contains(name) + .await + .map_err(AppError::from_lance)?; + if !existed { + return Ok(false); + } + self.datagen_registry + .write() + .await + .remove(name) + .await + .map_err(AppError::from_lance)?; + self.datagen_stores.lock().await.pop(name); + Ok(true) + } + + /// Look up a datagen store by name, lazily loading it from object storage on + /// a local cache miss. Existence is resolved against the durable registry, so + /// an evicted-but-registered store reopens rather than 404ing. See + /// [`Self::get_or_open_rollout_store`] for the full rationale. + pub async fn get_or_open_datagen_store( + &self, + name: &str, + ) -> Result>, AppError> { + Self::validate_name(name)?; + if let Some(store) = self.datagen_stores.lock().await.get(name) { + return Ok(store.clone()); + } + + let exists = self + .datagen_registry + .write() + .await + .contains(name) + .await + .map_err(AppError::from_lance)?; + if !exists { + return Err(AppError::NotFound(format!( + "Datagen store '{}' does not exist", + name + ))); + } + + let uri = self.datagen_uri(name); + let opened = DatagenStore::open_existing_with_options(&uri, self.datagen_store_options()) + .await + .map_err(AppError::from_lance)?; + let opened = Arc::new(RwLock::new(opened)); + + let mut cache = self.datagen_stores.lock().await; + if let Some(existing) = cache.get(name) { + return Ok(existing.clone()); + } + cache.put(name.to_string(), opened.clone()); + Ok(opened) + } + /// Spawn the single, process-wide WAL-cleanup sweeper. /// /// This replaces the former one-timer-per-store model, which does not scale diff --git a/crates/lance-context/src/lib.rs b/crates/lance-context/src/lib.rs index 93d6649..063ccfe 100644 --- a/crates/lance-context/src/lib.rs +++ b/crates/lance-context/src/lib.rs @@ -6,27 +6,34 @@ pub use lance_context_core::{ datagen_event_id, datagen_log_schema, datagen_trajectory, fold_datagen_events, CompactionConfig, CompactionMetrics, CompactionStats, Context, ContextEntry, ContextNamespace, ContextRecord, ContextStoreOptions, DatagenBlobValue, DatagenEvent, DatagenEventType, - DatagenFailure, DatagenFieldState, DatagenItemStatus, DatagenStepCursor, DatagenStore, - DatagenStoreOptions, DatagenTerminal, DatagenTrajectory, DatagenValue, FoldedDatagenItem, - IdIndexType, LifecycleQueryOptions, MetadataFilter, PartitionInfo, PartitionSelector, - PartitionSpec, RecordFilters, Relationship, RetrieveResult, RolloutFilters, RolloutRecord, - SearchResult, Snapshot, StateMetadata, DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, - LIFECYCLE_CONTRADICTED, + 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, }; pub use lance_context_api::{ - AddRecordRequest, AddRecordsResponse, AddRolloutRequest, AddRolloutsResponse, CompactRequest, - CompactResponse, CompactStatsResponse, ContextError, ContextResult, ContextStoreApi, - CreateRolloutStoreRequest, DeleteRecordResponse, RecordDto, RelationshipDto, RetrieveRequest, + AddDatagenEventsRequest, AddDatagenEventsResponse, AddRecordRequest, AddRecordsResponse, + AddRolloutRequest, AddRolloutsResponse, CompactRequest, CompactResponse, CompactStatsResponse, + ContextError, ContextResult, ContextStoreApi, CreateDatagenStoreRequest, + CreateRolloutStoreRequest, DatagenEventDto, DatagenFailureDto, DatagenFieldStateDto, + DatagenRootItemStatusesResponse, DatagenStepCursorDto, DatagenStoreApi, DatagenValueDto, + DeleteRecordResponse, FoldedDatagenItemDto, RecordDto, RelationshipDto, RetrieveRequest, RetrieveResponse, RetrieveResultDto, RolloutRecordDto, RolloutStoreApi, SearchResultDto, UpsertRecordRequest, UpsertRecordResponse, }; #[cfg(feature = "remote")] -pub use lance_context_client::{ClientError, RemoteContextStore, RemoteRolloutStore}; +pub use lance_context_client::{ + ClientError, RemoteContextStore, RemoteDatagenStore, RemoteRolloutStore, +}; mod unified; pub use unified::ContextStore; mod unified_rollout; pub use unified_rollout::RolloutStore; + +mod unified_datagen; +pub use unified_datagen::DatagenStore; diff --git a/crates/lance-context/src/unified_datagen.rs b/crates/lance-context/src/unified_datagen.rs new file mode 100644 index 0000000..b43861b --- /dev/null +++ b/crates/lance-context/src/unified_datagen.rs @@ -0,0 +1,130 @@ +use lance_context_api::{ + AddDatagenEventsResponse, ContextError, ContextResult, DatagenEventDto, DatagenFailureDto, + DatagenRootItemStatusesResponse, DatagenStoreApi, FoldedDatagenItemDto, +}; +use lance_context_core::{DatagenStore as LocalStore, DatagenStoreOptions}; + +#[cfg(feature = "remote")] +use lance_context_client::RemoteDatagenStore; + +/// A datagen checkpoint store that is either an in-process Lance dataset +/// (`Local`) or a handle to a remote server (`Remote`). Mirrors +/// [`crate::RolloutStore`] but for the datagen delta-log schema. +pub enum DatagenStore { + Local(Box), + #[cfg(feature = "remote")] + Remote(RemoteDatagenStore), +} + +impl DatagenStore { + pub async fn open(uri: &str) -> Result { + Self::open_with_options(uri, None).await + } + + pub async fn open_with_options( + uri: &str, + storage_options: Option>, + ) -> Result { + let options = DatagenStoreOptions { + storage_options, + // Embedded single-process use writes to the fallback shard; a + // multi-writer embedded deployment threads a per-writer id through + // the core `DatagenStore` directly. + shard_id: None, + merge_after_generations: None, + cleanup_interval_secs: None, + }; + let store = LocalStore::open_with_options(uri, options) + .await + .map_err(|e| ContextError::Internal(e.to_string()))?; + Ok(Self::Local(Box::new(store))) + } + + #[cfg(feature = "remote")] + pub async fn connect(base_url: &str, store_name: &str) -> Result { + let store = RemoteDatagenStore::connect(base_url, store_name) + .await + .map_err(|e| ContextError::Internal(e.to_string()))?; + Ok(Self::Remote(store)) + } + + #[cfg(feature = "remote")] + pub async fn connect_or_create( + base_url: &str, + req: &lance_context_api::CreateDatagenStoreRequest, + ) -> Result { + let store = RemoteDatagenStore::connect_or_create(base_url, req) + .await + .map_err(|e| ContextError::Internal(e.to_string()))?; + Ok(Self::Remote(store)) + } +} + +macro_rules! dispatch_mut { + ($self:expr, $method:ident $(, $arg:expr)*) => { + match $self { + DatagenStore::Local(s) => DatagenStoreApi::$method(s.as_mut() $(, $arg)*).await, + #[cfg(feature = "remote")] + DatagenStore::Remote(s) => DatagenStoreApi::$method(s $(, $arg)*).await, + } + }; +} + +macro_rules! dispatch_ref { + ($self:expr, $method:ident $(, $arg:expr)*) => { + match $self { + DatagenStore::Local(s) => DatagenStoreApi::$method(s.as_ref() $(, $arg)*).await, + #[cfg(feature = "remote")] + DatagenStore::Remote(s) => DatagenStoreApi::$method(s $(, $arg)*).await, + } + }; +} + +macro_rules! dispatch_sync { + ($self:expr, $method:ident $(, $arg:expr)*) => { + match $self { + DatagenStore::Local(s) => DatagenStoreApi::$method(s.as_ref() $(, $arg)*), + #[cfg(feature = "remote")] + DatagenStore::Remote(s) => DatagenStoreApi::$method(s $(, $arg)*), + } + }; +} + +impl DatagenStoreApi for DatagenStore { + async fn append( + &mut self, + events: &[DatagenEventDto], + ) -> ContextResult { + dispatch_mut!(self, append, events) + } + + async fn append_checkpoint( + &mut self, + events: &[DatagenEventDto], + ) -> ContextResult { + dispatch_mut!(self, append_checkpoint, events) + } + + async fn fold_item(&self, item_id: &str) -> ContextResult> { + dispatch_ref!(self, fold_item, item_id) + } + + async fn root_item_statuses( + &self, + root_item_ids: &[String], + ) -> ContextResult { + dispatch_ref!(self, root_item_statuses, root_item_ids) + } + + async fn item_failures(&self, item_id: &str) -> ContextResult> { + dispatch_ref!(self, item_failures, item_id) + } + + async fn get_blob(&self, event_id: &str) -> ContextResult>> { + dispatch_ref!(self, get_blob, event_id) + } + + fn version(&self) -> u64 { + dispatch_sync!(self, version) + } +} diff --git a/python/python/lance_context/__init__.py b/python/python/lance_context/__init__.py index bbb5ec5..7072b2b 100644 --- a/python/python/lance_context/__init__.py +++ b/python/python/lance_context/__init__.py @@ -5,10 +5,12 @@ AsyncRolloutStore, Context, ContextNamespace, + DatagenStore, EmbeddingProvider, RemoteContext, RolloutStore, __version__, + datagen_event_id, generate_id, ) from .embeddings import ( # pyright: ignore[reportMissingImports] @@ -20,10 +22,12 @@ "AsyncRolloutStore", "Context", "ContextNamespace", + "DatagenStore", "EmbeddingProvider", "MultiModalEmbeddingProvider", "RemoteContext", "RolloutStore", "__version__", + "datagen_event_id", "generate_id", ] diff --git a/python/python/lance_context/api.py b/python/python/lance_context/api.py index 11a6998..90f3b75 100644 --- a/python/python/lance_context/api.py +++ b/python/python/lance_context/api.py @@ -3,7 +3,7 @@ import asyncio import json import warnings -from collections.abc import Iterable, Mapping +from collections.abc import Iterable, Mapping, Sequence from datetime import datetime from io import BytesIO from typing import TYPE_CHECKING, Any @@ -15,12 +15,18 @@ from ._internal import ( # pyright: ignore[reportMissingImports] ContextNamespace as _ContextNamespace, ) +from ._internal import ( # pyright: ignore[reportMissingImports] + DatagenStore as _DatagenStore, +) from ._internal import ( # pyright: ignore[reportMissingImports] RemoteContext as _RemoteContext, ) from ._internal import ( # pyright: ignore[reportMissingImports] RolloutStore as _RolloutStore, ) +from ._internal import ( # pyright: ignore[reportMissingImports] + datagen_event_id as _datagen_event_id, +) from ._internal import ( # pyright: ignore[reportMissingImports] generate_id as _generate_id, ) @@ -32,10 +38,12 @@ "AsyncRolloutStore", "Context", "ContextNamespace", + "DatagenStore", "EmbeddingProvider", "RemoteContext", "RolloutStore", "__version__", + "datagen_event_id", "generate_id", ] @@ -51,6 +59,15 @@ def generate_id() -> str: return _generate_id() +def datagen_event_id(item_id: str, checkpoint_id: str, ordinal: int) -> str: + """Compute the deterministic event id for a datagen checkpoint field write. + + The id is a pure function of `item_id`, `checkpoint_id`, and `ordinal`, so a + retried append produces the same id and dedupes against the log. + """ + return _datagen_event_id(item_id, checkpoint_id, ordinal) + + _ARROW_STREAM_MIME = "application/vnd.apache.arrow.stream" _DEFAULT_INGEST_BATCH_SIZE = 1000 _MISSING = object() @@ -2738,3 +2755,89 @@ async def checkout(self, version: int) -> None: def __repr__(self) -> str: return f"AsyncRolloutStore(version={self._sync.version()})" + + +class DatagenStore: + """Synchronous append-only datagen delta-log store. + + Wraps the native store: events are appended as plain dicts and item state is + recovered by folding an item's events. Open an embedded log with :meth:`open`. + + Event dicts carry the datagen event schema (``event_id``, ``item_id``, + ``root_item_id``, ``event_type``, ``item_seq``, ...); values use the tagged + ``{"kind": ..., "value": ...}`` shape. Folded items and failures come back as + plain dicts. + """ + + def __init__(self, sync_store: _DatagenStore) -> None: + self._sync = sync_store + + @classmethod + def open( + cls, + uri: str, + *, + storage_options: Mapping[str, str] | None = None, + shard_id: str | None = None, + ) -> "DatagenStore": + """Open (or create) an embedded datagen log at ``uri``. + + ``shard_id`` gives this writer a stable identity for multi-writer fencing. + """ + opts = dict(storage_options) if storage_options else None + return cls(_DatagenStore.open(uri, opts, shard_id)) + + @classmethod + def connect(cls, base_url: str, name: str) -> "DatagenStore": + """Connect to an existing datagen store on a remote server.""" + return cls(_DatagenStore.connect(base_url, name)) + + @classmethod + def connect_or_create( + cls, + base_url: str, + name: str, + *, + storage_options: Mapping[str, str] | None = None, + ) -> "DatagenStore": + """Connect to a remote datagen store, creating it if absent.""" + opts = dict(storage_options) if storage_options else None + return cls(_DatagenStore.connect_or_create(base_url, name, opts)) + + def version(self) -> int: + """Return the current store version (base dataset version).""" + return self._sync.version() + + def append_checkpoint(self, events: Iterable[Mapping[str, Any]]) -> int: + """Append one completed step boundary atomically. + + ``events`` share one item/checkpoint/writer attempt and contain exactly + one ``STEP_COMPLETED``. Returns the new store version. + """ + return self._sync.append_checkpoint(list(events)) + + def append(self, events: Iterable[Mapping[str, Any]]) -> int: + """Append raw events as one MemWAL generation. Returns the new store version.""" + return self._sync.append(list(events)) + + def fold_item(self, item_id: str) -> dict[str, Any] | None: + """Fold an item's events into its latest state, or ``None`` if never started.""" + return self._sync.fold_item(item_id) + + def root_item_statuses(self, root_item_ids: Sequence[str]) -> dict[str, str]: + """Classify each root item id by folded lifecycle status. + + Missing ids (never started) are absent from the returned dict. + """ + return self._sync.root_item_statuses(list(root_item_ids)) + + def item_failures(self, item_id: str) -> list[dict[str, Any]]: + """All failure records for an item (the failure lens), oldest first.""" + return self._sync.item_failures(item_id) + + 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 __repr__(self) -> str: + return f"DatagenStore(version={self._sync.version()})" diff --git a/python/src/lib.rs b/python/src/lib.rs index 7491ba4..9cc9283 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -12,8 +12,10 @@ use serde_json::Value; use tokio::runtime::Runtime; use lance_context::{ - AddRolloutRequest, CreateRolloutStoreRequest, RolloutRecordDto, - RolloutStore as UnifiedRolloutStore, RolloutStoreApi, + AddRolloutRequest, CreateDatagenStoreRequest, CreateRolloutStoreRequest, DatagenEventDto, + DatagenFailureDto, DatagenFieldStateDto, DatagenStepCursorDto, + DatagenStore as UnifiedDatagenStore, DatagenStoreApi, DatagenValueDto, FoldedDatagenItemDto, + RolloutRecordDto, RolloutStore as UnifiedRolloutStore, RolloutStoreApi, }; use lance_context_api::{ AddRecordRequest, CompactRequest, CompactResponse, CompactStatsResponse, ContextStoreApi, @@ -23,12 +25,13 @@ use lance_context_api::{ use lance_context_client::RemoteContextStore; use lance_context_core::serde::CONTENT_TYPE_TEXT; use lance_context_core::{ - CompactionConfig, CompactionMetrics, CompactionStats, Context as RustContext, - ContextNamespace as RustContextNamespace, ContextRecord, ContextStore, ContextStoreOptions, - DistanceMetric, EvalConfig, EvalQuerySet, ExportConfig, ExportTask, GroupBy, IdIndexType, - LifecycleQueryOptions, PartitionInfo, PartitionSelector, PartitionSpec, PreferenceForm, - ReadProjection, RecordFilters, RecordPatch, Relationship, RetrievalMode, RetrieveResult, - SearchResult, SplitConfig, StateMetadata, LIFECYCLE_ACTIVE, + 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, + PartitionSelector, PartitionSpec, PreferenceForm, ReadProjection, RecordFilters, RecordPatch, + Relationship, RetrievalMode, RetrieveResult, SearchResult, SplitConfig, StateMetadata, + DATAGEN_SCHEMA_VERSION, LIFECYCLE_ACTIVE, }; const DEFAULT_BINARY_CONTENT_TYPE: &str = "application/octet-stream"; @@ -2548,13 +2551,409 @@ fn rollout_record_to_json(record: &RolloutRecordDto) -> PyResult { serde_json::to_string(record).map_err(to_py_err) } +// --------------------------------------------------------------------------- +// Datagen store binding (append-only delta-log, fold model) +// --------------------------------------------------------------------------- + +/// A single embedded datagen checkpoint log (`open`). +/// +/// Events cross the FFI boundary as plain dicts (not pickle): the executor +/// builds each `DatagenEvent` field-by-field, `append_checkpoint` persists one +/// atomic step boundary, and reads (`fold_item`) return the folded item state as +/// a dict. Mirrors the concurrency contract of the rollout store: open a fresh +/// handle per writer, never share a handle across concurrent appends. +#[pyclass] +struct DatagenStore { + store: UnifiedDatagenStore, + runtime: Arc, +} + +impl DatagenStore { + fn from_store(store: UnifiedDatagenStore, runtime: Arc) -> Self { + Self { store, runtime } + } +} + +#[pymethods] +impl DatagenStore { + /// Open (or create) an embedded datagen log at `uri`. `shard_id` gives this + /// writer instance a stable identity for multi-writer fencing. + #[classmethod] + #[pyo3(signature = (uri, storage_options = None, shard_id = None))] + fn open( + _cls: &Bound<'_, PyType>, + py: Python<'_>, + uri: &str, + storage_options: Option>, + shard_id: Option, + ) -> PyResult { + let _ = shard_id; + let runtime = Arc::new(Runtime::new().map_err(to_py_err)?); + let store_res = py.allow_threads(|| { + runtime.block_on(UnifiedDatagenStore::open_with_options(uri, storage_options)) + }); + let store = store_res.map_err(to_py_err)?; + Ok(Self::from_store(store, runtime)) + } + + /// Connect to an existing datagen store on a remote server. + #[classmethod] + fn connect( + _cls: &Bound<'_, PyType>, + py: Python<'_>, + base_url: &str, + name: &str, + ) -> PyResult { + let runtime = Arc::new(Runtime::new().map_err(to_py_err)?); + let store_res = + py.allow_threads(|| runtime.block_on(UnifiedDatagenStore::connect(base_url, name))); + let store = store_res.map_err(to_py_err)?; + Ok(Self::from_store(store, runtime)) + } + + /// Connect to a remote datagen store, creating it if it does not exist. + #[classmethod] + #[pyo3(signature = (base_url, name, storage_options = None))] + fn connect_or_create( + _cls: &Bound<'_, PyType>, + py: Python<'_>, + base_url: &str, + name: &str, + storage_options: Option>, + ) -> PyResult { + let req = CreateDatagenStoreRequest { + name: name.to_string(), + storage_options, + }; + let runtime = Arc::new(Runtime::new().map_err(to_py_err)?); + let store_res = py.allow_threads(|| { + runtime.block_on(UnifiedDatagenStore::connect_or_create(base_url, &req)) + }); + let store = store_res.map_err(to_py_err)?; + Ok(Self::from_store(store, runtime)) + } + + /// Current store version (base dataset version). + fn version(&self) -> u64 { + self.store.version() + } + + /// Append one completed step boundary atomically. `events` is a list of + /// event dicts sharing one item/checkpoint/writer attempt, with exactly one + /// STEP_COMPLETED. Returns the new store version. + fn append_checkpoint(&mut self, py: Python<'_>, events: &Bound<'_, PyList>) -> PyResult { + let parsed = events_from_pylist(events)?; + let resp = py + .allow_threads(|| self.runtime.block_on(self.store.append_checkpoint(&parsed))) + .map_err(to_py_err)?; + Ok(resp.version) + } + + /// Append raw events as one MemWAL generation (no single-STEP_COMPLETED + /// constraint). Returns the new store version. + fn append(&mut self, py: Python<'_>, events: &Bound<'_, PyList>) -> PyResult { + let parsed = events_from_pylist(events)?; + let resp = py + .allow_threads(|| self.runtime.block_on(self.store.append(&parsed))) + .map_err(to_py_err)?; + Ok(resp.version) + } + + /// Fold an item's events into its latest state, or `None` if never started. + fn fold_item(&self, py: Python<'_>, item_id: &str) -> PyResult> { + let item = py + .allow_threads(|| self.runtime.block_on(self.store.fold_item(item_id))) + .map_err(to_py_err)?; + match item { + None => Ok(None), + Some(item) => Ok(Some(folded_item_to_py(py, &item)?)), + } + } + + /// Classify each root item id by folded lifecycle status. Missing ids (never + /// started) are absent from the returned dict. + fn root_item_statuses(&self, py: Python<'_>, root_item_ids: Vec) -> PyResult { + let statuses = py + .allow_threads(|| { + self.runtime + .block_on(self.store.root_item_statuses(&root_item_ids)) + }) + .map_err(to_py_err)?; + let dict = PyDict::new(py); + for (item_id, status) in statuses.statuses.iter() { + dict.set_item(item_id, status)?; + } + Ok(dict.into_pyobject(py)?.unbind().into()) + } + + /// All failure records for an item (the failure lens), oldest first. + fn item_failures(&self, py: Python<'_>, item_id: &str) -> PyResult { + let failures = py + .allow_threads(|| self.runtime.block_on(self.store.item_failures(item_id))) + .map_err(to_py_err)?; + let list = PyList::empty(py); + for failure in &failures { + list.append(failure_to_py(py, failure)?)?; + } + Ok(list.into_pyobject(py)?.unbind().into()) + } + + /// Materialize one FIELD_* event's blob bytes by event id, or `None`. + fn get_blob(&self, py: Python<'_>, event_id: &str) -> PyResult>> { + let bytes = py + .allow_threads(|| self.runtime.block_on(self.store.get_blob(event_id))) + .map_err(to_py_err)?; + Ok(bytes.map(|b| PyBytes::new(py, &b).unbind())) + } +} + +/// Deterministic idempotency key for a checkpoint event, so retried batches +/// dedup instead of double-appending. +#[pyfunction] +fn datagen_event_id(item_id: &str, checkpoint_id: &str, ordinal: u32) -> String { + core_datagen_event_id(item_id, checkpoint_id, ordinal) +} + +fn events_from_pylist(events: &Bound<'_, PyList>) -> PyResult> { + events + .iter() + .enumerate() + .map(|(index, item)| { + let dict = item + .downcast::() + .map_err(|_| PyTypeError::new_err(format!("events[{index}] must be a dict")))?; + event_from_dict(dict, index) + }) + .collect() +} + +fn event_from_dict(dict: &Bound<'_, PyDict>, index: usize) -> PyResult { + let value = optional_item(dict, "value")? + .map(|value| { + let value_dict = value + .downcast::() + .map_err(|_| PyTypeError::new_err("event value must be a dict"))?; + value_from_dict(value_dict) + }) + .transpose()?; + let query_tags = optional_item(dict, "query_tags")? + .map(|value| py_any_to_json(&value)) + .transpose()?; + let event_ts = optional_item(dict, "event_ts")? + .map(|value| parse_optional_datetime(Some(value.extract::()?), "event_ts")) + .transpose()? + .flatten(); + + Ok(DatagenEventDto { + event_id: required_item(dict, "event_id", index)?.extract::()?, + item_id: required_item(dict, "item_id", index)?.extract::()?, + root_item_id: required_item(dict, "root_item_id", index)?.extract::()?, + parent_item_id: optional_item(dict, "parent_item_id")? + .map(|value| value.extract::()) + .transpose()?, + item_seq: required_item(dict, "item_seq", index)?.extract::()?, + checkpoint_id: required_item(dict, "checkpoint_id", index)?.extract::()?, + event_type: required_item(dict, "event_type", index)?.extract::()?, + step_name: optional_item(dict, "step_name")? + .map(|value| value.extract::()) + .transpose()?, + step_kind: optional_item(dict, "step_kind")? + .map(|value| value.extract::()) + .transpose()?, + step_index: optional_item(dict, "step_index")? + .map(|value| value.extract::()) + .transpose()?, + enclosing_step: optional_item(dict, "enclosing_step")? + .map(|value| value.extract::()) + .transpose()?, + selector_step: optional_item(dict, "selector_step")? + .map(|value| value.extract::()) + .transpose()?, + attempt: optional_item(dict, "attempt")? + .map(|value| value.extract::()) + .transpose()? + .unwrap_or(0), + run_id: required_item(dict, "run_id", index)?.extract::()?, + writer_epoch: required_item(dict, "writer_epoch", index)?.extract::()?, + field_name: optional_item(dict, "field_name")? + .map(|value| value.extract::()) + .transpose()?, + field_type: optional_item(dict, "field_type")? + .map(|value| value.extract::()) + .transpose()?, + codec_version: optional_item(dict, "codec_version")? + .map(|value| value.extract::()) + .transpose()?, + value, + query_tags, + status: optional_item(dict, "status")? + .map(|value| value.extract::()) + .transpose()?, + error_type: optional_item(dict, "error_type")? + .map(|value| value.extract::()) + .transpose()?, + error_dump: optional_item(dict, "error_dump")? + .map(|value| value.extract::()) + .transpose()?, + traceback: optional_item(dict, "traceback")? + .map(|value| value.extract::()) + .transpose()?, + event_ts, + schema_version: optional_item(dict, "schema_version")? + .map(|value| value.extract::()) + .transpose()? + .unwrap_or(DATAGEN_SCHEMA_VERSION), + }) +} + +fn value_from_dict(dict: &Bound<'_, PyDict>) -> PyResult { + let kind = dict + .get_item("kind")? + .ok_or_else(|| PyRuntimeError::new_err("event value is missing 'kind'"))? + .extract::()?; + let inner = || { + dict.get_item("value")? + .ok_or_else(|| PyRuntimeError::new_err("event value is missing 'value'")) + }; + let mut dto = DatagenValueDto { + kind: kind.clone(), + value: None, + bytes: None, + size: None, + checksum: None, + }; + match kind.as_str() { + "int" => dto.value = Some(Value::from(inner()?.extract::()?)), + "float" => dto.value = Some(Value::from(inner()?.extract::()?)), + "bool" => dto.value = Some(Value::from(inner()?.extract::()?)), + "str" => dto.value = Some(Value::from(inner()?.extract::()?)), + "json" => dto.value = Some(py_any_to_json(&inner()?)?), + "blob" => { + let bytes = dict + .get_item("bytes")? + .filter(|value| !value.is_none()) + .map(|value| value.extract::>()) + .transpose()?; + let size = match dict.get_item("size")? { + Some(value) if !value.is_none() => Some(value.extract::()?), + _ => bytes.as_ref().map(|b| b.len() as i64), + }; + let checksum = dict + .get_item("checksum")? + .filter(|value| !value.is_none()) + .map(|value| value.extract::()) + .transpose()?; + dto.bytes = bytes; + dto.size = size; + dto.checksum = checksum; + } + other => { + return Err(PyRuntimeError::new_err(format!( + "unsupported datagen value kind '{other}'" + ))) + } + } + Ok(dto) +} + +fn py_any_to_json(value: &Bound<'_, PyAny>) -> PyResult { + let py = value.py(); + let json = PyModule::import(py, "json")?; + let text = json.call_method1("dumps", (value,))?.extract::()?; + serde_json::from_str(&text).map_err(to_py_err) +} + +fn value_to_py(py: Python<'_>, value: &DatagenValueDto) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("kind", &value.kind)?; + if value.kind == "blob" { + dict.set_item("bytes", value.bytes.as_ref().map(|b| PyBytes::new(py, b)))?; + dict.set_item("size", value.size)?; + dict.set_item("checksum", value.checksum.clone())?; + } else if let Some(inner) = &value.value { + dict.set_item("value", json_value_to_py(py, inner)?)?; + } + Ok(dict.into_pyobject(py)?.unbind().into()) +} + +fn folded_item_to_py(py: Python<'_>, item: &FoldedDatagenItemDto) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("item_id", &item.item_id)?; + dict.set_item("root_item_id", &item.root_item_id)?; + dict.set_item("parent_item_id", item.parent_item_id.clone())?; + dict.set_item("status", &item.status)?; + dict.set_item("last_item_seq", item.last_item_seq)?; + dict.set_item("last_attempt", item.last_attempt)?; + + let fields = PyDict::new(py); + for (name, state) in &item.fields { + fields.set_item(name, field_state_to_py(py, state)?)?; + } + dict.set_item("fields", fields)?; + + let trajectory = PyList::empty(py); + for cursor in &item.trajectory { + trajectory.append(cursor_to_py(py, cursor)?)?; + } + dict.set_item("trajectory", trajectory)?; + + dict.set_item( + "query_tags", + match &item.query_tags { + Some(tags) => Some(json_value_to_py(py, tags)?), + None => None, + }, + )?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + +fn field_state_to_py(py: Python<'_>, state: &DatagenFieldStateDto) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("mode", &state.mode)?; + if let Some(value) = &state.value { + dict.set_item("value", value_to_py(py, value)?)?; + } + if !state.values.is_empty() { + let list = PyList::empty(py); + for value in &state.values { + list.append(value_to_py(py, value)?)?; + } + dict.set_item("values", list)?; + } + Ok(dict.into_pyobject(py)?.unbind().into()) +} + +fn cursor_to_py(py: Python<'_>, cursor: &DatagenStepCursorDto) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("step_name", &cursor.step_name)?; + dict.set_item("step_kind", &cursor.step_kind)?; + dict.set_item("step_index", cursor.step_index)?; + dict.set_item("enclosing_step", cursor.enclosing_step.clone())?; + dict.set_item("selector_step", cursor.selector_step.clone())?; + dict.set_item("item_seq", cursor.item_seq)?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + +fn failure_to_py(py: Python<'_>, failure: &DatagenFailureDto) -> PyResult { + let dict = PyDict::new(py); + dict.set_item("at", cursor_to_py(py, &failure.at)?)?; + dict.set_item("run_id", &failure.run_id)?; + dict.set_item("attempt", failure.attempt)?; + dict.set_item("error_type", &failure.error_type)?; + dict.set_item("error_dump", failure.error_dump.clone())?; + dict.set_item("traceback", failure.traceback.clone())?; + Ok(dict.into_pyobject(py)?.unbind().into()) +} + #[pymodule] fn _internal(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_function(wrap_pyfunction!(version, m)?)?; m.add_function(wrap_pyfunction!(generate_id, m)?)?; + m.add_function(wrap_pyfunction!(datagen_event_id, m)?)?; m.add_class::()?; m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add_class::()?; Ok(()) }