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::events_decode::try_decode_event_payload;
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::events::v1::TurnDispatchedEvent;
#[derive(Debug, Clone)]
pub(crate) struct TurnDispatchRow {
pub partition: String,
pub position: u64,
pub turn_id: Option<String>,
pub occurrence: String,
pub visibility: &'static str,
pub visibility_source: &'static str,
pub source_turn_id: String,
pub edge_asserted_visibility: &'static str,
}
#[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("occurrence", DataType::Utf8, false),
Field::new("visibility", DataType::Utf8, false),
Field::new("visibility_source", DataType::Utf8, false),
Field::new("source_turn_id", DataType::Utf8, false),
Field::new("edge_asserted_visibility", DataType::Utf8, false),
]))
}
pub(crate) fn decode_turn_dispatch_batch(
rows: &[TurnDispatchRow],
) -> Result<RecordBatch, ArrowError> {
let mut partition_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut position_b = UInt64Builder::with_capacity(rows.len());
let mut turn_id_b = StringBuilder::with_capacity(rows.len(), rows.len() * 36);
let mut occurrence_b = StringBuilder::with_capacity(rows.len(), rows.len() * 32);
let mut visibility_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut visibility_source_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut source_turn_id_b = StringBuilder::with_capacity(rows.len(), rows.len() * 36);
let mut edge_asserted_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
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(),
}
occurrence_b.append_value(&row.occurrence);
visibility_b.append_value(row.visibility);
visibility_source_b.append_value(row.visibility_source);
source_turn_id_b.append_value(&row.source_turn_id);
edge_asserted_b.append_value(row.edge_asserted_visibility);
}
let columns: Vec<ArrayRef> = vec![
Arc::new(partition_b.finish()),
Arc::new(position_b.finish()),
Arc::new(turn_id_b.finish()),
Arc::new(occurrence_b.finish()),
Arc::new(visibility_b.finish()),
Arc::new(visibility_source_b.finish()),
Arc::new(source_turn_id_b.finish()),
Arc::new(edge_asserted_b.finish()),
];
RecordBatch::try_new(schema(), columns)
}
#[must_use]
pub(crate) fn decode_turn_dispatch_events(
partition: &str,
events: &[(u64, Event)],
) -> Vec<TurnDispatchRow> {
crate::decode::decode_typed_kind_events(
partition,
events,
&[kinds::TURN_DISPATCHED],
"turn_dispatch",
try_decode_event_payload::<TurnDispatchedEvent>,
|partition, position, turn_id, event: TurnDispatchedEvent| TurnDispatchRow {
partition,
position,
turn_id,
occurrence: event.occurrence,
visibility: polyc_proto::audience_display::recorded_visibility_label(event.visibility),
visibility_source: polyc_proto::audience_display::recorded_source_label(
event.visibility_source,
),
source_turn_id: event.source_turn_id,
edge_asserted_visibility: polyc_proto::audience_display::asserted_visibility_label(
event.edge_asserted_visibility,
),
},
)
}
#[cfg(test)]
mod tests {
use arrow::array::Array as _;
use buffa::Message as _;
use polyc_proto::proto::polychrome::events::v1::{
RecordedTurnVisibility, RecordedVisibilitySource,
};
use uuid::Uuid;
use super::*;
#[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",
"occurrence",
"visibility",
"visibility_source",
"source_turn_id",
"edge_asserted_visibility",
]
);
let expect = [
("partition", DataType::Utf8, false),
("position", DataType::UInt64, false),
("turn_id", DataType::Utf8, true),
("occurrence", DataType::Utf8, false),
("visibility", DataType::Utf8, false),
("visibility_source", DataType::Utf8, false),
("source_turn_id", DataType::Utf8, false),
("edge_asserted_visibility", 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_scheduled_origin_dispatch() {
let turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_2345);
let event = TurnDispatchedEvent {
occurrence: "daily-standup-28461600".to_owned(),
..Default::default()
};
let events = vec![(
3,
Event::new(
kinds::tagged(kinds::TURN_DISPATCHED, &turn),
event.encode_to_vec(),
),
)];
let decoded = decode_turn_dispatch_events("conv-rt", &events);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].partition, "conv-rt");
assert_eq!(decoded[0].position, 3);
assert_eq!(decoded[0].turn_id, Some(turn.to_string()));
assert_eq!(decoded[0].occurrence, "daily-standup-28461600");
let batch = decode_turn_dispatch_batch(&decoded).expect("batch build");
assert_eq!(batch.num_rows(), 1);
assert_eq!(batch.schema(), schema());
let occurrence = batch
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(occurrence.value(0), "daily-standup-28461600");
}
#[test]
fn an_inherited_stamp_names_the_turn_it_came_from() {
use polyc_proto::proto::polychrome::agent::v1::ConversationVisibility;
let turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_7777);
let source = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_8888);
let event = TurnDispatchedEvent {
visibility: buffa::EnumValue::Known(
RecordedTurnVisibility::RECORDED_TURN_VISIBILITY_DIRECT,
),
visibility_source: buffa::EnumValue::Known(
RecordedVisibilitySource::RECORDED_VISIBILITY_SOURCE_INHERITED,
),
source_turn_id: source.to_string(),
edge_asserted_visibility: buffa::EnumValue::Known(
ConversationVisibility::CONVERSATION_VISIBILITY_UNKNOWN,
),
..Default::default()
};
let events = vec![(
9,
Event::new(
kinds::tagged(kinds::TURN_DISPATCHED, &turn),
event.encode_to_vec(),
),
)];
let decoded = decode_turn_dispatch_events("conv-resume", &events);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].visibility, "direct");
assert_eq!(decoded[0].visibility_source, "inherited");
assert_eq!(decoded[0].source_turn_id, source.to_string());
assert_eq!(decoded[0].edge_asserted_visibility, "unknown");
let batch = decode_turn_dispatch_batch(&decoded).expect("batch build");
let column = |index: usize| {
batch
.column(index)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("string column")
.value(0)
.to_owned()
};
assert_eq!(column(4), "direct");
assert_eq!(column(5), "inherited");
assert_eq!(column(6), source.to_string());
assert_eq!(column(7), "unknown");
}
#[test]
fn a_record_without_a_source_reads_as_unrecorded() {
let turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_6666);
let event = TurnDispatchedEvent {
visibility: buffa::EnumValue::Known(
RecordedTurnVisibility::RECORDED_TURN_VISIBILITY_DIRECT,
),
..Default::default()
};
let events = vec![(
2,
Event::new(
kinds::tagged(kinds::TURN_DISPATCHED, &turn),
event.encode_to_vec(),
),
)];
let decoded = decode_turn_dispatch_events("conv-legacy", &events);
assert_eq!(decoded[0].visibility_source, "unrecorded");
assert_eq!(decoded[0].source_turn_id, "");
}
#[test]
fn an_unnameable_audience_value_reads_as_unrecognized() {
let turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_5555);
let event = TurnDispatchedEvent {
visibility: buffa::EnumValue::Unknown(97),
visibility_source: buffa::EnumValue::Unknown(98),
edge_asserted_visibility: buffa::EnumValue::Unknown(99),
..Default::default()
};
let events = vec![(
4,
Event::new(
kinds::tagged(kinds::TURN_DISPATCHED, &turn),
event.encode_to_vec(),
),
)];
let decoded = decode_turn_dispatch_events("conv-future", &events);
assert_eq!(decoded[0].visibility, "unrecognized");
assert_eq!(decoded[0].visibility_source, "unrecognized");
assert_eq!(decoded[0].edge_asserted_visibility, "unrecognized");
}
#[test]
fn chat_turn_has_empty_not_null_occurrence() {
let turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_9999);
let events = vec![(
1,
Event::new(kinds::tagged(kinds::TURN_DISPATCHED, &turn), Vec::new()),
)];
let decoded = decode_turn_dispatch_events("conv-chat", &events);
assert_eq!(decoded.len(), 1, "an empty payload must decode, not skip");
assert_eq!(decoded[0].occurrence, "");
}
#[test]
fn undecodable_non_empty_payload_is_skipped() {
let events = vec![(
1,
Event::new(kinds::TURN_DISPATCHED, vec![0xFF, 0xFE, 0xFD]),
)];
let decoded = decode_turn_dispatch_events("conv-corrupt", &events);
assert_eq!(decoded.len(), 0);
}
#[test]
fn unrelated_kind_is_not_decoded_as_a_dispatch() {
let events = vec![(1, Event::new(kinds::USAGE, Vec::new()))];
let decoded = decode_turn_dispatch_events("conv-unrelated", &events);
assert_eq!(decoded.len(), 0);
}
}