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
20 changes: 20 additions & 0 deletions crates/lance-context-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -663,6 +663,16 @@ pub struct AddRolloutRequest {
#[serde(default = "default_content_type")]
pub content_type: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_input_string: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_output_string: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rationale: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub problem_text: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub user_metadata: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_tokens: Option<Vec<i32>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_tokens: Option<Vec<i32>>,
Expand Down Expand Up @@ -742,6 +752,16 @@ pub struct RolloutRecordDto {
pub content: Option<String>,
pub content_type: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_input_string: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_output_string: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rationale: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub problem_text: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub user_metadata: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_tokens: Option<Vec<i32>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_tokens: Option<Vec<i32>>,
Expand Down
10 changes: 10 additions & 0 deletions crates/lance-context-core/src/api_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -456,6 +456,11 @@ fn rollout_record_from_add_request(r: &AddRolloutRequest) -> RolloutRecord {
created_at: r.created_at.unwrap_or_else(Utc::now),
content: r.content.clone(),
content_type: r.content_type.clone(),
model_input_string: r.model_input_string.clone(),
model_output_string: r.model_output_string.clone(),
rationale: r.rationale.clone(),
problem_text: r.problem_text.clone(),
user_metadata: r.user_metadata.clone(),
input_tokens: r.input_tokens.clone(),
output_tokens: r.output_tokens.clone(),
num_input_tokens: r.num_input_tokens,
Expand Down Expand Up @@ -498,6 +503,11 @@ pub fn rollout_record_to_dto(r: RolloutRecord) -> RolloutRecordDto {
created_at: r.created_at,
content: r.content,
content_type: r.content_type,
model_input_string: r.model_input_string,
model_output_string: r.model_output_string,
rationale: r.rationale,
problem_text: r.problem_text,
user_metadata: r.user_metadata,
input_tokens: r.input_tokens,
output_tokens: r.output_tokens,
num_input_tokens: r.num_input_tokens,
Expand Down
10 changes: 10 additions & 0 deletions crates/lance-context-core/src/rollout.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,16 @@ pub struct RolloutRecord {
pub content: Option<String>,
pub content_type: String,

// Oversized message fields offloaded via the claim-check write path. Each is
// its own nullable, individually-projectable column (not packed into
// `content`/`binary_payload`) so a reader can select one without
// materializing the rest.
pub model_input_string: Option<String>,
pub model_output_string: Option<String>,
pub rationale: Option<String>,
pub problem_text: Option<String>,
pub user_metadata: Option<String>,

// Tokens.
pub input_tokens: Option<Vec<i32>>,
pub output_tokens: Option<Vec<i32>>,
Expand Down
58 changes: 58 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1408,6 +1408,11 @@ impl RolloutStore {
let mut created_at_builder = TimestampMicrosecondBuilder::with_capacity(records.len());
let mut content_builder = LargeStringBuilder::new();
let mut content_type_builder = StringBuilder::new();
let mut model_input_string_builder = LargeStringBuilder::new();
let mut model_output_string_builder = LargeStringBuilder::new();
let mut rationale_builder = LargeStringBuilder::new();
let mut problem_text_builder = LargeStringBuilder::new();
let mut user_metadata_builder = LargeStringBuilder::new();
let mut input_tokens_builder = ListBuilder::new(Int32Builder::new());
let mut output_tokens_builder = ListBuilder::new(Int32Builder::new());
let mut num_input_tokens_builder = Int32Builder::new();
Expand Down Expand Up @@ -1442,6 +1447,11 @@ impl RolloutStore {
created_at_builder.append_value(record.created_at.timestamp_micros());
content_builder.append_option(record.content.as_deref());
content_type_builder.append_value(&record.content_type);
model_input_string_builder.append_option(record.model_input_string.as_deref());
model_output_string_builder.append_option(record.model_output_string.as_deref());
rationale_builder.append_option(record.rationale.as_deref());
problem_text_builder.append_option(record.problem_text.as_deref());
user_metadata_builder.append_option(record.user_metadata.as_deref());
append_i32_list(&mut input_tokens_builder, record.input_tokens.as_deref());
append_i32_list(&mut output_tokens_builder, record.output_tokens.as_deref());
num_input_tokens_builder.append_option(record.num_input_tokens);
Expand Down Expand Up @@ -1521,6 +1531,26 @@ impl RolloutStore {
"content_type".to_string(),
Arc::new(content_type_builder.finish()),
);
arrays_by_name.insert(
"model_input_string".to_string(),
Arc::new(model_input_string_builder.finish()),
);
arrays_by_name.insert(
"model_output_string".to_string(),
Arc::new(model_output_string_builder.finish()),
);
arrays_by_name.insert(
"rationale".to_string(),
Arc::new(rationale_builder.finish()),
);
arrays_by_name.insert(
"problem_text".to_string(),
Arc::new(problem_text_builder.finish()),
);
arrays_by_name.insert(
"user_metadata".to_string(),
Arc::new(user_metadata_builder.finish()),
);
arrays_by_name.insert(
"input_tokens".to_string(),
Arc::new(input_tokens_builder.finish()),
Expand Down Expand Up @@ -1698,6 +1728,12 @@ pub fn rollout_schema() -> Schema {
// Message content.
Field::new("content", DataType::LargeUtf8, true),
Field::new("content_type", DataType::Utf8, false),
// Claim-check offloaded message fields.
Field::new("model_input_string", DataType::LargeUtf8, true),
Field::new("model_output_string", DataType::LargeUtf8, true),
Field::new("rationale", DataType::LargeUtf8, true),
Field::new("problem_text", DataType::LargeUtf8, true),
Field::new("user_metadata", DataType::LargeUtf8, true),
// Tokens.
list_field("input_tokens", DataType::Int32),
list_field("output_tokens", DataType::Int32),
Expand Down Expand Up @@ -1788,6 +1824,13 @@ fn batch_to_rollout_records(batch: &RecordBatch) -> LanceResult<Vec<RolloutRecor
let created_at_array = column_as::<TimestampMicrosecondArray>(batch, "created_at")?;
let content_array = column_as_optional::<LargeStringArray>(batch, "content");
let content_type_array = column_as::<StringArray>(batch, "content_type")?;
let model_input_string_array =
column_as_optional::<LargeStringArray>(batch, "model_input_string");
let model_output_string_array =
column_as_optional::<LargeStringArray>(batch, "model_output_string");
let rationale_array = column_as_optional::<LargeStringArray>(batch, "rationale");
let problem_text_array = column_as_optional::<LargeStringArray>(batch, "problem_text");
let user_metadata_array = column_as_optional::<LargeStringArray>(batch, "user_metadata");
let input_tokens_array = column_as_optional::<ListArray>(batch, "input_tokens");
let output_tokens_array = column_as_optional::<ListArray>(batch, "output_tokens");
let num_input_tokens_array = column_as_optional::<Int32Array>(batch, "num_input_tokens");
Expand Down Expand Up @@ -1863,6 +1906,11 @@ fn batch_to_rollout_records(batch: &RecordBatch) -> LanceResult<Vec<RolloutRecor
created_at,
content: optional_large_string(content_array, row),
content_type: content_type_array.value(row).to_string(),
model_input_string: optional_large_string(model_input_string_array, row),
model_output_string: optional_large_string(model_output_string_array, row),
rationale: optional_large_string(rationale_array, row),
problem_text: optional_large_string(problem_text_array, row),
user_metadata: optional_large_string(user_metadata_array, row),
input_tokens: optional_i32_list(input_tokens_array, row)?,
output_tokens: optional_i32_list(output_tokens_array, row)?,
num_input_tokens: optional_i32(num_input_tokens_array, row),
Expand Down Expand Up @@ -2070,6 +2118,11 @@ mod tests {
created_at: Utc.timestamp_micros(1_700_000_000_000_000).unwrap(),
content: Some("the answer is 42".to_string()),
content_type: "text/plain".to_string(),
model_input_string: None,
model_output_string: None,
rationale: None,
problem_text: None,
user_metadata: None,
input_tokens: Some(vec![10, 11, 12]),
output_tokens: Some(vec![20, 21]),
num_input_tokens: Some(3),
Expand Down Expand Up @@ -2110,6 +2163,11 @@ mod tests {
created_at: Utc.timestamp_micros(1_700_000_000_500_000).unwrap(),
content: None,
content_type: "application/octet-stream".to_string(),
model_input_string: None,
model_output_string: None,
rationale: None,
problem_text: None,
user_metadata: None,
input_tokens: None,
output_tokens: None,
num_input_tokens: None,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,11 @@ fn rec(id: &str) -> RolloutRecord {
created_at: chrono::Utc::now(),
content: Some("x".to_string()),
content_type: "text/plain".to_string(),
model_input_string: None,
model_output_string: None,
rationale: None,
problem_text: None,
user_metadata: None,
input_tokens: None,
output_tokens: None,
num_input_tokens: None,
Expand Down
5 changes: 5 additions & 0 deletions crates/lance-context-master/src/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -511,6 +511,11 @@ mod tests {
created_at: Utc::now(),
content: Some("answer".to_string()),
content_type: if with_blob { "image/png" } else { "text/plain" }.to_string(),
model_input_string: None,
model_output_string: None,
rationale: None,
problem_text: None,
user_metadata: None,
input_tokens: None,
output_tokens: Some(vec![1, 2]),
num_input_tokens: None,
Expand Down
5 changes: 5 additions & 0 deletions crates/lance-context-master/src/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -761,6 +761,11 @@ mod tests {
created_at: Utc.timestamp_micros(1_700_000_000_000_000).unwrap(),
content: Some("x".to_string()),
content_type: "text/plain".to_string(),
model_input_string: None,
model_output_string: None,
rationale: None,
problem_text: None,
user_metadata: None,
input_tokens: None,
output_tokens: None,
num_input_tokens: None,
Expand Down
10 changes: 10 additions & 0 deletions crates/lance-context-server/src/routes/rollouts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -530,6 +530,11 @@ fn rollout_record_from_add_request(r: &AddRolloutRequest) -> RolloutRecord {
created_at: r.created_at.unwrap_or_else(Utc::now),
content: r.content.clone(),
content_type: r.content_type.clone(),
model_input_string: r.model_input_string.clone(),
model_output_string: r.model_output_string.clone(),
rationale: r.rationale.clone(),
problem_text: r.problem_text.clone(),
user_metadata: r.user_metadata.clone(),
input_tokens: r.input_tokens.clone(),
output_tokens: r.output_tokens.clone(),
num_input_tokens: r.num_input_tokens,
Expand Down Expand Up @@ -571,6 +576,11 @@ fn rollout_record_to_dto(r: RolloutRecord) -> RolloutRecordDto {
created_at: r.created_at,
content: r.content,
content_type: r.content_type,
model_input_string: r.model_input_string,
model_output_string: r.model_output_string,
rationale: r.rationale,
problem_text: r.problem_text,
user_metadata: r.user_metadata,
input_tokens: r.input_tokens,
output_tokens: r.output_tokens,
num_input_tokens: r.num_input_tokens,
Expand Down
10 changes: 10 additions & 0 deletions specs/rollout-schema-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,16 @@ One row per **message** in a rollout (assistant turn, tool call, grade, or artif
| `content` | `LargeUtf8` | message text |
| `content_type` | `Utf8` | MIME type |

**Claim-check offloaded message fields** *(oversized rollout-message fields, each an individually-projectable nullable column)*

| Column | Arrow type | Purpose |
|---|---|---|
| `model_input_string` | `LargeUtf8` | rendered model prompt string |
| `model_output_string` | `LargeUtf8` | raw model completion string |
| `rationale` | `LargeUtf8` | grader rationale |
| `problem_text` | `LargeUtf8` | source problem text |
| `user_metadata` | `LargeUtf8` | harness-supplied per-message metadata blob |

**Tokens** *(first-class variable-length arrays)*

| Column | Arrow type | Purpose |
Expand Down
Loading