use std::sync::Arc;
use arrow::array::{ArrayRef, BinaryBuilder, StringBuilder, UInt64Builder};
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use arrow::error::ArrowError;
use arrow::record_batch::RecordBatch;
use polyc_crypto::approval::{
VerifiedRoutineCreated, VerifiedRoutineDeleted, VerifiedRoutinePaused, VerifiedRoutineResumed,
VerifiedRoutineScopeChanged, verify_routine_created, verify_routine_deleted,
verify_routine_paused, verify_routine_resumed, verify_routine_scope_changed,
};
use polyc_eventlog::Event;
use polyc_proto::kinds;
#[derive(Debug, Clone)]
pub(crate) struct RoutineLifecycleRow {
pub partition: String,
pub position: u64,
pub turn_id: Option<String>,
pub phase: String,
pub routine: String,
pub actor_persona: String,
pub conversation_id: String,
pub at_ms: u64,
pub reason: Option<String>,
pub scope: Option<String>,
pub channel: Option<String>,
pub signer_public_key: Vec<u8>,
}
#[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("phase", DataType::Utf8, false),
Field::new("routine", DataType::Utf8, false),
Field::new("actor_persona", DataType::Utf8, false),
Field::new("conversation_id", DataType::Utf8, false),
Field::new("at_ms", DataType::UInt64, false),
Field::new("reason", DataType::Utf8, true),
Field::new("scope", DataType::Utf8, true),
Field::new("channel", DataType::Utf8, true),
Field::new("signer_public_key", DataType::Binary, false),
]))
}
pub(crate) fn decode_routine_lifecycle_batch(
rows: &[RoutineLifecycleRow],
) -> 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 phase_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut routine_b = StringBuilder::with_capacity(rows.len(), rows.len() * 24);
let mut actor_persona_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut conversation_id_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut at_ms_b = UInt64Builder::with_capacity(rows.len());
let mut reason_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut scope_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut channel_b = StringBuilder::with_capacity(rows.len(), rows.len() * 4);
let mut signer_public_key_b = BinaryBuilder::with_capacity(rows.len(), rows.len() * 32);
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(),
}
phase_b.append_value(&row.phase);
routine_b.append_value(&row.routine);
actor_persona_b.append_value(&row.actor_persona);
conversation_id_b.append_value(&row.conversation_id);
at_ms_b.append_value(row.at_ms);
match &row.reason {
Some(v) => reason_b.append_value(v),
None => reason_b.append_null(),
}
match &row.scope {
Some(v) => scope_b.append_value(v),
None => scope_b.append_null(),
}
match &row.channel {
Some(v) => channel_b.append_value(v),
None => channel_b.append_null(),
}
signer_public_key_b.append_value(&row.signer_public_key);
}
let columns: Vec<ArrayRef> = vec![
Arc::new(partition_b.finish()),
Arc::new(position_b.finish()),
Arc::new(turn_id_b.finish()),
Arc::new(phase_b.finish()),
Arc::new(routine_b.finish()),
Arc::new(actor_persona_b.finish()),
Arc::new(conversation_id_b.finish()),
Arc::new(at_ms_b.finish()),
Arc::new(reason_b.finish()),
Arc::new(scope_b.finish()),
Arc::new(channel_b.finish()),
Arc::new(signer_public_key_b.finish()),
];
RecordBatch::try_new(schema(), columns)
}
fn is_trusted(key: &[u8], trusted_signers: &[Vec<u8>]) -> bool {
trusted_signers.iter().any(|k| k.as_slice() == key)
}
fn warn_verification_failed(partition: &str, position: u64, kind_base: &str) {
tracing::warn!(
partition,
position,
kind = kind_base,
table = "routine_lifecycle",
"routine_lifecycle: dropping row — payload failed verify_routine_* \
(malformed payload or a signature that doesn't check out)"
);
}
fn warn_untrusted_signer(
partition: &str,
position: u64,
kind_base: &str,
signer_public_key: &[u8],
) {
tracing::warn!(
partition,
position,
kind = kind_base,
table = "routine_lifecycle",
signer_public_key = %polyc_crypto::hex::lower(signer_public_key),
"routine_lifecycle: dropping row — signer_public_key is not in trusted_signers"
);
}
fn verify_trusted<V>(
partition: &str,
position: u64,
kind_base: &str,
payload: &[u8],
trusted_signers: &[Vec<u8>],
verify: impl FnOnce(&[u8]) -> Option<V>,
signer_public_key: impl Fn(&V) -> &[u8],
) -> Option<V> {
let Some(v) = verify(payload) else {
warn_verification_failed(partition, position, kind_base);
return None;
};
if is_trusted(signer_public_key(&v), trusted_signers) {
Some(v)
} else {
warn_untrusted_signer(partition, position, kind_base, signer_public_key(&v));
None
}
}
struct PhaseSpec<V> {
phase: &'static str,
verify: fn(&[u8]) -> Option<V>,
signer_public_key: fn(&V) -> &[u8],
#[allow(clippy::type_complexity)]
fields: fn(
V,
) -> (
String,
String,
String,
u64,
Option<String>,
Option<String>,
Option<String>,
Vec<u8>,
),
}
fn decode_phase_row<V>(
partition: &str,
position: u64,
turn_id: Option<String>,
payload: &[u8],
trusted_signers: &[Vec<u8>],
spec: &PhaseSpec<V>,
) -> Option<RoutineLifecycleRow> {
let v = verify_trusted(
partition,
position,
spec.phase,
payload,
trusted_signers,
spec.verify,
spec.signer_public_key,
)?;
let (routine, actor_persona, conversation_id, at_ms, reason, scope, channel, signer_public_key) =
(spec.fields)(v);
Some(RoutineLifecycleRow {
partition: partition.to_string(),
position,
turn_id,
phase: spec.phase.to_string(),
routine,
actor_persona,
conversation_id,
at_ms,
reason,
scope,
channel,
signer_public_key,
})
}
#[allow(clippy::type_complexity)] fn created_fields(
v: VerifiedRoutineCreated,
) -> (
String,
String,
String,
u64,
Option<String>,
Option<String>,
Option<String>,
Vec<u8>,
) {
(
v.routine,
v.creator_persona,
v.conversation_id,
v.created_at_ms,
None,
None,
None,
v.signer_public_key,
)
}
#[allow(clippy::type_complexity)] fn paused_fields(
v: VerifiedRoutinePaused,
) -> (
String,
String,
String,
u64,
Option<String>,
Option<String>,
Option<String>,
Vec<u8>,
) {
(
v.routine,
v.actor_persona,
v.conversation_id,
v.paused_at_ms,
v.reason,
None,
Some(v.channel),
v.signer_public_key,
)
}
#[allow(clippy::type_complexity)] fn resumed_fields(
v: VerifiedRoutineResumed,
) -> (
String,
String,
String,
u64,
Option<String>,
Option<String>,
Option<String>,
Vec<u8>,
) {
(
v.routine,
v.actor_persona,
v.conversation_id,
v.resumed_at_ms,
None,
None,
Some(v.channel),
v.signer_public_key,
)
}
#[allow(clippy::type_complexity)] fn deleted_fields(
v: VerifiedRoutineDeleted,
) -> (
String,
String,
String,
u64,
Option<String>,
Option<String>,
Option<String>,
Vec<u8>,
) {
(
v.routine,
v.actor_persona,
v.conversation_id,
v.deleted_at_ms,
None,
None,
Some(v.channel),
v.signer_public_key,
)
}
#[allow(clippy::type_complexity)] fn scope_changed_fields(
v: VerifiedRoutineScopeChanged,
) -> (
String,
String,
String,
u64,
Option<String>,
Option<String>,
Option<String>,
Vec<u8>,
) {
(
v.routine,
v.actor_persona,
v.conversation_id,
v.changed_at_ms,
None,
Some(v.scope),
None,
v.signer_public_key,
)
}
#[must_use]
pub(crate) fn decode_routine_lifecycle_events(
partition: &str,
events: &[(u64, Event)],
trusted_signers: &[Vec<u8>],
) -> Vec<RoutineLifecycleRow> {
events
.iter()
.filter_map(|(position, event)| {
let (base, turn_id) = kinds::parse(&event.kind);
let turn_id = turn_id.map(|id| id.to_string());
if base == kinds::ROUTINE_CREATED {
decode_phase_row(
partition,
*position,
turn_id,
&event.payload,
trusted_signers,
&PhaseSpec {
phase: kinds::ROUTINE_CREATED,
verify: verify_routine_created,
signer_public_key: |v| v.signer_public_key.as_slice(),
fields: created_fields,
},
)
} else if base == kinds::ROUTINE_PAUSED {
decode_phase_row(
partition,
*position,
turn_id,
&event.payload,
trusted_signers,
&PhaseSpec {
phase: kinds::ROUTINE_PAUSED,
verify: verify_routine_paused,
signer_public_key: |v| v.signer_public_key.as_slice(),
fields: paused_fields,
},
)
} else if base == kinds::ROUTINE_RESUMED {
decode_phase_row(
partition,
*position,
turn_id,
&event.payload,
trusted_signers,
&PhaseSpec {
phase: kinds::ROUTINE_RESUMED,
verify: verify_routine_resumed,
signer_public_key: |v| v.signer_public_key.as_slice(),
fields: resumed_fields,
},
)
} else if base == kinds::ROUTINE_DELETED {
decode_phase_row(
partition,
*position,
turn_id,
&event.payload,
trusted_signers,
&PhaseSpec {
phase: kinds::ROUTINE_DELETED,
verify: verify_routine_deleted,
signer_public_key: |v| v.signer_public_key.as_slice(),
fields: deleted_fields,
},
)
} else if base == kinds::ROUTINE_SCOPE_CHANGED {
decode_phase_row(
partition,
*position,
turn_id,
&event.payload,
trusted_signers,
&PhaseSpec {
phase: kinds::ROUTINE_SCOPE_CHANGED,
verify: verify_routine_scope_changed,
signer_public_key: |v| v.signer_public_key.as_slice(),
fields: scope_changed_fields,
},
)
} else {
None
}
})
.collect()
}
#[cfg(test)]
mod tests {
use arrow::array::Array as _;
use polyc_crypto::approval::{
ApprovalSigner, routine_created_payload, routine_deleted_payload, routine_paused_payload,
routine_resumed_payload, routine_scope_changed_payload,
};
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",
"phase",
"routine",
"actor_persona",
"conversation_id",
"at_ms",
"reason",
"scope",
"channel",
"signer_public_key",
]
);
let expect = [
("partition", DataType::Utf8, false),
("position", DataType::UInt64, false),
("turn_id", DataType::Utf8, true),
("phase", DataType::Utf8, false),
("routine", DataType::Utf8, false),
("actor_persona", DataType::Utf8, false),
("conversation_id", DataType::Utf8, false),
("at_ms", DataType::UInt64, false),
("reason", DataType::Utf8, true),
("scope", DataType::Utf8, true),
("channel", DataType::Utf8, true),
("signer_public_key", DataType::Binary, 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]
#[allow(clippy::too_many_lines)] fn decode_round_trips_all_five_signed_lifecycle_kinds() {
let signer = ApprovalSigner::from_seed(1);
let trusted_signers = vec![signer.public_key_bytes()];
let (created, _, _) = routine_created_payload(
"daily-standup",
"persona-1",
"conv-1",
"call-1",
"hash-1",
1_000,
&signer,
);
let (paused, _, _) = routine_paused_payload(
"daily-standup",
"persona-1",
"conv-2",
"call-2",
"hash-2",
"chat",
2_000,
Some("rotating out old announcements"),
&signer,
);
let (resumed, _, _) = routine_resumed_payload(
"daily-standup",
"persona-1",
"conv-3",
"call-3",
"hash-3",
"chat",
3_000,
&signer,
);
let (deleted, _, _) = routine_deleted_payload(
"daily-standup",
"persona-1",
"conv-4",
"call-4",
"hash-4",
"chat",
4_000,
&signer,
);
let (scope_changed, _, _) = routine_scope_changed_payload(
"daily-standup",
"persona-1",
"conv-5",
"call-5",
"hash-5",
"public",
5_000,
&signer,
);
let events = vec![
(10, Event::new(kinds::ROUTINE_CREATED, created)),
(11, Event::new(kinds::ROUTINE_PAUSED, paused)),
(12, Event::new(kinds::ROUTINE_RESUMED, resumed)),
(13, Event::new(kinds::ROUTINE_DELETED, deleted)),
(14, Event::new(kinds::ROUTINE_SCOPE_CHANGED, scope_changed)),
];
let decoded =
decode_routine_lifecycle_events("routine-scheduler", &events, &trusted_signers);
assert_eq!(decoded.len(), 5);
assert_eq!(decoded[0].phase, "routine_created");
assert_eq!(decoded[0].routine, "daily-standup");
assert_eq!(decoded[0].actor_persona, "persona-1");
assert_eq!(decoded[0].conversation_id, "conv-1");
assert_eq!(decoded[0].at_ms, 1_000);
assert_eq!(decoded[0].reason, None);
assert_eq!(decoded[0].scope, None);
assert_eq!(decoded[0].channel, None);
assert_eq!(decoded[1].phase, "routine_paused");
assert_eq!(decoded[1].at_ms, 2_000);
assert_eq!(
decoded[1].reason.as_deref(),
Some("rotating out old announcements")
);
assert_eq!(decoded[1].scope, None);
assert_eq!(decoded[1].channel.as_deref(), Some("chat"));
assert_eq!(decoded[2].phase, "routine_resumed");
assert_eq!(decoded[2].at_ms, 3_000);
assert_eq!(decoded[2].reason, None);
assert_eq!(decoded[2].scope, None);
assert_eq!(decoded[2].channel.as_deref(), Some("chat"));
assert_eq!(decoded[3].phase, "routine_deleted");
assert_eq!(decoded[3].at_ms, 4_000);
assert_eq!(decoded[3].reason, None);
assert_eq!(decoded[3].scope, None);
assert_eq!(decoded[3].channel.as_deref(), Some("chat"));
assert_eq!(decoded[4].phase, "routine_scope_changed");
assert_eq!(decoded[4].routine, "daily-standup");
assert_eq!(decoded[4].actor_persona, "persona-1");
assert_eq!(decoded[4].conversation_id, "conv-5");
assert_eq!(decoded[4].at_ms, 5_000);
assert_eq!(decoded[4].reason, None);
assert_eq!(decoded[4].scope.as_deref(), Some("public"));
assert_eq!(decoded[4].channel, None);
let batch = decode_routine_lifecycle_batch(&decoded).expect("batch build");
assert_eq!(batch.num_rows(), 5);
assert_eq!(batch.schema(), schema());
let phase = batch
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(phase.value(0), "routine_created");
assert_eq!(phase.value(3), "routine_deleted");
assert_eq!(phase.value(4), "routine_scope_changed");
let scope = batch
.column(9)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert!(scope.is_null(0));
assert_eq!(scope.value(4), "public");
}
#[test]
fn scope_changed_carries_the_target_scope() {
let signer = ApprovalSigner::from_seed(6);
let trusted_signers = vec![signer.public_key_bytes()];
let (payload, _, _) = routine_scope_changed_payload(
"weekly-digest",
"persona-2",
"conv-b",
"call-6",
"hash-6",
"private",
6_000,
&signer,
);
let events = vec![(1, Event::new(kinds::ROUTINE_SCOPE_CHANGED, payload))];
let decoded =
decode_routine_lifecycle_events("routine-scheduler", &events, &trusted_signers);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].phase, "routine_scope_changed");
assert_eq!(decoded[0].scope.as_deref(), Some("private"));
assert_eq!(decoded[0].reason, None);
}
#[test]
fn paused_without_a_reason_has_null_reason() {
let signer = ApprovalSigner::from_seed(2);
let trusted_signers = vec![signer.public_key_bytes()];
let (paused, _, _) = routine_paused_payload(
"weekly-digest",
"persona-2",
"conv-a",
"call-7",
"hash-7",
"chat",
5_000,
None,
&signer,
);
let events = vec![(1, Event::new(kinds::ROUTINE_PAUSED, paused))];
let decoded =
decode_routine_lifecycle_events("routine-scheduler", &events, &trusted_signers);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].reason, None);
}
#[test]
fn malformed_payload_drops_the_row() {
let events = vec![(1, Event::new(kinds::ROUTINE_CREATED, b"not json".to_vec()))];
let decoded = decode_routine_lifecycle_events("routine-scheduler", &events, &[]);
assert_eq!(decoded.len(), 0);
}
#[test]
fn untrusted_signer_drops_the_row() {
let untrusted = ApprovalSigner::from_seed(3);
let trusted = ApprovalSigner::from_seed(4);
let (created, _, _) = routine_created_payload(
"daily-standup",
"persona-1",
"conv-1",
"call-8",
"hash-8",
1_000,
&untrusted,
);
let events = vec![(1, Event::new(kinds::ROUTINE_CREATED, created))];
let decoded = decode_routine_lifecycle_events(
"routine-scheduler",
&events,
&[trusted.public_key_bytes()],
);
assert_eq!(
decoded.len(),
0,
"an untrusted-signer lifecycle record must not surface as a row"
);
}
#[test]
fn tampered_payload_drops_the_row() {
let signer = ApprovalSigner::from_seed(5);
let (created, _, _) = routine_created_payload(
"daily-standup",
"persona-1",
"conv-1",
"call-9",
"hash-9",
1_000,
&signer,
);
let mut tampered: serde_json::Value = serde_json::from_slice(&created).unwrap();
tampered["routine"] = serde_json::json!("evil-routine");
let events = vec![(
1,
Event::new(kinds::ROUTINE_CREATED, tampered.to_string().into_bytes()),
)];
let decoded = decode_routine_lifecycle_events(
"routine-scheduler",
&events,
&[signer.public_key_bytes()],
);
assert_eq!(decoded.len(), 0);
}
#[test]
#[tracing_test::traced_test]
fn untrusted_signer_drop_is_logged() {
let untrusted = ApprovalSigner::from_seed(9);
let trusted = ApprovalSigner::from_seed(10);
let (created, _, _) = routine_created_payload(
"daily-standup",
"persona-1",
"conv-1",
"call-10",
"hash-10",
1_000,
&untrusted,
);
let events = vec![(42, Event::new(kinds::ROUTINE_CREATED, created))];
let decoded = decode_routine_lifecycle_events(
"routine-scheduler",
&events,
&[trusted.public_key_bytes()],
);
assert_eq!(decoded.len(), 0);
assert!(
logs_contain("routine_lifecycle"),
"the drop must be logged, not silent"
);
assert!(
logs_contain("routine-scheduler"),
"the log line must name the partition"
);
assert!(
logs_contain("trusted_signers"),
"the log line must name the drop reason"
);
assert!(
logs_contain(&polyc_crypto::hex::lower(&untrusted.public_key_bytes())),
"the log line must name the rejected signer's key"
);
}
#[test]
#[tracing_test::traced_test]
fn verification_failure_drop_is_logged() {
let events = vec![(7, Event::new(kinds::ROUTINE_CREATED, b"not json".to_vec()))];
let decoded = decode_routine_lifecycle_events("routine-scheduler", &events, &[]);
assert_eq!(decoded.len(), 0);
assert!(
logs_contain("routine_lifecycle"),
"the drop must be logged, not silent"
);
assert!(
logs_contain("routine-scheduler"),
"the log line must name the partition"
);
assert!(
logs_contain("verify_routine_"),
"the log line must name the verification-failure drop reason"
);
}
#[test]
fn unrelated_kind_is_not_decoded() {
let events = vec![(1, Event::new(kinds::ROUTINE_FIRED, Vec::new()))];
let decoded = decode_routine_lifecycle_events("routine-scheduler", &events, &[]);
assert_eq!(decoded.len(), 0);
}
#[test]
fn empty_events_yield_zero_rows() {
let decoded = decode_routine_lifecycle_events("routine-scheduler", &[], &[]);
assert_eq!(decoded.len(), 0);
}
#[test]
fn retired_signer_key_still_decodes() {
let retired = ApprovalSigner::from_seed(7);
let current = ApprovalSigner::from_seed(8);
let (created, _, _) = routine_created_payload(
"daily-standup",
"persona-1",
"conv-1",
"call-11",
"hash-11",
1_000,
&retired,
);
let events = vec![(1, Event::new(kinds::ROUTINE_CREATED, created))];
let trusted_signers = vec![current.public_key_bytes(), retired.public_key_bytes()];
let decoded =
decode_routine_lifecycle_events("routine-scheduler", &events, &trusted_signers);
assert_eq!(
decoded.len(),
1,
"an event signed by a retired-but-trusted key must still decode"
);
assert_eq!(decoded[0].signer_public_key, retired.public_key_bytes());
}
#[test]
fn rotation_keeps_both_old_and_new_signer_events_visible() {
let key_a = ApprovalSigner::from_seed(9);
let key_b = ApprovalSigner::from_seed(10);
let (created_a, _, _) = routine_created_payload(
"daily-standup",
"persona-1",
"conv-1",
"call-12",
"hash-12",
1_000,
&key_a,
);
let (created_b, _, _) = routine_created_payload(
"weekly-digest",
"persona-2",
"conv-2",
"call-13",
"hash-13",
2_000,
&key_b,
);
let events = vec![
(1, Event::new(kinds::ROUTINE_CREATED, created_a)),
(2, Event::new(kinds::ROUTINE_CREATED, created_b)),
];
let trusted_signers = vec![key_b.public_key_bytes(), key_a.public_key_bytes()];
let decoded =
decode_routine_lifecycle_events("routine-scheduler", &events, &trusted_signers);
assert_eq!(
decoded.len(),
2,
"both the pre-rotation and post-rotation signer's events must decode"
);
assert_eq!(decoded[0].routine, "daily-standup");
assert_eq!(decoded[1].routine, "weekly-digest");
}
}