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_eventlog::Event;
use polyc_proto::kinds;
#[cfg(test)]
const KNOWN_TRANSITIONS: [&str; 3] = ["linked", "renewed", "revoked"];
#[cfg(test)]
fn is_known_transition(transition: &str) -> bool {
KNOWN_TRANSITIONS.contains(&transition)
}
#[derive(Debug, Clone)]
pub(crate) struct WalletLinkLifecycleRow {
pub partition: String,
pub position: u64,
pub turn_id: Option<String>,
pub transition: String,
pub subject: String,
pub wallet_address: String,
pub currency: String,
pub chain_id: String,
pub limit_base_units: String,
pub limit_human: String,
pub period_secs: String,
pub expiry_unix: String,
pub recipients: String,
pub conversation_id: String,
pub timestamp: 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("transition", DataType::Utf8, false),
Field::new("subject", DataType::Utf8, false),
Field::new("wallet_address", DataType::Utf8, false),
Field::new("currency", DataType::Utf8, false),
Field::new("chain_id", DataType::Utf8, false),
Field::new("limit_base_units", DataType::Utf8, false),
Field::new("limit_human", DataType::Utf8, false),
Field::new("period_secs", DataType::Utf8, false),
Field::new("expiry_unix", DataType::Utf8, false),
Field::new("recipients", DataType::Utf8, false),
Field::new("conversation_id", DataType::Utf8, false),
Field::new("timestamp", DataType::Utf8, false),
Field::new("signer_public_key", DataType::Binary, false),
]))
}
pub(crate) fn decode_wallet_link_lifecycle_batch(
rows: &[WalletLinkLifecycleRow],
) -> 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 transition_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut subject_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut wallet_address_b = StringBuilder::with_capacity(rows.len(), rows.len() * 42);
let mut currency_b = StringBuilder::with_capacity(rows.len(), rows.len() * 42);
let mut chain_id_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut limit_base_units_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut limit_human_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut period_secs_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut expiry_unix_b = StringBuilder::with_capacity(rows.len(), rows.len() * 12);
let mut recipients_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut conversation_id_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut timestamp_b = StringBuilder::with_capacity(rows.len(), rows.len() * 12);
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(),
}
transition_b.append_value(&row.transition);
subject_b.append_value(&row.subject);
wallet_address_b.append_value(&row.wallet_address);
currency_b.append_value(&row.currency);
chain_id_b.append_value(&row.chain_id);
limit_base_units_b.append_value(&row.limit_base_units);
limit_human_b.append_value(&row.limit_human);
period_secs_b.append_value(&row.period_secs);
expiry_unix_b.append_value(&row.expiry_unix);
recipients_b.append_value(&row.recipients);
conversation_id_b.append_value(&row.conversation_id);
timestamp_b.append_value(&row.timestamp);
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(transition_b.finish()),
Arc::new(subject_b.finish()),
Arc::new(wallet_address_b.finish()),
Arc::new(currency_b.finish()),
Arc::new(chain_id_b.finish()),
Arc::new(limit_base_units_b.finish()),
Arc::new(limit_human_b.finish()),
Arc::new(period_secs_b.finish()),
Arc::new(expiry_unix_b.finish()),
Arc::new(recipients_b.finish()),
Arc::new(conversation_id_b.finish()),
Arc::new(timestamp_b.finish()),
Arc::new(signer_public_key_b.finish()),
];
RecordBatch::try_new(schema(), columns)
}
#[must_use]
pub(crate) fn decode_wallet_link_lifecycle_events(
partition: &str,
events: &[(u64, Event)],
trusted_signers: &[Vec<u8>],
) -> Vec<WalletLinkLifecycleRow> {
events
.iter()
.filter_map(|(position, event)| {
let (base, turn_id) = kinds::parse(&event.kind);
if base != kinds::WALLET_LINK_LIFECYCLE {
return None;
}
let verified = polyc_facts::verified_wallet_link_lifecycle_events(
std::slice::from_ref(event),
trusted_signers,
)
.next()?;
Some(WalletLinkLifecycleRow {
partition: partition.to_string(),
position: *position,
turn_id: turn_id.map(|id| id.to_string()),
transition: verified.transition,
subject: verified.subject,
wallet_address: verified.wallet_address,
currency: verified.currency,
chain_id: verified.chain_id,
limit_base_units: verified.limit_base_units,
limit_human: verified.limit_human,
period_secs: verified.period_secs,
expiry_unix: verified.expiry_unix,
recipients: verified.recipients,
conversation_id: verified.conversation_id,
timestamp: verified.timestamp,
signer_public_key: verified.signer_public_key,
})
})
.collect()
}
#[cfg(test)]
mod tests {
use arrow::array::Array as _;
use polyc_crypto::approval::{
ApprovalSigner, WalletLinkLifecyclePayload, wallet_link_lifecycle_payload,
};
use super::*;
fn signed_lifecycle_event(
signer: &ApprovalSigner,
transition: &str,
conversation_id: &str,
) -> Vec<u8> {
let (payload, sig, pk) = wallet_link_lifecycle_payload(
&WalletLinkLifecyclePayload {
kind: kinds::WALLET_LINK_LIFECYCLE,
transition,
subject: "persona-1",
wallet_address: "0x1111111111111111111111111111111111111111",
currency: "0x2222222222222222222222222222222222222222",
chain_id: "42431",
limit_base_units: "5000000",
limit_human: "5",
period_secs: "86400",
expiry_unix: "1780086400",
recipients: "",
conversation_id,
timestamp: "1780000000",
},
signer,
);
let _ = (sig, pk);
payload
}
#[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",
"transition",
"subject",
"wallet_address",
"currency",
"chain_id",
"limit_base_units",
"limit_human",
"period_secs",
"expiry_unix",
"recipients",
"conversation_id",
"timestamp",
"signer_public_key",
]
);
}
#[test]
fn decode_wallet_link_lifecycle_batch_round_trips() {
let rows = vec![WalletLinkLifecycleRow {
partition: "conv-a".to_string(),
position: 5,
turn_id: None,
transition: "linked".to_string(),
subject: "persona-1".to_string(),
wallet_address: "0x1111111111111111111111111111111111111111".to_string(),
currency: "0x2222222222222222222222222222222222222222".to_string(),
chain_id: "42431".to_string(),
limit_base_units: "5000000".to_string(),
limit_human: "5".to_string(),
period_secs: "86400".to_string(),
expiry_unix: "1780086400".to_string(),
recipients: String::new(),
conversation_id: "conv-a".to_string(),
timestamp: "1780000000".to_string(),
signer_public_key: vec![1, 2, 3, 4],
}];
let batch = decode_wallet_link_lifecycle_batch(&rows).expect("batch build");
assert_eq!(batch.num_rows(), 1);
assert_eq!(batch.schema(), schema());
let transition = batch
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(transition.value(0), "linked");
}
#[test]
fn decode_wallet_link_lifecycle_events_verifies_and_decodes() {
let signer = ApprovalSigner::from_seed(1);
let bytes = signed_lifecycle_event(&signer, "linked", "conv-real");
let events = vec![
(1, Event::new(kinds::TURN_START, Vec::new())),
(2, Event::new(kinds::WALLET_LINK_LIFECYCLE, bytes)),
];
let trusted_signers = vec![signer.public_key_bytes()];
let decoded = decode_wallet_link_lifecycle_events("conv-real", &events, &trusted_signers);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].partition, "conv-real");
assert_eq!(decoded[0].position, 2);
assert_eq!(decoded[0].turn_id, None);
assert_eq!(decoded[0].transition, "linked");
assert_eq!(decoded[0].conversation_id, "conv-real");
assert_eq!(decoded[0].signer_public_key, signer.public_key_bytes());
}
#[test]
fn decode_wallet_link_lifecycle_events_bare_kind_has_no_turn_id() {
let signer = ApprovalSigner::from_seed(5);
let bytes = signed_lifecycle_event(&signer, "revoked", "conv-bare");
let events = vec![(7, Event::new(kinds::WALLET_LINK_LIFECYCLE, bytes))];
let trusted_signers = vec![signer.public_key_bytes()];
let decoded = decode_wallet_link_lifecycle_events("conv-bare", &events, &trusted_signers);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].turn_id, None);
}
#[test]
fn decode_wallet_link_lifecycle_events_keeps_an_unknown_transition_tag() {
let signer = ApprovalSigner::from_seed(2);
let bytes = signed_lifecycle_event(&signer, "some_future_transition", "conv-unknown");
let events = vec![(1, Event::new(kinds::WALLET_LINK_LIFECYCLE, bytes))];
let trusted_signers = vec![signer.public_key_bytes()];
let decoded =
decode_wallet_link_lifecycle_events("conv-unknown", &events, &trusted_signers);
assert_eq!(
decoded.len(),
1,
"an unknown transition tag must not be dropped"
);
assert_eq!(decoded[0].transition, "some_future_transition");
assert!(!is_known_transition(&decoded[0].transition));
}
#[test]
fn decode_wallet_link_lifecycle_events_drops_an_event_from_an_untrusted_signer() {
let trusted = ApprovalSigner::from_seed(3);
let untrusted = ApprovalSigner::from_seed(4);
let forged = signed_lifecycle_event(&untrusted, "linked", "conv-forged");
let events = vec![(1, Event::new(kinds::WALLET_LINK_LIFECYCLE, forged))];
let trusted_signers = vec![trusted.public_key_bytes()];
let decoded = decode_wallet_link_lifecycle_events("conv-forged", &events, &trusted_signers);
assert_eq!(
decoded.len(),
0,
"an event signed by a key outside trusted_signers must never appear as a row"
);
}
#[test]
fn decode_wallet_link_lifecycle_events_drops_a_malformed_payload() {
let events = vec![(
1,
Event::new(kinds::WALLET_LINK_LIFECYCLE, vec![0xFF, 0xFE, 0xFD]),
)];
let decoded = decode_wallet_link_lifecycle_events("conv-corrupt", &events, &[]);
assert_eq!(decoded.len(), 0);
}
#[test]
fn unrelated_kind_is_not_decoded_as_a_wallet_link_lifecycle_event() {
let signer = ApprovalSigner::from_seed(6);
let bytes = signed_lifecycle_event(&signer, "linked", "conv-unrelated");
let events = vec![(1, Event::new(kinds::USAGE, bytes))];
let trusted_signers = vec![signer.public_key_bytes()];
let decoded =
decode_wallet_link_lifecycle_events("conv-unrelated", &events, &trusted_signers);
assert_eq!(
decoded.len(),
0,
"a usage-kind event must never decode as a wallet-link lifecycle event"
);
}
#[test]
fn known_transitions_cover_every_recording_seam_tag() {
for tag in ["linked", "renewed", "revoked"] {
assert!(
is_known_transition(tag),
"{tag} must be a known transition tag"
);
}
assert!(!is_known_transition("not_a_real_tag"));
}
}