use std::collections::HashMap;
use aion_core::{ActivityEvent, ActivityEventKind};
use serde_json::{Value, json};
use super::patch;
use super::slot::{RecordedBase, StreamSlot};
use super::wire::{self, ENVELOPE_DELTA_KEY, ENVELOPE_DELTA_VERSION};
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum Compaction {
NotApplicable,
BaseRecorded {
base: String,
},
Compacted {
base: String,
full_bytes: usize,
delta_bytes: usize,
},
FullRetained {
base: String,
reason: FullRetainedReason,
},
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum FullRetainedReason {
NotByteReversible,
NotSmaller {
full_bytes: usize,
delta_bytes: usize,
},
NotSerializable,
}
#[derive(Debug, Default)]
pub struct EnvelopeDeltaEncoder {
bases: HashMap<StreamSlot, RecordedBase>,
}
impl EnvelopeDeltaEncoder {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn compact(&mut self, event: &mut ActivityEvent) -> Compaction {
if event.ephemeral {
return Compaction::NotApplicable;
}
let slot = StreamSlot::of(event);
let worker_seq = event.worker_seq;
let ActivityEventKind::Raw { value, .. } = &mut event.kind else {
return Compaction::NotApplicable;
};
if wire::delta_document(value).is_some() {
return Compaction::NotApplicable;
}
let Some(response_id) = wire::envelope_response_id(value) else {
return Compaction::NotApplicable;
};
let response_id = response_id.to_owned();
match self.bases.get(&slot) {
Some(base) if base.response_id == response_id => match build_delta(base, value) {
Ok(delta) => {
*value = delta.document;
Compaction::Compacted {
base: response_id,
full_bytes: delta.full_bytes,
delta_bytes: delta.delta_bytes,
}
}
Err(reason) => Compaction::FullRetained {
base: response_id,
reason,
},
},
_ => match RecordedBase::record_parts(&response_id, worker_seq, value) {
Some(base) => {
self.bases.insert(slot, base);
Compaction::BaseRecorded { base: response_id }
}
None => Compaction::FullRetained {
base: response_id,
reason: FullRetainedReason::NotSerializable,
},
},
}
}
}
struct BuiltDelta {
document: Value,
full_bytes: usize,
delta_bytes: usize,
}
fn build_delta(base: &RecordedBase, frame: &Value) -> Result<BuiltDelta, FullRetainedReason> {
let (Value::Object(base_object), Value::Object(frame_object)) = (&base.value, frame) else {
return Err(FullRetainedReason::NotByteReversible);
};
let diff = patch::diff(base_object, frame_object);
let document = json!({
ENVELOPE_DELTA_KEY: {
"v": ENVELOPE_DELTA_VERSION,
"base": base.response_id,
"base_worker_seq": base.worker_seq,
"base_digest": base.digest,
"set": Value::Object(diff.set.clone()),
"unset": diff.unset.clone(),
}
});
let reconstructed =
patch::apply(&base.value, &diff).map_err(|_| FullRetainedReason::NotByteReversible)?;
let (Ok(frame_bytes), Ok(reconstructed_bytes)) = (
serde_json::to_vec(frame),
serde_json::to_vec(&reconstructed),
) else {
return Err(FullRetainedReason::NotSerializable);
};
if frame_bytes != reconstructed_bytes {
return Err(FullRetainedReason::NotByteReversible);
}
let Ok(document_bytes) = serde_json::to_vec(&document) else {
return Err(FullRetainedReason::NotSerializable);
};
let full_bytes = frame_bytes.len();
let delta_bytes = document_bytes.len();
if delta_bytes >= full_bytes {
return Err(FullRetainedReason::NotSmaller {
full_bytes,
delta_bytes,
});
}
Ok(BuiltDelta {
document,
full_bytes,
delta_bytes,
})
}