use std::sync::Arc;
use arrow::array::{ArrayRef, StringBuilder, UInt64Builder};
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use arrow::error::ArrowError;
use arrow::record_batch::RecordBatch;
use polyc_eventlog::Event;
use polyc_proto::kinds;
#[derive(Debug, Clone)]
pub(crate) struct RoutineSetupRow {
pub partition: String,
pub position: u64,
pub turn_id: Option<String>,
pub routine_uid: String,
}
#[must_use]
pub(crate) fn schema() -> SchemaRef {
Arc::new(Schema::new(vec![
Field::new("partition", DataType::Utf8, false),
Field::new("position", DataType::UInt64, false),
Field::new("turn_id", DataType::Utf8, true),
Field::new("routine_uid", DataType::Utf8, false),
]))
}
pub(crate) fn decode_routine_setup_batch(
rows: &[RoutineSetupRow],
) -> Result<RecordBatch, ArrowError> {
let mut partition_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut position_b = UInt64Builder::with_capacity(rows.len());
let mut turn_id_b = StringBuilder::with_capacity(rows.len(), rows.len() * 36);
let mut routine_uid_b = StringBuilder::with_capacity(rows.len(), rows.len() * 36);
for row in rows {
partition_b.append_value(&row.partition);
position_b.append_value(row.position);
match &row.turn_id {
Some(id) => turn_id_b.append_value(id),
None => turn_id_b.append_null(),
}
routine_uid_b.append_value(&row.routine_uid);
}
let columns: Vec<ArrayRef> = vec![
Arc::new(partition_b.finish()),
Arc::new(position_b.finish()),
Arc::new(turn_id_b.finish()),
Arc::new(routine_uid_b.finish()),
];
RecordBatch::try_new(schema(), columns)
}
fn decode_routine_uid(payload: &[u8]) -> Result<String, &'static str> {
let value: serde_json::Value = serde_json::from_slice(payload).map_err(|_| "not valid JSON")?;
value
.get("routine_uid")
.and_then(serde_json::Value::as_str)
.filter(|uid| !uid.is_empty())
.map(std::borrow::ToOwned::to_owned)
.ok_or("no non-empty routine_uid field")
}
#[must_use]
pub(crate) fn decode_routine_setup_events(
partition: &str,
events: &[(u64, Event)],
) -> Vec<RoutineSetupRow> {
crate::decode::decode_typed_kind_events(
partition,
events,
&[kinds::ROUTINE_SETUP_COMPLETED],
"routine_setup",
decode_routine_uid,
|partition, position, turn_id, routine_uid| RoutineSetupRow {
partition,
position,
turn_id,
routine_uid,
},
)
}
#[cfg(test)]
mod tests {
use arrow::array::Array as _;
use super::*;
fn marker(routine_uid: &str) -> Event {
let payload = serde_json::json!({ "routine_uid": routine_uid }).to_string();
Event::trusted(kinds::ROUTINE_SETUP_COMPLETED, payload.into_bytes())
}
#[test]
fn schema_shape() {
let schema = schema();
let names: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
names,
vec!["partition", "position", "turn_id", "routine_uid"]
);
let expect = [
("partition", DataType::Utf8, false),
("position", DataType::UInt64, false),
("turn_id", DataType::Utf8, true),
("routine_uid", DataType::Utf8, false),
];
for (field, (name, ty, nullable)) in schema.fields().iter().zip(expect) {
assert_eq!(field.name(), name);
assert_eq!(field.data_type(), &ty);
assert_eq!(field.is_nullable(), nullable);
}
}
#[test]
fn decode_round_trips_a_setup_completed_marker() {
let events = vec![(7, marker("uid-standup"))];
let decoded = decode_routine_setup_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].partition, "routine-scheduler");
assert_eq!(decoded[0].position, 7);
assert_eq!(decoded[0].turn_id, None);
assert_eq!(decoded[0].routine_uid, "uid-standup");
let batch = decode_routine_setup_batch(&decoded).expect("batch build");
assert_eq!(batch.num_rows(), 1);
assert_eq!(batch.schema(), schema());
let routine_uid = batch
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(routine_uid.value(0), "uid-standup");
}
#[test]
fn empty_or_missing_routine_uid_drops_the_row() {
let events = vec![
(1, marker("")),
(
2,
Event::trusted(kinds::ROUTINE_SETUP_COMPLETED, b"{}".to_vec()),
),
];
let decoded = decode_routine_setup_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 0);
}
#[test]
fn unrelated_kind_is_not_decoded_as_setup_state() {
let events = vec![(1, Event::new(kinds::USAGE, Vec::new()))];
let decoded = decode_routine_setup_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 0);
}
#[test]
fn structurally_malformed_payload_drops_the_row() {
let events = vec![(
1,
Event::trusted(kinds::ROUTINE_SETUP_COMPLETED, vec![0xFF, 0xFE]),
)];
let decoded = decode_routine_setup_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 0);
}
}