use std::collections::HashMap;
use std::fmt;
use aion_core::{ActivityEvent, ActivityEventKind};
use serde_json::Value;
use super::patch::{self, Patch};
use super::slot::{RecordedBase, StreamSlot};
use super::wire::{self, ENVELOPE_DELTA_VERSION, EnvelopeDelta};
#[derive(Clone, Debug, PartialEq)]
pub enum ResolvedEnvelope {
Verbatim,
Reconstructed(Value),
Unresolved(UnresolvedDelta),
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct UnresolvedDelta {
pub base: String,
pub base_worker_seq: Option<u64>,
pub reason: UnresolvedReason,
}
impl fmt::Display for UnresolvedDelta {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "envelope delta for turn {:?}", self.base)?;
if let Some(worker_seq) = self.base_worker_seq {
write!(formatter, " (base at worker_seq {worker_seq})")?;
}
write!(formatter, " could not be resolved: {}", self.reason)
}
}
#[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)]
pub enum UnresolvedReason {
#[error(
"its base record was not in the range this reader has seen (trimmed by retention, or before the requested window)"
)]
BaseNotSeen,
#[error(
"the base record present here has digest {observed}, but the delta was built against {expected}"
)]
BaseDigestMismatch {
expected: String,
observed: String,
},
#[error("it declares delta format version {version}, and this reader implements {supported}")]
UnsupportedVersion {
version: u64,
supported: u64,
},
#[error("the delta document is malformed: {detail}")]
Malformed {
detail: String,
},
#[error("the delta did not apply to the base: {detail}")]
PatchFailed {
detail: String,
},
}
#[derive(Debug, Default)]
pub struct EnvelopeDeltaDecoder {
bases: HashMap<StreamSlot, RecordedBase>,
}
impl EnvelopeDeltaDecoder {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn resolve(&mut self, event: &ActivityEvent) -> ResolvedEnvelope {
let slot = StreamSlot::of(event);
let ActivityEventKind::Raw { value, .. } = &event.kind else {
return ResolvedEnvelope::Verbatim;
};
if let Some(document) = wire::delta_document(value) {
return self.resolve_delta(&slot, document);
}
let Some(response_id) = wire::envelope_response_id(value) else {
return ResolvedEnvelope::Verbatim;
};
if event.ephemeral {
return ResolvedEnvelope::Verbatim;
}
self.record_base(slot, response_id, event.worker_seq, value);
ResolvedEnvelope::Verbatim
}
pub fn resolve_event(
&mut self,
event: &ActivityEvent,
) -> (ActivityEvent, Option<UnresolvedDelta>) {
match self.resolve(event) {
ResolvedEnvelope::Verbatim => (event.clone(), None),
ResolvedEnvelope::Reconstructed(envelope) => {
let mut resolved = event.clone();
if let ActivityEventKind::Raw { value, .. } = &mut resolved.kind {
*value = envelope;
}
(resolved, None)
}
ResolvedEnvelope::Unresolved(unresolved) => (event.clone(), Some(unresolved)),
}
}
fn record_base(&mut self, slot: StreamSlot, response_id: &str, worker_seq: u64, value: &Value) {
if self
.bases
.get(&slot)
.is_some_and(|held| held.response_id == response_id)
{
return;
}
if let Some(base) = RecordedBase::record_parts(response_id, worker_seq, value) {
self.bases.insert(slot, base);
}
}
fn resolve_delta(&self, slot: &StreamSlot, document: &Value) -> ResolvedEnvelope {
let delta: EnvelopeDelta = match serde_json::from_value(document.clone()) {
Ok(delta) => delta,
Err(error) => {
return unresolved(
document
.get("base")
.and_then(Value::as_str)
.unwrap_or_default(),
document.get("base_worker_seq").and_then(Value::as_u64),
UnresolvedReason::Malformed {
detail: error.to_string(),
},
);
}
};
if delta.version != ENVELOPE_DELTA_VERSION {
return unresolved(
&delta.base,
Some(delta.base_worker_seq),
UnresolvedReason::UnsupportedVersion {
version: delta.version,
supported: ENVELOPE_DELTA_VERSION,
},
);
}
let Some(base) = self
.bases
.get(slot)
.filter(|held| held.response_id == delta.base)
else {
return unresolved(
&delta.base,
Some(delta.base_worker_seq),
UnresolvedReason::BaseNotSeen,
);
};
if base.digest != delta.base_digest {
return unresolved(
&delta.base,
Some(delta.base_worker_seq),
UnresolvedReason::BaseDigestMismatch {
expected: delta.base_digest,
observed: base.digest.clone(),
},
);
}
let reconstruction = patch::apply(
&base.value,
&Patch {
set: delta.set,
unset: delta.unset,
},
);
match reconstruction {
Ok(envelope) => ResolvedEnvelope::Reconstructed(envelope),
Err(error) => unresolved(
&delta.base,
Some(delta.base_worker_seq),
UnresolvedReason::PatchFailed {
detail: error.to_string(),
},
),
}
}
}
fn unresolved(
base: &str,
base_worker_seq: Option<u64>,
reason: UnresolvedReason,
) -> ResolvedEnvelope {
ResolvedEnvelope::Unresolved(UnresolvedDelta {
base: base.to_owned(),
base_worker_seq,
reason,
})
}