use std::collections::{HashMap, HashSet};
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_model::Event;
use polyc_proto::events_decode::try_decode_event_payload;
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::events::v1::{
RoutineFireOutcome, RoutineFireOutcomeEvent, RoutineFiredEvent,
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RoutineLifecycleFact {
pub position: u64,
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>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RoutineSetupFact {
pub position: u64,
pub routine_uid: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RoutineFireFact {
pub position: u64,
pub routine: String,
pub routine_uid: String,
pub occurrence: String,
pub scheduled_at_ms: u64,
pub fired_at_ms: u64,
pub outcome: Option<String>,
pub grant_drift_tools: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RoutineLifecycleFacts {
pub lifecycle: Vec<RoutineLifecycleFact>,
pub setup: Vec<RoutineSetupFact>,
pub fires: Vec<RoutineFireFact>,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum RoutineLifecycleError {
#[error("position {position} appears more than once in one source prefix")]
RepeatedPosition {
position: u64,
},
#[error("the routine-scheduler partition carries an undeclared kind at position {position}")]
UndeclaredKind {
position: u64,
},
}
struct LifecycleFields {
routine: String,
actor_persona: String,
conversation_id: String,
at_ms: u64,
reason: Option<String>,
scope: Option<String>,
channel: Option<String>,
signer_public_key: Vec<u8>,
}
impl From<VerifiedRoutineCreated> for LifecycleFields {
fn from(v: VerifiedRoutineCreated) -> Self {
Self {
routine: v.routine,
actor_persona: v.creator_persona,
conversation_id: v.conversation_id,
at_ms: v.created_at_ms,
reason: None,
scope: None,
channel: None,
signer_public_key: v.signer_public_key,
}
}
}
impl From<VerifiedRoutinePaused> for LifecycleFields {
fn from(v: VerifiedRoutinePaused) -> Self {
Self {
routine: v.routine,
actor_persona: v.actor_persona,
conversation_id: v.conversation_id,
at_ms: v.paused_at_ms,
reason: v.reason,
scope: None,
channel: Some(v.channel),
signer_public_key: v.signer_public_key,
}
}
}
impl From<VerifiedRoutineResumed> for LifecycleFields {
fn from(v: VerifiedRoutineResumed) -> Self {
Self {
routine: v.routine,
actor_persona: v.actor_persona,
conversation_id: v.conversation_id,
at_ms: v.resumed_at_ms,
reason: None,
scope: None,
channel: Some(v.channel),
signer_public_key: v.signer_public_key,
}
}
}
impl From<VerifiedRoutineDeleted> for LifecycleFields {
fn from(v: VerifiedRoutineDeleted) -> Self {
Self {
routine: v.routine,
actor_persona: v.actor_persona,
conversation_id: v.conversation_id,
at_ms: v.deleted_at_ms,
reason: None,
scope: None,
channel: Some(v.channel),
signer_public_key: v.signer_public_key,
}
}
}
impl From<VerifiedRoutineScopeChanged> for LifecycleFields {
fn from(v: VerifiedRoutineScopeChanged) -> Self {
Self {
routine: v.routine,
actor_persona: v.actor_persona,
conversation_id: v.conversation_id,
at_ms: v.changed_at_ms,
reason: None,
scope: Some(v.scope),
channel: None,
signer_public_key: v.signer_public_key,
}
}
}
pub fn fold_routine_lifecycle(
events: &[(u64, Event)],
trusted_signers: &[Vec<u8>],
) -> Result<RoutineLifecycleFacts, RoutineLifecycleError> {
refuse_repeated_positions(events)?;
let outcomes = fold_fire_outcomes(events);
let mut facts = RoutineLifecycleFacts::default();
for (position, event) in events {
let (base, _) = kinds::parse(&event.kind);
if fold_lifecycle_kind(&mut facts, *position, base, &event.payload, trusted_signers) {
continue;
}
match base {
kinds::ROUTINE_SETUP_COMPLETED => {
fold_setup_marker(&mut facts, *position, &event.payload);
}
kinds::ROUTINE_FIRED => {
fold_fired_marker(&mut facts, *position, &event.payload, &outcomes);
}
kinds::ROUTINE_FIRE_OUTCOME => {}
_ => {
return Err(RoutineLifecycleError::UndeclaredKind {
position: *position,
});
}
}
}
Ok(facts)
}
fn fold_lifecycle_kind(
facts: &mut RoutineLifecycleFacts,
position: u64,
base: &str,
payload: &[u8],
trusted_signers: &[Vec<u8>],
) -> bool {
match base {
kinds::ROUTINE_CREATED => fold_lifecycle_event(
facts,
position,
kinds::ROUTINE_CREATED,
payload,
trusted_signers,
verify_routine_created,
|v| v.signer_public_key.as_slice(),
LifecycleFields::from,
),
kinds::ROUTINE_PAUSED => fold_lifecycle_event(
facts,
position,
kinds::ROUTINE_PAUSED,
payload,
trusted_signers,
verify_routine_paused,
|v| v.signer_public_key.as_slice(),
LifecycleFields::from,
),
kinds::ROUTINE_RESUMED => fold_lifecycle_event(
facts,
position,
kinds::ROUTINE_RESUMED,
payload,
trusted_signers,
verify_routine_resumed,
|v| v.signer_public_key.as_slice(),
LifecycleFields::from,
),
kinds::ROUTINE_DELETED => fold_lifecycle_event(
facts,
position,
kinds::ROUTINE_DELETED,
payload,
trusted_signers,
verify_routine_deleted,
|v| v.signer_public_key.as_slice(),
LifecycleFields::from,
),
kinds::ROUTINE_SCOPE_CHANGED => fold_lifecycle_event(
facts,
position,
kinds::ROUTINE_SCOPE_CHANGED,
payload,
trusted_signers,
verify_routine_scope_changed,
|v| v.signer_public_key.as_slice(),
LifecycleFields::from,
),
_ => return false,
}
true
}
#[allow(clippy::too_many_arguments)] fn fold_lifecycle_event<V>(
facts: &mut RoutineLifecycleFacts,
position: u64,
phase: &'static str,
payload: &[u8],
trusted_signers: &[Vec<u8>],
verify: fn(&[u8]) -> Option<V>,
signer_public_key: fn(&V) -> &[u8],
fields: fn(V) -> LifecycleFields,
) {
let Some(verified) = verify(payload) else {
tracing::warn!(
position,
kind = phase,
table = "routine_lifecycle",
"routine_lifecycle: dropping row — payload failed verify_routine_* \
(malformed payload or a signature that doesn't check out)"
);
return;
};
if !trusted_signers
.iter()
.any(|key| key.as_slice() == signer_public_key(&verified))
{
tracing::warn!(
position,
kind = phase,
table = "routine_lifecycle",
signer_public_key = %polyc_crypto::hex::lower(signer_public_key(&verified)),
"routine_lifecycle: dropping row — signer_public_key is not in trusted_signers"
);
return;
}
let fields = fields(verified);
facts.lifecycle.push(RoutineLifecycleFact {
position,
phase: phase.to_owned(),
routine: fields.routine,
actor_persona: fields.actor_persona,
conversation_id: fields.conversation_id,
at_ms: fields.at_ms,
reason: fields.reason,
scope: fields.scope,
channel: fields.channel,
signer_public_key: fields.signer_public_key,
});
}
fn fold_fire_outcomes(events: &[(u64, Event)]) -> HashMap<(String, String), (String, String)> {
let mut outcomes = HashMap::new();
for (_position, event) in events {
let (base, _) = kinds::parse(&event.kind);
if base != kinds::ROUTINE_FIRE_OUTCOME {
continue;
}
let Ok(decoded) = try_decode_event_payload::<RoutineFireOutcomeEvent>(&event.payload)
else {
tracing::warn!(
kind = kinds::ROUTINE_FIRE_OUTCOME,
table = "fires",
"fires: dropping outcome — payload failed to decode"
);
continue;
};
let Some(label) = outcome_str(decoded.outcome.as_known().unwrap_or_default()) else {
continue;
};
let drift =
serde_json::to_string(&decoded.grant_drift_tools).unwrap_or_else(|_| "[]".to_owned());
outcomes.insert(
(decoded.routine, decoded.occurrence),
(label.to_owned(), drift),
);
}
outcomes
}
fn fold_fired_marker(
facts: &mut RoutineLifecycleFacts,
position: u64,
payload: &[u8],
outcomes: &HashMap<(String, String), (String, String)>,
) {
let Ok(decoded) = try_decode_event_payload::<RoutineFiredEvent>(payload) else {
tracing::warn!(
position,
kind = kinds::ROUTINE_FIRED,
table = "fires",
"fires: dropping row — payload failed to decode"
);
return;
};
let matched = outcomes.get(&(decoded.routine.clone(), decoded.occurrence.clone()));
facts.fires.push(RoutineFireFact {
position,
routine: decoded.routine,
routine_uid: decoded.routine_uid,
occurrence: decoded.occurrence,
scheduled_at_ms: decoded.scheduled_at_ms,
fired_at_ms: decoded.fired_at_ms,
outcome: matched.map(|(label, _)| label.clone()),
grant_drift_tools: matched.map(|(_, drift)| drift.clone()),
});
}
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,
}
}
fn fold_setup_marker(facts: &mut RoutineLifecycleFacts, position: u64, payload: &[u8]) {
let routine_uid = serde_json::from_slice::<serde_json::Value>(payload)
.ok()
.and_then(|value| {
value
.get("routine_uid")
.and_then(serde_json::Value::as_str)
.map(str::to_owned)
})
.filter(|uid| !uid.is_empty());
match routine_uid {
Some(routine_uid) => facts.setup.push(RoutineSetupFact {
position,
routine_uid,
}),
None => {
tracing::warn!(
position,
kind = kinds::ROUTINE_SETUP_COMPLETED,
table = "routine_setup",
"routine_setup: dropping row — marker has no non-empty routine_uid"
);
}
}
}
fn refuse_repeated_positions(events: &[(u64, Event)]) -> Result<(), RoutineLifecycleError> {
let mut seen: HashSet<u64> = HashSet::with_capacity(events.len());
for (position, _) in events {
if !seen.insert(*position) {
return Err(RoutineLifecycleError::RepeatedPosition {
position: *position,
});
}
}
Ok(())
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs, clippy::unwrap_used)]
use polyc_crypto::approval::{
ApprovalSigner, routine_created_payload, routine_deleted_payload, routine_paused_payload,
routine_resumed_payload, routine_scope_changed_payload,
};
use super::*;
fn marker(routine_uid: &str) -> Event {
let payload = serde_json::json!({ "routine_uid": routine_uid }).to_string();
Event::new(kinds::ROUTINE_SETUP_COMPLETED, payload.into_bytes())
}
fn created(signer: &ApprovalSigner, routine: &str) -> Vec<u8> {
routine_created_payload(
routine,
"persona-owner",
"conv-1",
"tool-call-1",
"args-hash-1",
1_750_000_000_000,
signer,
)
.0
}
#[test]
fn every_declared_kind_folds_to_its_table() {
let signer = ApprovalSigner::from_seed(1);
let trusted = vec![signer.public_key_bytes()];
let events = vec![
(
1,
Event::new(kinds::ROUTINE_CREATED, created(&signer, "r-1")),
),
(2, marker("uid-1")),
(
3,
Event::new(
kinds::ROUTINE_PAUSED,
routine_paused_payload(
"r-1",
"persona-owner",
"conv-1",
"tool-call-2",
"args-hash-2",
"chat",
1_750_000_000_100,
Some("maintenance"),
&signer,
)
.0,
),
),
(
4,
Event::new(
kinds::ROUTINE_RESUMED,
routine_resumed_payload(
"r-1",
"persona-owner",
"conv-1",
"tool-call-3",
"args-hash-3",
"rpc",
1_750_000_000_200,
&signer,
)
.0,
),
),
(
5,
Event::new(
kinds::ROUTINE_DELETED,
routine_deleted_payload(
"r-1",
"persona-owner",
"conv-1",
"tool-call-4",
"args-hash-4",
"chat",
1_750_000_000_300,
&signer,
)
.0,
),
),
(
6,
Event::new(
kinds::ROUTINE_SCOPE_CHANGED,
routine_scope_changed_payload(
"r-1",
"persona-owner",
"conv-1",
"tool-call-5",
"args-hash-5",
"public",
1_750_000_000_400,
&signer,
)
.0,
),
),
];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert_eq!(facts.lifecycle.len(), 5);
assert_eq!(facts.setup.len(), 1);
assert_eq!(facts.lifecycle[0].phase, kinds::ROUTINE_CREATED);
assert_eq!(facts.lifecycle[0].routine, "r-1");
assert_eq!(facts.lifecycle[0].actor_persona, "persona-owner");
assert_eq!(facts.lifecycle[0].at_ms, 1_750_000_000_000);
assert_eq!(facts.lifecycle[1].phase, kinds::ROUTINE_PAUSED);
assert_eq!(facts.lifecycle[1].reason.as_deref(), Some("maintenance"));
assert_eq!(facts.lifecycle[1].channel.as_deref(), Some("chat"));
assert!(facts.lifecycle[1].scope.is_none());
assert_eq!(facts.lifecycle[2].channel.as_deref(), Some("rpc"));
assert_eq!(facts.lifecycle[3].phase, kinds::ROUTINE_DELETED);
assert_eq!(facts.lifecycle[4].scope.as_deref(), Some("public"));
assert_eq!(facts.setup[0].routine_uid, "uid-1");
assert_eq!(facts.setup[0].position, 2);
}
#[test]
fn a_marker_before_the_created_event_still_folds() {
let signer = ApprovalSigner::from_seed(1);
let trusted = vec![signer.public_key_bytes()];
let events = vec![
(1, marker("uid-1")),
(
2,
Event::new(kinds::ROUTINE_CREATED, created(&signer, "r-1")),
),
];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert_eq!(facts.setup.len(), 1);
assert_eq!(facts.setup[0].position, 1);
assert_eq!(facts.lifecycle.len(), 1);
assert_eq!(facts.lifecycle[0].position, 2);
}
#[test]
fn an_untrusted_signer_produces_no_row() {
let signer = ApprovalSigner::from_seed(1);
let foreign = ApprovalSigner::from_seed(9);
let trusted = vec![signer.public_key_bytes()];
let events = vec![(
1,
Event::new(kinds::ROUTINE_CREATED, created(&foreign, "r-evil")),
)];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert!(facts.lifecycle.is_empty());
assert!(facts.setup.is_empty());
}
#[test]
fn a_forged_record_produces_no_row() {
let signer = ApprovalSigner::from_seed(1);
let trusted = vec![signer.public_key_bytes()];
let genuine = created(&signer, "r-1");
let mut value: serde_json::Value = serde_json::from_slice(&genuine).unwrap();
value["routine"] = serde_json::json!("r-forged");
let events = vec![(
1,
Event::new(kinds::ROUTINE_CREATED, value.to_string().into_bytes()),
)];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert!(facts.lifecycle.is_empty(), "a forged record yields no row");
}
#[test]
fn a_fired_marker_joins_its_outcome() {
use buffa::Message as _;
let trusted: Vec<Vec<u8>> = Vec::new();
let marker = RoutineFiredEvent {
routine: "daily-standup".to_owned(),
occurrence: "daily-standup-1".to_owned(),
scheduled_at_ms: 1_750_000_000_000,
fired_at_ms: 1_750_000_000_500,
routine_uid: "uid-1".to_owned(),
..Default::default()
};
let outcome = RoutineFireOutcomeEvent {
routine: marker.routine.clone(),
occurrence: marker.occurrence.clone(),
outcome: RoutineFireOutcome::Ok.into(),
fired_at_ms: marker.fired_at_ms,
grant_drift_tools: vec!["fs_write".to_owned()],
..Default::default()
};
let events = vec![
(1, Event::new(kinds::ROUTINE_FIRED, marker.encode_to_vec())),
(
2,
Event::new(kinds::ROUTINE_FIRE_OUTCOME, outcome.encode_to_vec()),
),
];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert_eq!(facts.fires.len(), 1);
assert_eq!(facts.fires[0].routine, "daily-standup");
assert_eq!(facts.fires[0].occurrence, "daily-standup-1");
assert_eq!(facts.fires[0].routine_uid, "uid-1");
assert_eq!(facts.fires[0].outcome.as_deref(), Some("ok"));
assert_eq!(
facts.fires[0].grant_drift_tools.as_deref(),
Some(r#"["fs_write"]"#)
);
assert!(facts.lifecycle.is_empty());
assert!(facts.setup.is_empty());
}
#[test]
fn a_fired_marker_without_outcome_still_folds() {
use buffa::Message as _;
let trusted: Vec<Vec<u8>> = Vec::new();
let marker = RoutineFiredEvent {
routine: "daily-standup".to_owned(),
occurrence: "daily-standup-1".to_owned(),
scheduled_at_ms: 1,
fired_at_ms: 2,
routine_uid: "uid-1".to_owned(),
..Default::default()
};
let events = vec![(1, Event::new(kinds::ROUTINE_FIRED, marker.encode_to_vec()))];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert_eq!(facts.fires.len(), 1);
assert!(facts.fires[0].outcome.is_none());
assert!(facts.fires[0].grant_drift_tools.is_none());
}
#[test]
fn a_malformed_fired_marker_produces_no_row() {
let trusted: Vec<Vec<u8>> = Vec::new();
let events = vec![(1, Event::new(kinds::ROUTINE_FIRED, b"{}".to_vec()))];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert!(facts.fires.is_empty());
}
#[test]
fn an_undeclared_kind_refuses_the_generation() {
let signer = ApprovalSigner::from_seed(1);
let trusted = vec![signer.public_key_bytes()];
let events = vec![
(
1,
Event::new(kinds::ROUTINE_CREATED, created(&signer, "r-1")),
),
(2, Event::new("turn_start", Vec::new())),
];
let result = fold_routine_lifecycle(&events, &trusted);
assert_eq!(
result,
Err(RoutineLifecycleError::UndeclaredKind { position: 2 })
);
}
#[test]
fn a_tagged_lifecycle_kind_folds_as_its_base() {
let signer = ApprovalSigner::from_seed(1);
let trusted = vec![signer.public_key_bytes()];
let tagged = format!(
"{}:9b8d6e62-2a44-4ef2-8ec1-934d598b2410",
kinds::ROUTINE_CREATED
);
let events = vec![(1, Event::new(tagged, created(&signer, "r-1")))];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert_eq!(facts.lifecycle.len(), 1);
assert_eq!(facts.lifecycle[0].phase, kinds::ROUTINE_CREATED);
}
#[test]
fn a_uid_less_marker_produces_no_row() {
let trusted: Vec<Vec<u8>> = Vec::new();
let events = vec![
(1, marker("")),
(
2,
Event::new(kinds::ROUTINE_SETUP_COMPLETED, b"{}".to_vec()),
),
(
3,
Event::new(kinds::ROUTINE_SETUP_COMPLETED, vec![0xFF, 0xFE]),
),
];
let facts = fold_routine_lifecycle(&events, &trusted).unwrap();
assert!(facts.setup.is_empty());
}
#[test]
fn a_repeated_position_refuses_the_generation() {
let trusted: Vec<Vec<u8>> = Vec::new();
let events = vec![(1, marker("uid-1")), (1, marker("uid-2"))];
let result = fold_routine_lifecycle(&events, &trusted);
assert_eq!(
result,
Err(RoutineLifecycleError::RepeatedPosition { position: 1 })
);
}
}