use std::collections::HashMap;
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::proto::polychrome::events::v1::{
RoutineFireOutcome, RoutineFireOutcomeEvent, RoutineFiredEvent,
};
#[derive(Debug, Clone)]
pub(crate) struct FireRow {
pub partition: String,
pub position: u64,
pub turn_id: Option<String>,
pub routine: String,
pub occurrence: String,
pub scheduled_at_ms: u64,
pub fired_at_ms: u64,
pub routine_uid: String,
pub outcome: Option<String>,
pub grant_drift_tools: Option<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", DataType::Utf8, false),
Field::new("occurrence", DataType::Utf8, false),
Field::new("scheduled_at_ms", DataType::UInt64, false),
Field::new("fired_at_ms", DataType::UInt64, false),
Field::new("routine_uid", DataType::Utf8, false),
Field::new("outcome", DataType::Utf8, true),
Field::new("grant_drift_tools", DataType::Utf8, true),
]))
}
pub(crate) fn decode_fires_batch(rows: &[FireRow]) -> 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_b = StringBuilder::with_capacity(rows.len(), rows.len() * 24);
let mut occurrence_b = StringBuilder::with_capacity(rows.len(), rows.len() * 32);
let mut scheduled_at_ms_b = UInt64Builder::with_capacity(rows.len());
let mut fired_at_ms_b = UInt64Builder::with_capacity(rows.len());
let mut routine_uid_b = StringBuilder::with_capacity(rows.len(), rows.len() * 36);
let mut outcome_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut grant_drift_tools_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(),
}
routine_b.append_value(&row.routine);
occurrence_b.append_value(&row.occurrence);
scheduled_at_ms_b.append_value(row.scheduled_at_ms);
fired_at_ms_b.append_value(row.fired_at_ms);
routine_uid_b.append_value(&row.routine_uid);
match &row.outcome {
Some(outcome) => outcome_b.append_value(outcome),
None => outcome_b.append_null(),
}
match &row.grant_drift_tools {
Some(drift) => grant_drift_tools_b.append_value(drift),
None => grant_drift_tools_b.append_null(),
}
}
let columns: Vec<ArrayRef> = vec![
Arc::new(partition_b.finish()),
Arc::new(position_b.finish()),
Arc::new(turn_id_b.finish()),
Arc::new(routine_b.finish()),
Arc::new(occurrence_b.finish()),
Arc::new(scheduled_at_ms_b.finish()),
Arc::new(fired_at_ms_b.finish()),
Arc::new(routine_uid_b.finish()),
Arc::new(outcome_b.finish()),
Arc::new(grant_drift_tools_b.finish()),
];
RecordBatch::try_new(schema(), columns)
}
#[must_use]
const fn outcome_str(outcome: RoutineFireOutcome) -> Option<&'static str> {
match outcome {
RoutineFireOutcome::Ok => Some("ok"),
RoutineFireOutcome::Paused => Some("paused"),
RoutineFireOutcome::RetryableError => Some("retryable_error"),
RoutineFireOutcome::TerminalError => Some("terminal_error"),
RoutineFireOutcome::TimedOut => Some("timed_out"),
RoutineFireOutcome::StoppedUngranted => Some("stopped_ungranted"),
RoutineFireOutcome::Unspecified => None,
}
}
#[must_use]
fn fold_outcome_events(
partition: &str,
events: &[(u64, Event)],
) -> HashMap<(String, String), (&'static str, String)> {
crate::decode::decode_typed_kind_events(
partition,
events,
&[polyc_proto::kinds::ROUTINE_FIRE_OUTCOME],
"fires",
try_decode_event_payload::<RoutineFireOutcomeEvent>,
|_partition, _position, _turn_id, decoded: RoutineFireOutcomeEvent| decoded,
)
.into_iter()
.filter_map(|decoded| {
let outcome = decoded.outcome.as_known().unwrap_or_default();
let drift =
serde_json::to_string(&decoded.grant_drift_tools).unwrap_or_else(|_| "[]".to_owned());
outcome_str(outcome).map(|s| ((decoded.routine, decoded.occurrence), (s, drift)))
})
.collect()
}
#[must_use]
pub(crate) fn decode_fires_events(partition: &str, events: &[(u64, Event)]) -> Vec<FireRow> {
let outcomes = fold_outcome_events(partition, events);
crate::decode::decode_typed_kind_events(
partition,
events,
&[polyc_proto::kinds::ROUTINE_FIRED],
"fires",
try_decode_event_payload::<RoutineFiredEvent>,
|partition, position, turn_id, decoded: RoutineFiredEvent| {
let matched = outcomes.get(&(decoded.routine.clone(), decoded.occurrence.clone()));
FireRow {
partition,
position,
turn_id,
routine: decoded.routine,
occurrence: decoded.occurrence,
scheduled_at_ms: decoded.scheduled_at_ms,
fired_at_ms: decoded.fired_at_ms,
routine_uid: decoded.routine_uid,
outcome: matched.map(|(s, _)| (*s).to_owned()),
grant_drift_tools: matched.map(|(_, drift)| drift.clone()),
}
},
)
}
#[cfg(test)]
mod tests {
use arrow::array::Array as _;
use buffa::Message as _;
use polyc_proto::kinds;
use super::*;
fn sample_event() -> RoutineFiredEvent {
RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-28461600".to_string(),
scheduled_at_ms: 1_753_300_800_000,
fired_at_ms: 1_753_300_805_000,
routine_uid: "uid-daily-standup".to_string(),
..Default::default()
}
}
#[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",
"occurrence",
"scheduled_at_ms",
"fired_at_ms",
"routine_uid",
"outcome",
"grant_drift_tools",
]
);
let expect = [
("partition", DataType::Utf8, false),
("position", DataType::UInt64, false),
("turn_id", DataType::Utf8, true),
("routine", DataType::Utf8, false),
("occurrence", DataType::Utf8, false),
("scheduled_at_ms", DataType::UInt64, false),
("fired_at_ms", DataType::UInt64, false),
("routine_uid", DataType::Utf8, false),
("outcome", DataType::Utf8, true),
("grant_drift_tools", DataType::Utf8, true),
];
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_routine_fired_marker() {
let ev = sample_event();
let events = vec![(9, Event::trusted(kinds::ROUTINE_FIRED, ev.encode_to_vec()))];
let decoded = decode_fires_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].partition, "routine-scheduler");
assert_eq!(decoded[0].position, 9);
assert_eq!(decoded[0].turn_id, None);
assert_eq!(decoded[0].routine, "daily-standup");
assert_eq!(decoded[0].occurrence, "daily-standup-28461600");
assert_eq!(decoded[0].scheduled_at_ms, 1_753_300_800_000);
assert_eq!(decoded[0].fired_at_ms, 1_753_300_805_000);
assert_eq!(decoded[0].routine_uid, "uid-daily-standup");
assert_eq!(
decoded[0].outcome, None,
"no sibling routine_fire_outcome event was appended"
);
assert_eq!(
decoded[0].grant_drift_tools, None,
"grant_drift_tools is NULL exactly when outcome is NULL"
);
let batch = decode_fires_batch(&decoded).expect("batch build");
assert_eq!(batch.num_rows(), 1);
assert_eq!(batch.schema(), schema());
let routine = batch
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(routine.value(0), "daily-standup");
let scheduled_at_ms = batch
.column(5)
.as_any()
.downcast_ref::<arrow::array::UInt64Array>()
.unwrap();
assert_eq!(scheduled_at_ms.value(0), 1_753_300_800_000);
let routine_uid = batch
.column(7)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(routine_uid.value(0), "uid-daily-standup");
}
#[test]
fn decode_folds_a_sibling_outcome_event_into_its_markers_row() {
let marker = sample_event();
let outcome_ev = RoutineFireOutcomeEvent {
routine: marker.routine.clone(),
occurrence: marker.occurrence.clone(),
outcome: RoutineFireOutcome::Ok.into(),
fired_at_ms: marker.fired_at_ms,
..Default::default()
};
let events = vec![
(
9,
Event::trusted(kinds::ROUTINE_FIRED, marker.encode_to_vec()),
),
(
10,
Event::trusted(kinds::ROUTINE_FIRE_OUTCOME, outcome_ev.encode_to_vec()),
),
];
let decoded = decode_fires_events("routine-scheduler", &events);
assert_eq!(
decoded.len(),
1,
"the outcome event mints no row of its own"
);
assert_eq!(decoded[0].outcome.as_deref(), Some("ok"));
assert_eq!(
decoded[0].grant_drift_tools.as_deref(),
Some("[]"),
"an outcome event with no drifted tool renders an empty JSON array, not NULL"
);
}
#[test]
fn decode_folds_a_stopped_ungranted_abort_and_its_drifted_tools() {
let marker = sample_event();
let outcome_ev = RoutineFireOutcomeEvent {
routine: marker.routine.clone(),
occurrence: marker.occurrence.clone(),
outcome: RoutineFireOutcome::StoppedUngranted.into(),
fired_at_ms: marker.fired_at_ms,
grant_drift_tools: vec!["fs_write".to_string(), "http_post".to_string()],
..Default::default()
};
let events = vec![
(
9,
Event::trusted(kinds::ROUTINE_FIRED, marker.encode_to_vec()),
),
(
10,
Event::trusted(kinds::ROUTINE_FIRE_OUTCOME, outcome_ev.encode_to_vec()),
),
];
let decoded = decode_fires_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].outcome.as_deref(), Some("stopped_ungranted"));
assert_eq!(
decoded[0].grant_drift_tools.as_deref(),
Some(r#"["fs_write","http_post"]"#)
);
let batch = decode_fires_batch(&decoded).expect("batch build");
let drift = batch
.column(9)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(drift.value(0), r#"["fs_write","http_post"]"#);
}
#[test]
fn decode_leaves_outcome_null_when_no_sibling_event_exists() {
let marker = sample_event();
let events = vec![(
9,
Event::trusted(kinds::ROUTINE_FIRED, marker.encode_to_vec()),
)];
let decoded = decode_fires_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].outcome, None);
}
#[test]
fn unrelated_kind_is_not_decoded_as_a_fire() {
let events = vec![(1, Event::new(kinds::USAGE, Vec::new()))];
let decoded = decode_fires_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_FIRED, vec![0xFF, 0xFE, 0xFD]),
)];
let decoded = decode_fires_events("routine-scheduler", &events);
assert_eq!(decoded.len(), 0);
}
#[test]
fn empty_events_yield_zero_rows() {
let decoded = decode_fires_events("routine-scheduler", &[]);
assert_eq!(decoded.len(), 0);
}
}