use std::sync::Arc;
use arrow::array::{ArrayRef, BinaryBuilder, BooleanBuilder, 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;
const PHASE_REQUEST: &str = "request";
const PHASE_RESPONSE: &str = "response";
const SIGNATURE_VERIFIED: &str = "verified";
const SIGNATURE_INVALID: &str = "invalid";
const SIGNATURE_LEGACY_UNVERIFIABLE: &str = "legacy_unverifiable";
#[derive(Debug, Clone)]
pub(crate) struct ApprovalRow {
pub partition: String,
pub position: u64,
pub turn_id: Option<String>,
pub phase: String,
pub request_id: String,
pub tool_name: Option<String>,
pub args_json: Option<String>,
pub request_reason: Option<String>,
pub request_sandbox_mode: Option<String>,
pub approved: Option<bool>,
pub response_reason: Option<String>,
pub signature_status: Option<String>,
pub signer_public_key: Option<Vec<u8>>,
pub modified_args_json: Option<String>,
pub approved_for_session: Option<bool>,
pub caller: Option<String>,
pub approver: Option<String>,
pub response_sandbox_mode: Option<String>,
pub injected_context: Option<String>,
pub routine_grant: Option<bool>,
pub tool_descriptor_hash: Option<String>,
pub grant_scope: 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("phase", DataType::Utf8, false),
Field::new("request_id", DataType::Utf8, false),
Field::new("tool_name", DataType::Utf8, true),
Field::new("args_json", DataType::Utf8, true),
Field::new("request_reason", DataType::Utf8, true),
Field::new("request_sandbox_mode", DataType::Utf8, true),
Field::new("approved", DataType::Boolean, true),
Field::new("response_reason", DataType::Utf8, true),
Field::new("signature_status", DataType::Utf8, true),
Field::new("signer_public_key", DataType::Binary, true),
Field::new("modified_args_json", DataType::Utf8, true),
Field::new("approved_for_session", DataType::Boolean, true),
Field::new("caller", DataType::Utf8, true),
Field::new("approver", DataType::Utf8, true),
Field::new("response_sandbox_mode", DataType::Utf8, true),
Field::new("injected_context", DataType::Utf8, true),
Field::new("routine_grant", DataType::Boolean, true),
Field::new("tool_descriptor_hash", DataType::Utf8, true),
Field::new("grant_scope", DataType::Utf8, true),
]))
}
#[allow(clippy::too_many_lines, clippy::similar_names)]
pub(crate) fn decode_approvals_batch(rows: &[ApprovalRow]) -> 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 phase_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut request_id_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut tool_name_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut args_json_b = StringBuilder::with_capacity(rows.len(), rows.len() * 32);
let mut request_reason_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut request_sandbox_mode_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut approved_b = BooleanBuilder::with_capacity(rows.len());
let mut response_reason_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut signature_status_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut signer_public_key_b = BinaryBuilder::with_capacity(rows.len(), rows.len() * 32);
let mut modified_args_json_b = StringBuilder::with_capacity(rows.len(), rows.len() * 32);
let mut approved_for_session_b = BooleanBuilder::with_capacity(rows.len());
let mut caller_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut approver_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut response_sandbox_mode_b = StringBuilder::with_capacity(rows.len(), rows.len() * 8);
let mut injected_context_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut routine_grant_b = BooleanBuilder::with_capacity(rows.len());
let mut tool_descriptor_hash_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
let mut grant_scope_b = StringBuilder::with_capacity(rows.len(), rows.len() * 16);
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);
request_id_b.append_value(&row.request_id);
match &row.tool_name {
Some(v) => tool_name_b.append_value(v),
None => tool_name_b.append_null(),
}
match &row.args_json {
Some(v) => args_json_b.append_value(v),
None => args_json_b.append_null(),
}
match &row.request_reason {
Some(v) => request_reason_b.append_value(v),
None => request_reason_b.append_null(),
}
match &row.request_sandbox_mode {
Some(v) => request_sandbox_mode_b.append_value(v),
None => request_sandbox_mode_b.append_null(),
}
match row.approved {
Some(v) => approved_b.append_value(v),
None => approved_b.append_null(),
}
match &row.response_reason {
Some(v) => response_reason_b.append_value(v),
None => response_reason_b.append_null(),
}
match &row.signature_status {
Some(v) => signature_status_b.append_value(v),
None => signature_status_b.append_null(),
}
match &row.signer_public_key {
Some(v) => signer_public_key_b.append_value(v),
None => signer_public_key_b.append_null(),
}
match &row.modified_args_json {
Some(v) => modified_args_json_b.append_value(v),
None => modified_args_json_b.append_null(),
}
match row.approved_for_session {
Some(v) => approved_for_session_b.append_value(v),
None => approved_for_session_b.append_null(),
}
match &row.caller {
Some(v) => caller_b.append_value(v),
None => caller_b.append_null(),
}
match &row.approver {
Some(v) => approver_b.append_value(v),
None => approver_b.append_null(),
}
match &row.response_sandbox_mode {
Some(v) => response_sandbox_mode_b.append_value(v),
None => response_sandbox_mode_b.append_null(),
}
match &row.injected_context {
Some(v) => injected_context_b.append_value(v),
None => injected_context_b.append_null(),
}
match row.routine_grant {
Some(v) => routine_grant_b.append_value(v),
None => routine_grant_b.append_null(),
}
match &row.tool_descriptor_hash {
Some(v) => tool_descriptor_hash_b.append_value(v),
None => tool_descriptor_hash_b.append_null(),
}
match &row.grant_scope {
Some(v) => grant_scope_b.append_value(v),
None => grant_scope_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(phase_b.finish()),
Arc::new(request_id_b.finish()),
Arc::new(tool_name_b.finish()),
Arc::new(args_json_b.finish()),
Arc::new(request_reason_b.finish()),
Arc::new(request_sandbox_mode_b.finish()),
Arc::new(approved_b.finish()),
Arc::new(response_reason_b.finish()),
Arc::new(signature_status_b.finish()),
Arc::new(signer_public_key_b.finish()),
Arc::new(modified_args_json_b.finish()),
Arc::new(approved_for_session_b.finish()),
Arc::new(caller_b.finish()),
Arc::new(approver_b.finish()),
Arc::new(response_sandbox_mode_b.finish()),
Arc::new(injected_context_b.finish()),
Arc::new(routine_grant_b.finish()),
Arc::new(tool_descriptor_hash_b.finish()),
Arc::new(grant_scope_b.finish()),
];
RecordBatch::try_new(schema(), columns)
}
const fn signature_status_str(status: polyc_facts::ApprovalSignatureStatus) -> &'static str {
match status {
polyc_facts::ApprovalSignatureStatus::Verified => SIGNATURE_VERIFIED,
polyc_facts::ApprovalSignatureStatus::Invalid => SIGNATURE_INVALID,
polyc_facts::ApprovalSignatureStatus::LegacyUnverifiable => SIGNATURE_LEGACY_UNVERIFIABLE,
}
}
#[must_use]
pub(crate) fn decode_approvals_events(
partition: &str,
events: &[(u64, Event)],
trusted_signers: &[Vec<u8>],
) -> Vec<ApprovalRow> {
events
.iter()
.filter_map(|(position, event)| {
let (_base, turn_id) = kinds::parse(&event.kind);
let turn_id = turn_id.map(|id| id.to_string());
match polyc_facts::fold_approval_event(event, trusted_signers)? {
polyc_facts::ApprovalFact::Request(req) => Some(ApprovalRow {
partition: partition.to_string(),
position: *position,
turn_id,
phase: PHASE_REQUEST.to_string(),
request_id: req.request_id,
tool_name: Some(req.tool_name),
args_json: Some(req.args_json),
request_reason: Some(req.reason),
request_sandbox_mode: Some(req.sandbox_mode),
approved: None,
response_reason: None,
signature_status: None,
signer_public_key: None,
modified_args_json: None,
approved_for_session: None,
caller: None,
approver: None,
response_sandbox_mode: None,
injected_context: None,
routine_grant: None,
tool_descriptor_hash: None,
grant_scope: None,
}),
polyc_facts::ApprovalFact::Response(resp) => Some(ApprovalRow {
partition: partition.to_string(),
position: *position,
turn_id,
phase: PHASE_RESPONSE.to_string(),
request_id: resp.request_id,
tool_name: resp.tool_name,
args_json: None,
request_reason: None,
request_sandbox_mode: None,
approved: Some(resp.approved),
response_reason: resp.response_reason,
signature_status: Some(signature_status_str(resp.signature_status).to_string()),
signer_public_key: resp.signer_public_key,
modified_args_json: resp.modified_args_json,
approved_for_session: resp.approved_for_session,
caller: resp.caller,
approver: resp.approver,
response_sandbox_mode: resp.sandbox_mode,
injected_context: resp.injected_context,
routine_grant: resp.routine_grant,
tool_descriptor_hash: resp.tool_descriptor_hash,
grant_scope: resp.grant_scope,
}),
}
})
.collect()
}
#[cfg(test)]
mod tests {
use arrow::array::Array as _;
use polyc_crypto::approval::{ApprovalSigner, request_payload, response_payload};
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",
"phase",
"request_id",
"tool_name",
"args_json",
"request_reason",
"request_sandbox_mode",
"approved",
"response_reason",
"signature_status",
"signer_public_key",
"modified_args_json",
"approved_for_session",
"caller",
"approver",
"response_sandbox_mode",
"injected_context",
"routine_grant",
"tool_descriptor_hash",
"grant_scope",
]
);
let expect = [
("partition", DataType::Utf8, false),
("position", DataType::UInt64, false),
("turn_id", DataType::Utf8, true),
("phase", DataType::Utf8, false),
("request_id", DataType::Utf8, false),
("tool_name", DataType::Utf8, true),
("args_json", DataType::Utf8, true),
("request_reason", DataType::Utf8, true),
("request_sandbox_mode", DataType::Utf8, true),
("approved", DataType::Boolean, true),
("response_reason", DataType::Utf8, true),
("signature_status", DataType::Utf8, true),
("signer_public_key", DataType::Binary, true),
("modified_args_json", DataType::Utf8, true),
("approved_for_session", DataType::Boolean, true),
("caller", DataType::Utf8, true),
("approver", DataType::Utf8, true),
("response_sandbox_mode", DataType::Utf8, true),
("injected_context", DataType::Utf8, true),
("routine_grant", DataType::Boolean, true),
("tool_descriptor_hash", DataType::Utf8, true),
("grant_scope", 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_real_signed_request_and_response_pair() {
let turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_4444);
let signer = ApprovalSigner::from_seed(1);
let request_bytes = request_payload(
"call-1",
"rm",
r#"{"path":"/tmp"}"#,
"default",
"",
&[],
"",
"",
"",
&[],
false,
);
let response_bytes = response_payload(
"call-1",
"rm",
r#"{"path":"/tmp"}"#,
"",
true,
false,
&[],
"caller-1",
"",
"default",
"ok",
"",
"conv-rt",
"nonce-1",
&turn.to_string(),
&signer,
)
.0;
let events = vec![
(
1,
Event::new(kinds::tagged(kinds::APPROVAL_REQUEST, &turn), request_bytes),
),
(
2,
Event::new(
kinds::tagged(kinds::APPROVAL_RESPONSE, &turn),
response_bytes,
),
),
];
let trusted_signers = vec![signer.public_key_bytes()];
let decoded = decode_approvals_events("conv-rt", &events, &trusted_signers);
assert_eq!(decoded.len(), 2);
let req = &decoded[0];
assert_eq!(req.phase, "request");
assert_eq!(req.request_id, "call-1");
assert_eq!(req.turn_id, Some(turn.to_string()));
assert_eq!(req.tool_name.as_deref(), Some("rm"));
assert_eq!(req.args_json.as_deref(), Some(r#"{"path":"/tmp"}"#));
assert_eq!(req.approved, None);
assert_eq!(req.signature_status, None);
let resp = &decoded[1];
assert_eq!(resp.phase, "response");
assert_eq!(resp.request_id, "call-1");
assert_eq!(resp.approved, Some(true));
assert_eq!(resp.response_reason.as_deref(), Some("ok"));
assert_eq!(resp.signature_status.as_deref(), Some("verified"));
assert_eq!(resp.signer_public_key, Some(signer.public_key_bytes()));
assert_eq!(resp.caller.as_deref(), Some("caller-1"));
assert_eq!(
resp.tool_name.as_deref(),
Some("rm"),
"a v2 response's own signed tool_name binding populates the shared column"
);
assert_eq!(resp.routine_grant, Some(false));
assert_eq!(resp.tool_descriptor_hash.as_deref(), Some(""));
assert_eq!(resp.grant_scope.as_deref(), Some(""));
let batch = decode_approvals_batch(&decoded).expect("batch build");
assert_eq!(batch.num_rows(), 2);
assert_eq!(batch.schema(), schema());
let phase = batch
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(phase.value(0), "request");
assert_eq!(phase.value(1), "response");
let tool_name = batch
.column(5)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert!(tool_name.is_valid(0));
assert!(
tool_name.is_valid(1),
"a v2 response's own signed tool_name binding populates the shared column"
);
assert_eq!(tool_name.value(1), "rm");
}
#[test]
fn signature_status_trusted_signer_verifies_and_row_is_present() {
let signer = ApprovalSigner::from_seed(2);
let bytes = response_payload(
"call-2",
"rm",
"{}",
"",
true,
false,
&[],
"caller-2",
"",
"",
"ok",
"",
"conv-a",
"n1",
"",
&signer,
)
.0;
let events = vec![(1, Event::new(kinds::APPROVAL_RESPONSE.to_owned(), bytes))];
let trusted_signers = vec![signer.public_key_bytes()];
let decoded = decode_approvals_events("conv-a", &events, &trusted_signers);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].signature_status.as_deref(), Some("verified"));
}
#[test]
fn signature_status_untrusted_signer_is_invalid_but_row_stays_present() {
let trusted = ApprovalSigner::from_seed(3);
let untrusted = ApprovalSigner::from_seed(4);
let bytes = response_payload(
"call-3",
"rm",
"{}",
"",
true,
false,
&[],
"caller-3",
"",
"",
"ok",
"",
"conv-b",
"n1",
"",
&untrusted,
)
.0;
let events = vec![(1, Event::new(kinds::APPROVAL_RESPONSE.to_owned(), bytes))];
let trusted_signers = vec![trusted.public_key_bytes()];
let decoded = decode_approvals_events("conv-b", &events, &trusted_signers);
assert_eq!(
decoded.len(),
1,
"an untrusted-signer response must still surface as a row, unlike payments"
);
assert_eq!(decoded[0].signature_status.as_deref(), Some("invalid"));
assert_eq!(
decoded[0].approved,
Some(true),
"claimed fields still shown"
);
}
#[test]
fn signature_status_legacy_shape_is_unverifiable_not_invalid() {
let legacy = br#"{"request_id":"call-4","approved":true,"reason":"ok"}"#.to_vec();
let events = vec![(1, Event::new(kinds::APPROVAL_RESPONSE.to_owned(), legacy))];
let decoded = decode_approvals_events("conv-c", &events, &[]);
assert_eq!(decoded.len(), 1);
assert_eq!(
decoded[0].signature_status.as_deref(),
Some("legacy_unverifiable")
);
assert_eq!(decoded[0].response_reason, None);
}
#[test]
fn structurally_malformed_response_drops_the_row() {
let events = vec![(
1,
Event::new(kinds::APPROVAL_RESPONSE.to_owned(), vec![0xFF, 0xFE, 0xFD]),
)];
let decoded = decode_approvals_events("conv-d", &events, &[]);
assert_eq!(decoded.len(), 0);
}
#[test]
fn unrelated_kind_is_not_decoded_as_an_approval() {
let events = vec![(1, Event::new(kinds::USAGE.to_owned(), Vec::new()))];
let decoded = decode_approvals_events("conv-e", &events, &[]);
assert_eq!(decoded.len(), 0);
}
#[test]
fn decode_reads_a_routine_grant_response() {
let signer = ApprovalSigner::from_seed(42);
let (payload, ..) = polyc_crypto::approval::routine_grant_payload(
"call-9",
"fs_write",
"{}",
"",
true,
"owner-1",
"",
"default",
"",
&[],
"conv-fire-9",
"nonce-9",
"",
"hash-xyz",
"blanket_below_high",
&signer,
);
let events = vec![(1, Event::new(kinds::APPROVAL_RESPONSE.to_owned(), payload))];
let trusted_signers = vec![signer.public_key_bytes()];
let decoded = decode_approvals_events("conv-fire-9", &events, &trusted_signers);
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].tool_name.as_deref(), Some("fs_write"));
assert_eq!(decoded[0].routine_grant, Some(true));
assert_eq!(decoded[0].tool_descriptor_hash.as_deref(), Some("hash-xyz"));
assert_eq!(
decoded[0].grant_scope.as_deref(),
Some("blanket_below_high")
);
let batch = decode_approvals_batch(&decoded).expect("batch build");
let routine_grant = batch
.column(19)
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.unwrap();
assert!(routine_grant.value(0));
let grant_scope = batch
.column(21)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(grant_scope.value(0), "blanket_below_high");
}
}