use serde::{Deserialize, Serialize};
use std::io::{self, Read};
use super::{Frame, Frames, RunRecord};
use crate::ContextCause;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RegionCommit {
pub region: String,
pub digest_before: String,
pub digest_after: String,
pub tokens_before: usize,
pub tokens_after: usize,
pub entries_before: usize,
pub entries_after: usize,
pub entries_added: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RegionTransition {
pub region: String,
pub digest_before: Option<String>,
pub digest_after: Option<String>,
pub tokens_before: Option<usize>,
pub tokens_after: Option<usize>,
pub token_delta: i64,
pub entries_before: Option<usize>,
pub entries_after: Option<usize>,
pub entries_added: usize,
pub entries_removed: usize,
}
impl From<RegionCommit> for RegionTransition {
fn from(commit: RegionCommit) -> Self {
Self {
region: commit.region,
digest_before: Some(commit.digest_before),
digest_after: Some(commit.digest_after),
tokens_before: Some(commit.tokens_before),
tokens_after: Some(commit.tokens_after),
token_delta: commit.tokens_after as i64 - commit.tokens_before as i64,
entries_before: Some(commit.entries_before),
entries_after: Some(commit.entries_after),
entries_added: commit.entries_added,
entries_removed: (commit.entries_before + commit.entries_added)
.saturating_sub(commit.entries_after),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ContextChangeRecord {
pub cause: ContextCause,
pub revision_before: Option<String>,
pub revision_after: Option<String>,
pub execution_id: Option<String>,
pub regions: Vec<RegionTransition>,
pub at: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IndexedChange {
pub position: u64,
pub record: ContextChangeRecord,
}
pub(super) fn change_of(record: &RunRecord) -> Option<ContextChangeRecord> {
match record {
RunRecord::ContextTransaction {
revision_before,
revision_after,
cause,
regions,
execution_id,
at,
} => Some(ContextChangeRecord {
cause: *cause,
revision_before: Some(revision_before.clone()),
revision_after: Some(revision_after.clone()),
execution_id: Some(execution_id.clone()).filter(|id| !id.is_empty()),
regions: regions
.iter()
.cloned()
.map(RegionTransition::from)
.collect(),
at: *at,
}),
RunRecord::ContextChange {
region,
cause,
entries_added,
entries_removed,
token_delta,
at,
} => Some(ContextChangeRecord {
cause: *cause,
revision_before: None,
revision_after: None,
execution_id: None,
regions: vec![RegionTransition {
region: region.clone(),
digest_before: None,
digest_after: None,
tokens_before: None,
tokens_after: None,
token_delta: *token_delta,
entries_before: None,
entries_after: None,
entries_added: *entries_added,
entries_removed: *entries_removed,
}],
at: *at,
}),
_ => None,
}
}
pub fn read_archive_changes(r: &mut dyn Read) -> io::Result<Vec<IndexedChange>> {
let (_, mut frames) = Frames::open(r)?;
let mut changes = Vec::new();
while let Ok(Some((position, frame))) = frames.next_frame() {
let Frame::Record(record) = frame else {
continue;
};
if let Some(record) = change_of(&record) {
changes.push(IndexedChange { position, record });
}
}
Ok(changes)
}
#[cfg(test)]
#[path = "transaction_tests.rs"]
mod tests;