Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions crates/lance-context-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -337,6 +337,15 @@ pub trait DatagenStoreApi {
item_id: &str,
) -> impl Future<Output = ContextResult<Vec<DatagenFailureDto>>> + Send;

/// Raw-dump every event for a root and its projected descendants, oldest
/// first. The transport-thin read the inspection tree is folded from
/// client-side, so the same `DatagenItemTree` assembly runs for embedded and
/// remote without duplicating fold logic on the server.
fn events_for_root(
&self,
root_item_id: &str,
) -> impl Future<Output = ContextResult<Vec<DatagenEventDto>>> + 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(
Expand Down Expand Up @@ -1161,6 +1170,10 @@ pub struct FoldedDatagenItemDto {
pub trajectory: Vec<DatagenStepCursorDto>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub query_tags: Option<Value>,
/// `field_name -> event_id` for the folded blob fields, so a caller can resolve a blob by field
/// name (via `load_blob`) without recomputing an `event_id`.
#[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
pub blob_event_ids: std::collections::BTreeMap<String, String>,
}

#[derive(Debug, Serialize, Deserialize)]
Expand Down Expand Up @@ -1193,6 +1206,11 @@ pub struct ListDatagenFailuresResponse {
pub failures: Vec<DatagenFailureDto>,
}

#[derive(Debug, Serialize, Deserialize)]
pub struct ListDatagenEventsResponse {
pub events: Vec<DatagenEventDto>,
}

// ---------------------------------------------------------------------------
// Error
// ---------------------------------------------------------------------------
Expand Down
25 changes: 25 additions & 0 deletions crates/lance-context-client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -567,6 +567,15 @@ impl DatagenStoreApi for RemoteDatagenStore {
Ok(resp.failures)
}

async fn events_for_root(&self, root_item_id: &str) -> ContextResult<Vec<DatagenEventDto>> {
let resp = self
.client
.datagen_events_for_root(&self.store_name, root_item_id)
.await
.map_err(to_ctx_err)?;
Ok(resp.events)
}

async fn get_blob(&self, event_id: &str) -> ContextResult<Option<Vec<u8>>> {
self.client
.fetch_datagen_blob(&self.store_name, event_id)
Expand Down Expand Up @@ -1313,6 +1322,22 @@ impl ContextClient {
Self::handle_response(resp).await
}

/// Fetch every raw event whose root item is `root_item_id`. The client
/// folds these into a tree via `DatagenItemTree::build`; the server does no
/// fold/tree work.
pub async fn datagen_events_for_root(
&self,
name: &str,
root_item_id: &str,
) -> Result<ListDatagenEventsResponse, ClientError> {
let resp = self
.http
.get(self.url(&format!("/datagen/{}/roots/{}/events", name, root_item_id)))
.send()
.await?;
Self::handle_response(resp).await
}

/// Materialize one FIELD_* event's offloaded blob bytes by event id.
/// Returns `None` when the event or its payload is absent (server 404).
pub async fn fetch_datagen_blob(
Expand Down
43 changes: 41 additions & 2 deletions crates/lance-context-core/src/api_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -789,6 +789,13 @@ impl DatagenStoreApi for DatagenStore {
Ok(failures.iter().map(failure_to_dto).collect())
}

async fn events_for_root(&self, root_item_id: &str) -> ContextResult<Vec<DatagenEventDto>> {
let events = DatagenStore::events_for_root(self, root_item_id)
.await
.map_err(to_ctx_err)?;
Ok(events.iter().map(datagen_event_to_dto).collect())
}

async fn get_blob(&self, event_id: &str) -> ContextResult<Option<Vec<u8>>> {
DatagenStore::get_blob(self, event_id)
.await
Expand All @@ -800,7 +807,7 @@ impl DatagenStoreApi for DatagenStore {
}
}

fn datagen_events_from_dtos(events: &[DatagenEventDto]) -> ContextResult<Vec<DatagenEvent>> {
pub fn datagen_events_from_dtos(events: &[DatagenEventDto]) -> ContextResult<Vec<DatagenEvent>> {
events.iter().map(datagen_event_from_dto).collect()
}

Expand Down Expand Up @@ -920,7 +927,38 @@ fn datagen_value_to_dto(value: &DatagenValue) -> DatagenValueDto {
dto
}

fn folded_item_to_dto(item: &FoldedDatagenItem) -> FoldedDatagenItemDto {
pub fn datagen_event_to_dto(event: &DatagenEvent) -> DatagenEventDto {
DatagenEventDto {
event_id: event.event_id.clone(),
item_id: event.item_id.clone(),
root_item_id: event.root_item_id.clone(),
parent_item_id: event.parent_item_id.clone(),
item_seq: event.item_seq,
checkpoint_id: event.checkpoint_id.clone(),
event_type: event.event_type.as_str().to_string(),
step_name: event.step_name.clone(),
step_kind: event.step_kind.map(|kind| kind.as_str().to_string()),
step_index: event.step_index,
enclosing_step: event.enclosing_step.clone(),
selector_step: event.selector_step.clone(),
attempt: event.attempt,
run_id: event.run_id.clone(),
writer_epoch: event.writer_epoch.clone(),
field_name: event.field_name.clone(),
field_type: event.field_type.clone(),
codec_version: event.codec_version,
value: event.value.as_ref().map(datagen_value_to_dto),
query_tags: event.query_tags.clone(),
status: event.status.map(|status| status.as_str().to_string()),
error_type: event.error_type.clone(),
error_dump: event.error_dump.clone(),
traceback: event.traceback.clone(),
event_ts: Some(event.event_ts),
schema_version: event.schema_version,
}
}

pub fn folded_item_to_dto(item: &FoldedDatagenItem) -> FoldedDatagenItemDto {
FoldedDatagenItemDto {
item_id: item.item_id.to_string(),
root_item_id: item.root_item_id.to_string(),
Expand All @@ -935,6 +973,7 @@ fn folded_item_to_dto(item: &FoldedDatagenItem) -> FoldedDatagenItemDto {
.collect(),
trajectory: item.trajectory.ordered.iter().map(cursor_to_dto).collect(),
query_tags: item.query_tags.clone(),
blob_event_ids: item.blob_event_ids.clone(),
}
}

Expand Down
Loading
Loading