use crate::communication_replay::{CommunicationConsumptionArtifact, CommunicationReplayMode};
use crate::determinism::EffectDeterminismTier;
use crate::effect::{CorruptionType, EffectTraceEntry};
use crate::session::{
AuthorityArtifact, AuthorityAuditEvent, AuthorityAuditRecord, AuthorityWitnessId,
FragmentOwnerId, OwnershipTerminalReason, SessionId,
};
use crate::trace::normalize_trace;
use crate::transfer_semantics::{DelegationAuditRecord, DelegationReceipt, DelegationStatus};
use crate::verification::Hash;
use crate::vm::{ObsEvent, SessionTerminalReason};
use serde::{de::DeserializeOwned, Deserialize, Serialize};
use serde_json::Value as JsonValue;
pub const SERIALIZATION_SCHEMA_VERSION: &str = "vm.serialization.v1";
fn default_serialization_schema_version() -> String {
SERIALIZATION_SCHEMA_VERSION.to_string()
}
fn normalize_serialization_schema_version(raw: &str) -> String {
if raw == "1" {
SERIALIZATION_SCHEMA_VERSION.to_string()
} else {
raw.to_string()
}
}
pub fn binary_encode<T: Serialize + ?Sized>(value: &T) -> Result<Vec<u8>, bincode::Error> {
bincode::serialize(value)
}
pub fn binary_decode<T: DeserializeOwned>(bytes: &[u8]) -> Result<T, bincode::Error> {
bincode::deserialize(bytes)
}
#[must_use]
pub fn binary_size<T: Serialize + ?Sized>(value: &T) -> usize {
bincode::serialized_size(value)
.ok()
.and_then(|bytes| usize::try_from(bytes).ok())
.unwrap_or(0)
}
fn deserialize_serialization_schema_version<'de, D>(deserializer: D) -> Result<String, D::Error>
where
D: serde::Deserializer<'de>,
{
#[derive(Deserialize)]
#[serde(untagged)]
enum SchemaVersionValue {
String(String),
Integer(u64),
}
let parsed = SchemaVersionValue::deserialize(deserializer)?;
Ok(match parsed {
SchemaVersionValue::String(version) => normalize_serialization_schema_version(&version),
SchemaVersionValue::Integer(version) => {
normalize_serialization_schema_version(&version.to_string())
}
})
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CanonicalTraceV1 {
#[serde(
default = "default_serialization_schema_version",
deserialize_with = "deserialize_serialization_schema_version"
)]
pub schema_version: String,
pub events: Vec<ObsEvent>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CanonicalReplayFragmentV1 {
#[serde(
default = "default_serialization_schema_version",
deserialize_with = "deserialize_serialization_schema_version"
)]
pub schema_version: String,
pub obs_trace: Vec<ObsEvent>,
pub effect_trace: Vec<EffectTraceEntry>,
pub crashed_sites: Vec<String>,
pub partitioned_edges: Vec<(String, String)>,
pub corrupted_edges: Vec<((String, String), CorruptionType)>,
pub timed_out_sites: Vec<(String, u64)>,
#[serde(default)]
pub effect_determinism_tier: EffectDeterminismTier,
#[serde(default)]
pub communication_replay_mode: CommunicationReplayMode,
#[serde(default)]
pub communication_replay_root: Option<Hash>,
#[serde(default)]
pub communication_consumption_artifacts: Vec<CommunicationConsumptionArtifact>,
#[serde(default)]
pub semantic_audit_log: Vec<SemanticAuditRecord>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum SemanticAuditRecord {
Authority {
tick: Option<u64>,
session: Option<SessionId>,
artifact: AuthorityArtifact,
event: AuthorityAuditEvent,
reason: Option<String>,
},
Delegation {
tick: u64,
session: SessionId,
receipt: DelegationReceipt,
status: DelegationStatus,
reason: Option<String>,
},
FailureBranch {
tick: u64,
session: SessionId,
coro_id: usize,
fault: crate::coroutine::Fault,
},
TimeoutIssued {
tick: u64,
site: String,
until_tick: u64,
witness_id: AuthorityWitnessId,
},
CancellationRequested {
tick: u64,
session: SessionId,
witness_id: AuthorityWitnessId,
owner_id: FragmentOwnerId,
reason: OwnershipTerminalReason,
},
Cancelled {
tick: u64,
session: SessionId,
witness_id: AuthorityWitnessId,
reason: OwnershipTerminalReason,
},
SessionTerminal {
tick: u64,
session: SessionId,
reason: SessionTerminalReason,
},
EffectObservation {
effect_id: u64,
ordering_key: u64,
session: Option<SessionId>,
effect_kind: String,
effect_interface: Option<String>,
effect_operation: Option<String>,
handler_identity: String,
inputs: JsonValue,
outputs: JsonValue,
},
}
#[must_use]
pub fn canonical_trace_v1(trace: &[ObsEvent]) -> CanonicalTraceV1 {
CanonicalTraceV1 {
schema_version: default_serialization_schema_version(),
events: normalize_trace(trace),
}
}
#[must_use]
pub fn canonical_effect_trace(trace: &[EffectTraceEntry]) -> Vec<EffectTraceEntry> {
let mut out = trace.to_vec();
out.sort_by(|lhs, rhs| {
(lhs.ordering_key, lhs.effect_id, &lhs.effect_kind).cmp(&(
rhs.ordering_key,
rhs.effect_id,
&rhs.effect_kind,
))
});
out
}
fn authority_artifact_session(artifact: &AuthorityArtifact) -> Option<SessionId> {
match artifact {
AuthorityArtifact::Readiness(witness) => Some(witness.session_id),
AuthorityArtifact::Cancellation(witness) => Some(witness.session_id),
AuthorityArtifact::Timeout(_) => None,
}
}
fn effect_entry_session(entry: &EffectTraceEntry) -> Option<SessionId> {
entry
.inputs
.get("session")
.and_then(JsonValue::as_u64)
.and_then(|sid| usize::try_from(sid).ok())
.or_else(|| {
entry
.inputs
.get("sid")
.and_then(JsonValue::as_u64)
.and_then(|sid| usize::try_from(sid).ok())
})
}
fn semantic_rank(record: &SemanticAuditRecord) -> u8 {
match record {
SemanticAuditRecord::Authority { .. } => 0,
SemanticAuditRecord::Delegation { .. } => 1,
SemanticAuditRecord::FailureBranch { .. } => 2,
SemanticAuditRecord::TimeoutIssued { .. } => 3,
SemanticAuditRecord::CancellationRequested { .. } => 4,
SemanticAuditRecord::Cancelled { .. } => 5,
SemanticAuditRecord::SessionTerminal { .. } => 6,
SemanticAuditRecord::EffectObservation { .. } => 7,
}
}
fn semantic_tick(record: &SemanticAuditRecord) -> u64 {
match record {
SemanticAuditRecord::Authority { tick, .. } => tick.unwrap_or(0),
SemanticAuditRecord::Delegation { tick, .. }
| SemanticAuditRecord::FailureBranch { tick, .. }
| SemanticAuditRecord::TimeoutIssued { tick, .. }
| SemanticAuditRecord::CancellationRequested { tick, .. }
| SemanticAuditRecord::Cancelled { tick, .. }
| SemanticAuditRecord::SessionTerminal { tick, .. } => *tick,
SemanticAuditRecord::EffectObservation { ordering_key, .. } => *ordering_key,
}
}
#[must_use]
pub fn canonical_semantic_audit_log(records: &[SemanticAuditRecord]) -> Vec<SemanticAuditRecord> {
let mut out = records.to_vec();
out.sort_by(|lhs, rhs| {
let lhs_key = (
semantic_tick(lhs),
semantic_rank(lhs),
serde_json::to_string(lhs).unwrap_or_default(),
);
let rhs_key = (
semantic_tick(rhs),
semantic_rank(rhs),
serde_json::to_string(rhs).unwrap_or_default(),
);
lhs_key.cmp(&rhs_key)
});
out
}
#[must_use]
pub fn semantic_audit_log_v1(
authority_audit_log: &[AuthorityAuditRecord],
delegation_audit_log: &[DelegationAuditRecord],
obs_trace: &[ObsEvent],
effect_trace: &[EffectTraceEntry],
) -> Vec<SemanticAuditRecord> {
let mut records = Vec::new();
records.extend(authority_audit_log.iter().cloned().map(|record| {
SemanticAuditRecord::Authority {
tick: record.tick,
session: authority_artifact_session(&record.artifact),
artifact: record.artifact,
event: record.event,
reason: record.reason,
}
}));
records.extend(delegation_audit_log.iter().cloned().map(|record| {
SemanticAuditRecord::Delegation {
tick: record.tick,
session: record.receipt.session,
receipt: record.receipt,
status: record.status,
reason: record.reason,
}
}));
records.extend(obs_trace.iter().filter_map(|event| match event {
ObsEvent::FailureBranchEntered {
tick,
session,
coro_id,
fault,
} => Some(SemanticAuditRecord::FailureBranch {
tick: *tick,
session: *session,
coro_id: *coro_id,
fault: fault.clone(),
}),
ObsEvent::TimeoutIssued {
tick,
site,
until_tick,
witness_id,
} => Some(SemanticAuditRecord::TimeoutIssued {
tick: *tick,
site: site.clone(),
until_tick: *until_tick,
witness_id: *witness_id,
}),
ObsEvent::CancellationRequested {
tick,
session,
witness_id,
owner_id,
reason,
} => Some(SemanticAuditRecord::CancellationRequested {
tick: *tick,
session: *session,
witness_id: *witness_id,
owner_id: owner_id.clone(),
reason: reason.clone(),
}),
ObsEvent::Cancelled {
tick,
session,
witness_id,
reason,
} => Some(SemanticAuditRecord::Cancelled {
tick: *tick,
session: *session,
witness_id: *witness_id,
reason: reason.clone(),
}),
ObsEvent::SessionTerminal {
tick,
session,
reason,
} => Some(SemanticAuditRecord::SessionTerminal {
tick: *tick,
session: *session,
reason: reason.clone(),
}),
_ => None,
}));
records.extend(effect_trace.iter().cloned().map(|entry| {
SemanticAuditRecord::EffectObservation {
effect_id: entry.effect_id,
ordering_key: entry.ordering_key,
session: effect_entry_session(&entry),
effect_kind: entry.effect_kind,
effect_interface: entry.effect_interface,
effect_operation: entry.effect_operation,
handler_identity: entry.handler_identity,
inputs: entry.inputs,
outputs: entry.outputs,
}
}));
canonical_semantic_audit_log(&records)
}
#[must_use]
#[allow(clippy::too_many_arguments)]
pub fn canonical_replay_fragment_v1(
obs_trace: &[ObsEvent],
effect_trace: &[EffectTraceEntry],
authority_audit_log: &[AuthorityAuditRecord],
delegation_audit_log: &[DelegationAuditRecord],
mut crashed_sites: Vec<String>,
mut partitioned_edges: Vec<(String, String)>,
mut corrupted_edges: Vec<((String, String), CorruptionType)>,
mut timed_out_sites: Vec<(String, u64)>,
effect_determinism_tier: EffectDeterminismTier,
communication_replay_mode: CommunicationReplayMode,
communication_replay_root: Option<Hash>,
communication_consumption_artifacts: Vec<CommunicationConsumptionArtifact>,
) -> CanonicalReplayFragmentV1 {
crashed_sites.sort_unstable();
crashed_sites.dedup();
partitioned_edges.sort_unstable();
partitioned_edges.dedup();
corrupted_edges.sort_by(|lhs, rhs| lhs.0.cmp(&rhs.0));
corrupted_edges.dedup();
timed_out_sites.sort_unstable();
CanonicalReplayFragmentV1 {
schema_version: default_serialization_schema_version(),
obs_trace: canonical_trace_v1(obs_trace).events,
effect_trace: canonical_effect_trace(effect_trace),
crashed_sites,
partitioned_edges,
corrupted_edges,
timed_out_sites,
effect_determinism_tier,
communication_replay_mode,
communication_replay_root,
communication_consumption_artifacts,
semantic_audit_log: semantic_audit_log_v1(
authority_audit_log,
delegation_audit_log,
obs_trace,
effect_trace,
),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::session::Edge;
#[test]
fn canonical_effect_trace_is_stably_sorted() {
let trace = vec![
EffectTraceEntry {
effect_id: 2,
effect_kind: "b".to_string(),
inputs: serde_json::json!({}),
outputs: serde_json::json!({}),
handler_identity: "h".to_string(),
effect_interface: None,
effect_operation: None,
ordering_key: 3,
topology: None,
},
EffectTraceEntry {
effect_id: 1,
effect_kind: "a".to_string(),
inputs: serde_json::json!({}),
outputs: serde_json::json!({}),
handler_identity: "h".to_string(),
effect_interface: None,
effect_operation: None,
ordering_key: 2,
topology: None,
},
];
let sorted = canonical_effect_trace(&trace);
assert_eq!(sorted[0].effect_id, 1);
assert_eq!(sorted[1].effect_id, 2);
}
#[test]
fn canonical_trace_payload_has_version() {
let trace = vec![ObsEvent::Sent {
tick: 1,
edge: Edge::new(1, "A", "B"),
session: 1,
from: "A".to_string(),
to: "B".to_string(),
label: "m".to_string(),
}];
let payload = canonical_trace_v1(&trace);
assert_eq!(payload.schema_version, SERIALIZATION_SCHEMA_VERSION);
assert_eq!(payload.events.len(), 1);
}
#[test]
fn legacy_numeric_schema_version_deserializes_to_string_identifier() {
let payload = serde_json::json!({
"schema_version": 1,
"events": []
});
let decoded: CanonicalTraceV1 =
serde_json::from_value(payload).expect("legacy schema version should deserialize");
assert_eq!(decoded.schema_version, SERIALIZATION_SCHEMA_VERSION);
}
}