use std::{
collections::{BTreeMap, BTreeSet},
num::NonZeroUsize,
};
use eredu_core::{RealtimeFrameSlot, RealtimeSlotCoordinate, RealtimeSpeechConfig};
use crate::generation::TokenDomain;
#[derive(Debug, Clone, Copy, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct RealtimePayloadOwnerIdentity(u64);
impl RealtimePayloadOwnerIdentity {
pub const fn new(value: u64) -> Option<Self> {
if value == 0 {
None
} else {
Some(Self(value))
}
}
pub const fn value(self) -> u64 {
self.0
}
}
#[derive(Debug, Clone, Copy, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct RealtimePayloadGeneration(u64);
impl RealtimePayloadGeneration {
pub const fn new(value: u64) -> Option<Self> {
if value == 0 {
None
} else {
Some(Self(value))
}
}
pub const fn value(self) -> u64 {
self.0
}
}
#[derive(Debug, Clone, Eq, PartialEq)]
pub struct RealtimePayloadContract {
schedule: RealtimeSpeechConfig,
batch: NonZeroUsize,
text_domain: TokenDomain,
audio_domain: TokenDomain,
generation: RealtimePayloadGeneration,
owner: RealtimePayloadOwnerIdentity,
}
impl RealtimePayloadContract {
pub fn new(
schedule: RealtimeSpeechConfig,
batch: usize,
text_domain: TokenDomain,
audio_domain: TokenDomain,
generation: RealtimePayloadGeneration,
owner: RealtimePayloadOwnerIdentity,
) -> Result<Self, RealtimePayloadContractError> {
let batch = NonZeroUsize::new(batch).ok_or(RealtimePayloadContractError::EmptyBatch)?;
Ok(Self {
schedule,
batch,
text_domain,
audio_domain,
generation,
owner,
})
}
pub const fn schedule(&self) -> &RealtimeSpeechConfig {
&self.schedule
}
pub const fn batch(&self) -> NonZeroUsize {
self.batch
}
pub const fn text_domain(&self) -> TokenDomain {
self.text_domain
}
pub const fn audio_domain(&self) -> TokenDomain {
self.audio_domain
}
pub const fn generation(&self) -> RealtimePayloadGeneration {
self.generation
}
pub const fn owner(&self) -> RealtimePayloadOwnerIdentity {
self.owner
}
pub fn slot_domain(
&self,
slot: RealtimeFrameSlot,
) -> Result<TokenDomain, RealtimePayloadContractError> {
match slot {
RealtimeFrameSlot::Text => Ok(self.text_domain),
RealtimeFrameSlot::Audio(codebook)
if codebook < self.schedule.total_audio_codebooks() =>
{
Ok(self.audio_domain)
}
_ => Err(RealtimePayloadContractError::InvalidSlot {
slot,
total_audio_codebooks: self.schedule.total_audio_codebooks(),
}),
}
}
pub fn validate(&self, contract: &Self) -> Result<(), RealtimePayloadContractError> {
if self.schedule != contract.schedule {
return Err(RealtimePayloadContractError::ScheduleMismatch);
}
if self.batch != contract.batch {
return Err(RealtimePayloadContractError::BatchMismatch);
}
if self.text_domain != contract.text_domain {
return Err(RealtimePayloadContractError::TextDomainMismatch);
}
if self.audio_domain != contract.audio_domain {
return Err(RealtimePayloadContractError::AudioDomainMismatch);
}
if self.generation != contract.generation {
return Err(RealtimePayloadContractError::GenerationMismatch);
}
if self.owner != contract.owner {
return Err(RealtimePayloadContractError::OwnerMismatch);
}
Ok(())
}
}
#[derive(Debug, Clone, Eq, PartialEq)]
pub struct RealtimePayloadEnvelope<P> {
contract: RealtimePayloadContract,
coordinate: RealtimeSlotCoordinate,
domain: TokenDomain,
payload: P,
}
impl<P> RealtimePayloadEnvelope<P> {
pub fn new(
contract: RealtimePayloadContract,
coordinate: RealtimeSlotCoordinate,
payload: P,
) -> Result<Self, RealtimePayloadContractError> {
let domain = contract.slot_domain(coordinate.slot())?;
Ok(Self {
contract,
coordinate,
domain,
payload,
})
}
pub const fn contract(&self) -> &RealtimePayloadContract {
&self.contract
}
pub const fn coordinate(&self) -> RealtimeSlotCoordinate {
self.coordinate
}
pub const fn domain(&self) -> TokenDomain {
self.domain
}
pub const fn payload(&self) -> &P {
&self.payload
}
pub fn validate(
&self,
contract: &RealtimePayloadContract,
coordinate: RealtimeSlotCoordinate,
) -> Result<(), RealtimePayloadContractError> {
self.contract.validate(contract)?;
if self.coordinate != coordinate {
return Err(RealtimePayloadContractError::CoordinateMismatch);
}
let domain = contract.slot_domain(coordinate.slot())?;
if self.domain != domain {
return Err(RealtimePayloadContractError::SlotDomainMismatch);
}
Ok(())
}
pub fn into_parts(
self,
) -> (
RealtimePayloadContract,
RealtimeSlotCoordinate,
TokenDomain,
P,
) {
(self.contract, self.coordinate, self.domain, self.payload)
}
}
#[derive(Debug, Clone, Eq, PartialEq, thiserror::Error)]
#[non_exhaustive]
pub enum RealtimePayloadContractError {
#[error("realtime payload contract batch is empty")]
EmptyBatch,
#[error("realtime payload contract schedule does not match")]
ScheduleMismatch,
#[error("realtime payload contract batch does not match")]
BatchMismatch,
#[error("realtime payload contract text-token domain does not match")]
TextDomainMismatch,
#[error("realtime payload contract audio-token domain does not match")]
AudioDomainMismatch,
#[error("realtime payload contract generation does not match")]
GenerationMismatch,
#[error("realtime payload contract owner does not match")]
OwnerMismatch,
#[error("realtime payload coordinate does not match")]
CoordinateMismatch,
#[error("realtime payload slot token domain does not match")]
SlotDomainMismatch,
#[error(
"realtime payload slot {slot:?} is outside text plus {total_audio_codebooks} audio codebooks"
)]
InvalidSlot {
slot: RealtimeFrameSlot,
total_audio_codebooks: usize,
},
}
#[derive(Debug, Clone, Eq, PartialEq)]
pub struct RealtimePayloadHistory<P> {
schedule: RealtimeSpeechConfig,
contract: Option<RealtimePayloadContract>,
payloads: BTreeMap<RealtimeSlotCoordinate, RealtimePayloadEnvelope<P>>,
}
impl<P> RealtimePayloadHistory<P> {
pub fn new(schedule: RealtimeSpeechConfig) -> Self {
Self {
schedule,
contract: None,
payloads: BTreeMap::new(),
}
}
pub fn with_contract(contract: RealtimePayloadContract) -> Self {
Self {
schedule: contract.schedule().clone(),
contract: Some(contract),
payloads: BTreeMap::new(),
}
}
pub const fn schedule(&self) -> &RealtimeSpeechConfig {
&self.schedule
}
pub const fn contract(&self) -> Option<&RealtimePayloadContract> {
self.contract.as_ref()
}
pub fn bind_or_validate_contract(
&mut self,
contract: &RealtimePayloadContract,
) -> Result<(), RealtimePayloadHistoryError> {
if &self.schedule != contract.schedule() {
return Err(RealtimePayloadHistoryError::PayloadContract(
RealtimePayloadContractError::ScheduleMismatch,
));
}
if let Some(current) = &self.contract {
current
.validate(contract)
.map_err(RealtimePayloadHistoryError::PayloadContract)
} else {
debug_assert!(self.payloads.is_empty());
self.contract = Some(contract.clone());
Ok(())
}
}
pub fn validate_successor(&self, successor: &Self) -> Result<(), RealtimePayloadHistoryError> {
self.validate_schedule(&successor.schedule)?;
match (&self.contract, &successor.contract) {
(Some(current), Some(candidate)) => current
.validate(candidate)
.map_err(RealtimePayloadHistoryError::PayloadContract),
(None, Some(_)) if self.payloads.is_empty() => Ok(()),
(None, None) => Ok(()),
(Some(_), None) | (None, Some(_)) => Err(RealtimePayloadHistoryError::UnboundContract),
}
}
pub fn len(&self) -> usize {
self.payloads.len()
}
pub fn is_empty(&self) -> bool {
self.payloads.is_empty()
}
pub fn validate_schedule(
&self,
schedule: &RealtimeSpeechConfig,
) -> Result<(), RealtimePayloadHistoryError> {
if &self.schedule == schedule {
Ok(())
} else {
Err(RealtimePayloadHistoryError::ScheduleMismatch)
}
}
pub fn delayed_coordinate(
&self,
schedule: &RealtimeSpeechConfig,
base_position: usize,
slot: RealtimeFrameSlot,
) -> Result<RealtimeSlotCoordinate, RealtimePayloadHistoryError> {
self.validate_schedule(schedule)?;
let delay = self.slot_delay(slot)?;
let position = base_position.checked_add(delay).ok_or(
RealtimePayloadHistoryError::CoordinateOverflow {
base_position,
delay,
},
)?;
Ok(RealtimeSlotCoordinate::new(position, slot))
}
pub fn insert(
&mut self,
schedule: &RealtimeSpeechConfig,
coordinate: RealtimeSlotCoordinate,
payload: P,
) -> Result<(), RealtimePayloadHistoryError> {
self.insert_many(schedule, [(coordinate, payload)])
}
pub fn insert_delayed(
&mut self,
schedule: &RealtimeSpeechConfig,
base_position: usize,
slot: RealtimeFrameSlot,
payload: P,
) -> Result<RealtimeSlotCoordinate, RealtimePayloadHistoryError> {
let coordinate = self.delayed_coordinate(schedule, base_position, slot)?;
self.insert(schedule, coordinate, payload)?;
Ok(coordinate)
}
pub fn insert_many(
&mut self,
schedule: &RealtimeSpeechConfig,
payloads: impl IntoIterator<Item = (RealtimeSlotCoordinate, P)>,
) -> Result<(), RealtimePayloadHistoryError> {
self.validate_schedule(schedule)?;
let contract = self
.contract
.as_ref()
.ok_or(RealtimePayloadHistoryError::UnboundContract)?
.clone();
let payloads = payloads.into_iter().collect::<Vec<_>>();
let mut pending = BTreeSet::new();
for (coordinate, _) in &payloads {
self.validate_coordinate(*coordinate)?;
if self.payloads.contains_key(coordinate) || !pending.insert(*coordinate) {
return Err(RealtimePayloadHistoryError::DuplicatePayload {
coordinate: *coordinate,
});
}
}
let envelopes = payloads
.into_iter()
.map(|(coordinate, payload)| {
RealtimePayloadEnvelope::new(contract.clone(), coordinate, payload)
.map(|envelope| (coordinate, envelope))
.map_err(RealtimePayloadHistoryError::PayloadContract)
})
.collect::<Result<Vec<_>, _>>()?;
self.payloads.extend(envelopes);
Ok(())
}
pub fn overwrite_many(
&mut self,
schedule: &RealtimeSpeechConfig,
payloads: impl IntoIterator<Item = (RealtimeSlotCoordinate, P)>,
) -> Result<(), RealtimePayloadHistoryError> {
self.validate_schedule(schedule)?;
let contract = self
.contract
.as_ref()
.ok_or(RealtimePayloadHistoryError::UnboundContract)?
.clone();
let payloads = payloads.into_iter().collect::<Vec<_>>();
for (coordinate, _) in &payloads {
self.validate_coordinate(*coordinate)?;
}
let envelopes = payloads
.into_iter()
.map(|(coordinate, payload)| {
RealtimePayloadEnvelope::new(contract.clone(), coordinate, payload)
.map(|envelope| (coordinate, envelope))
.map_err(RealtimePayloadHistoryError::PayloadContract)
})
.collect::<Result<Vec<_>, _>>()?;
self.payloads.extend(envelopes);
Ok(())
}
pub fn get(
&self,
schedule: &RealtimeSpeechConfig,
coordinate: RealtimeSlotCoordinate,
) -> Result<Option<&P>, RealtimePayloadHistoryError> {
self.validate_schedule(schedule)?;
self.validate_coordinate(coordinate)?;
Ok(self
.payloads
.get(&coordinate)
.map(RealtimePayloadEnvelope::payload))
}
pub fn envelope(
&self,
schedule: &RealtimeSpeechConfig,
coordinate: RealtimeSlotCoordinate,
) -> Result<Option<&RealtimePayloadEnvelope<P>>, RealtimePayloadHistoryError> {
self.validate_schedule(schedule)?;
self.validate_coordinate(coordinate)?;
Ok(self.payloads.get(&coordinate))
}
pub fn required(
&self,
schedule: &RealtimeSpeechConfig,
coordinate: RealtimeSlotCoordinate,
) -> Result<&P, RealtimePayloadHistoryError> {
self.get(schedule, coordinate)?
.ok_or(RealtimePayloadHistoryError::MissingPayload { coordinate })
}
pub fn resolve_required(
&self,
schedule: &RealtimeSpeechConfig,
coordinates: impl IntoIterator<Item = RealtimeSlotCoordinate>,
) -> Result<Vec<&P>, RealtimePayloadHistoryError> {
self.validate_schedule(schedule)?;
coordinates
.into_iter()
.map(|coordinate| self.required(schedule, coordinate))
.collect()
}
pub fn prune_for_next_frontier(
&mut self,
schedule: &RealtimeSpeechConfig,
next_frontier: usize,
) -> Result<usize, RealtimePayloadHistoryError> {
self.validate_schedule(schedule)?;
let retained_positions = schedule.max_delay().checked_add(2).ok_or(
RealtimePayloadHistoryError::RetentionWindowOverflow {
max_delay: schedule.max_delay(),
},
)?;
let minimum = next_frontier.saturating_sub(retained_positions);
let previous = self.payloads.len();
self.payloads
.retain(|coordinate, _| coordinate.position() >= minimum);
Ok(previous - self.payloads.len())
}
pub fn retained_values(&self) -> impl Iterator<Item = &P> {
self.payloads.values().map(RealtimePayloadEnvelope::payload)
}
pub fn envelopes(&self) -> impl Iterator<Item = &RealtimePayloadEnvelope<P>> {
self.payloads.values()
}
pub fn entries(&self) -> impl Iterator<Item = (RealtimeSlotCoordinate, &P)> {
self.payloads
.iter()
.map(|(coordinate, envelope)| (*coordinate, envelope.payload()))
}
fn validate_coordinate(
&self,
coordinate: RealtimeSlotCoordinate,
) -> Result<(), RealtimePayloadHistoryError> {
self.slot_delay(coordinate.slot()).map(|_| ())
}
fn slot_delay(&self, slot: RealtimeFrameSlot) -> Result<usize, RealtimePayloadHistoryError> {
match slot {
RealtimeFrameSlot::Text => Ok(self.schedule.text_delay()),
RealtimeFrameSlot::Audio(codebook) => {
self.schedule.audio_delays().get(codebook).copied().ok_or(
RealtimePayloadHistoryError::InvalidSlot {
slot,
total_audio_codebooks: self.schedule.total_audio_codebooks(),
},
)
}
_ => Err(RealtimePayloadHistoryError::InvalidSlot {
slot,
total_audio_codebooks: self.schedule.total_audio_codebooks(),
}),
}
}
}
#[derive(Debug, Clone, Eq, PartialEq, thiserror::Error)]
#[non_exhaustive]
pub enum RealtimePayloadHistoryError {
#[error("realtime payload history has no bound payload contract")]
UnboundContract,
#[error(transparent)]
PayloadContract(RealtimePayloadContractError),
#[error("realtime payload history does not match the normalized schedule")]
ScheduleMismatch,
#[error(
"realtime payload slot {slot:?} is outside text plus {total_audio_codebooks} audio codebooks"
)]
InvalidSlot {
slot: RealtimeFrameSlot,
total_audio_codebooks: usize,
},
#[error("realtime coordinate {coordinate:?} already contains a payload")]
DuplicatePayload {
coordinate: RealtimeSlotCoordinate,
},
#[error("realtime coordinate {coordinate:?} has no retained payload")]
MissingPayload {
coordinate: RealtimeSlotCoordinate,
},
#[error("realtime payload coordinate overflowed from base {base_position} plus delay {delay}")]
CoordinateOverflow {
base_position: usize,
delay: usize,
},
#[error("realtime payload retention window overflowed for maximum delay {max_delay}")]
RetentionWindowOverflow {
max_delay: usize,
},
}
#[cfg(test)]
mod tests {
use eredu_core::RealtimeFrameConvention;
use super::*;
fn schedule() -> RealtimeSpeechConfig {
RealtimeSpeechConfig::new(
2,
1,
1,
1,
0,
1,
RealtimeFrameConvention::FeedbackAlignedHistory,
vec![2, 0, 3],
)
.unwrap()
}
fn coordinate(position: usize, slot: RealtimeFrameSlot) -> RealtimeSlotCoordinate {
RealtimeSlotCoordinate::new(position, slot)
}
fn payload_contract(
schedule: RealtimeSpeechConfig,
batch: usize,
text_domain: TokenDomain,
audio_domain: TokenDomain,
generation: u64,
owner: u64,
) -> RealtimePayloadContract {
RealtimePayloadContract::new(
schedule,
batch,
text_domain,
audio_domain,
RealtimePayloadGeneration::new(generation).unwrap(),
RealtimePayloadOwnerIdentity::new(owner).unwrap(),
)
.unwrap()
}
fn exact_payload_contract() -> RealtimePayloadContract {
payload_contract(
schedule(),
2,
TokenDomain::new(32),
TokenDomain::new(16),
7,
3,
)
}
fn payload_history<P>() -> RealtimePayloadHistory<P> {
RealtimePayloadHistory::with_contract(exact_payload_contract())
}
#[test]
fn payload_contract_requires_nonempty_batch_and_typed_nonempty_identities() {
assert_eq!(
RealtimePayloadContract::new(
schedule(),
0,
TokenDomain::new(32),
TokenDomain::new(16),
RealtimePayloadGeneration::new(7).unwrap(),
RealtimePayloadOwnerIdentity::new(3).unwrap(),
),
Err(RealtimePayloadContractError::EmptyBatch)
);
assert_eq!(RealtimePayloadGeneration::new(0), None);
assert_eq!(RealtimePayloadOwnerIdentity::new(0), None);
let contract = exact_payload_contract();
assert_eq!(contract.schedule(), &schedule());
assert_eq!(contract.batch().get(), 2);
assert_eq!(contract.text_domain(), TokenDomain::new(32));
assert_eq!(contract.audio_domain(), TokenDomain::new(16));
assert_eq!(contract.generation().value(), 7);
assert_eq!(contract.owner().value(), 3);
}
#[test]
fn payload_contract_rejects_each_exact_identity_perturbation() {
let exact = exact_payload_contract();
let other_schedule = RealtimeSpeechConfig::new(
2,
1,
1,
1,
0,
1,
RealtimeFrameConvention::AbsoluteDelayedSlots,
vec![2, 0, 3],
)
.unwrap();
assert_eq!(
exact.validate(&payload_contract(
other_schedule,
2,
TokenDomain::new(32),
TokenDomain::new(16),
7,
3,
)),
Err(RealtimePayloadContractError::ScheduleMismatch)
);
assert_eq!(
exact.validate(&payload_contract(
schedule(),
3,
TokenDomain::new(32),
TokenDomain::new(16),
7,
3,
)),
Err(RealtimePayloadContractError::BatchMismatch)
);
assert_eq!(
exact.validate(&payload_contract(
schedule(),
2,
TokenDomain::new(33),
TokenDomain::new(16),
7,
3,
)),
Err(RealtimePayloadContractError::TextDomainMismatch)
);
assert_eq!(
exact.validate(&payload_contract(
schedule(),
2,
TokenDomain::new(32),
TokenDomain::new(17),
7,
3,
)),
Err(RealtimePayloadContractError::AudioDomainMismatch)
);
assert_eq!(
exact.validate(&payload_contract(
schedule(),
2,
TokenDomain::new(32),
TokenDomain::new(16),
8,
3,
)),
Err(RealtimePayloadContractError::GenerationMismatch)
);
assert_eq!(
exact.validate(&payload_contract(
schedule(),
2,
TokenDomain::new(32),
TokenDomain::new(16),
7,
4,
)),
Err(RealtimePayloadContractError::OwnerMismatch)
);
}
#[test]
fn payload_envelope_derives_slot_domain_and_validates_exact_coordinate() {
let contract = exact_payload_contract();
let text_coordinate = coordinate(5, RealtimeFrameSlot::Text);
let text = RealtimePayloadEnvelope::new(contract.clone(), text_coordinate, "text").unwrap();
assert_eq!(text.contract(), &contract);
assert_eq!(text.coordinate(), text_coordinate);
assert_eq!(text.domain(), TokenDomain::new(32));
assert_eq!(text.payload(), &"text");
assert_eq!(text.validate(&contract, text_coordinate), Ok(()));
let audio_coordinate = coordinate(5, RealtimeFrameSlot::Audio(1));
let audio =
RealtimePayloadEnvelope::new(contract.clone(), audio_coordinate, "audio").unwrap();
assert_eq!(audio.domain(), TokenDomain::new(16));
assert_eq!(audio.validate(&contract, audio_coordinate), Ok(()));
assert_eq!(
audio.validate(&contract, coordinate(6, RealtimeFrameSlot::Audio(1))),
Err(RealtimePayloadContractError::CoordinateMismatch)
);
assert_eq!(
RealtimePayloadEnvelope::new(
contract,
coordinate(5, RealtimeFrameSlot::Audio(2)),
"invalid",
),
Err(RealtimePayloadContractError::InvalidSlot {
slot: RealtimeFrameSlot::Audio(2),
total_audio_codebooks: 2,
})
);
}
#[test]
fn history_requires_first_contract_binding_and_stores_only_typed_envelopes() {
let schedule = schedule();
let text = coordinate(0, RealtimeFrameSlot::Text);
let audio = coordinate(0, RealtimeFrameSlot::Audio(0));
let mut history = RealtimePayloadHistory::new(schedule.clone());
assert!(history.contract().is_none());
assert_eq!(history.get(&schedule, text), Ok(None));
assert_eq!(
history.insert(&schedule, text, 3),
Err(RealtimePayloadHistoryError::UnboundContract)
);
assert!(history.is_empty());
let contract = exact_payload_contract();
history.bind_or_validate_contract(&contract).unwrap();
history
.insert_many(&schedule, [(text, 3), (audio, 5)])
.unwrap();
assert_eq!(history.contract(), Some(&contract));
let envelopes = history.envelopes().collect::<Vec<_>>();
assert_eq!(envelopes.len(), 2);
assert_eq!(envelopes[0].contract(), &contract);
assert_eq!(envelopes[0].coordinate(), text);
assert_eq!(envelopes[0].domain(), TokenDomain::new(32));
assert_eq!(envelopes[0].payload(), &3);
assert_eq!(envelopes[1].contract(), &contract);
assert_eq!(envelopes[1].coordinate(), audio);
assert_eq!(envelopes[1].domain(), TokenDomain::new(16));
assert_eq!(
history.envelope(&schedule, audio).unwrap(),
Some(envelopes[1])
);
}
#[test]
fn bound_history_rejects_every_contract_perturbation_without_mutation() {
let exact = exact_payload_contract();
let schedule = schedule();
let text = coordinate(0, RealtimeFrameSlot::Text);
let mut history = RealtimePayloadHistory::with_contract(exact.clone());
history.insert(&schedule, text, 3).unwrap();
let other_schedule = RealtimeSpeechConfig::new(
2,
1,
1,
1,
0,
1,
RealtimeFrameConvention::AbsoluteDelayedSlots,
vec![2, 0, 3],
)
.unwrap();
let perturbations = [
(
payload_contract(
other_schedule,
2,
TokenDomain::new(32),
TokenDomain::new(16),
7,
3,
),
RealtimePayloadContractError::ScheduleMismatch,
),
(
payload_contract(
schedule.clone(),
3,
TokenDomain::new(32),
TokenDomain::new(16),
7,
3,
),
RealtimePayloadContractError::BatchMismatch,
),
(
payload_contract(
schedule.clone(),
2,
TokenDomain::new(33),
TokenDomain::new(16),
7,
3,
),
RealtimePayloadContractError::TextDomainMismatch,
),
(
payload_contract(
schedule.clone(),
2,
TokenDomain::new(32),
TokenDomain::new(17),
7,
3,
),
RealtimePayloadContractError::AudioDomainMismatch,
),
(
payload_contract(
schedule.clone(),
2,
TokenDomain::new(32),
TokenDomain::new(16),
8,
3,
),
RealtimePayloadContractError::GenerationMismatch,
),
(
payload_contract(
schedule.clone(),
2,
TokenDomain::new(32),
TokenDomain::new(16),
7,
4,
),
RealtimePayloadContractError::OwnerMismatch,
),
];
for (candidate, expected) in perturbations {
assert_eq!(
history.bind_or_validate_contract(&candidate),
Err(RealtimePayloadHistoryError::PayloadContract(expected))
);
assert_eq!(history.contract(), Some(&exact));
assert_eq!(history.required(&schedule, text), Ok(&3));
assert_eq!(history.len(), 1);
}
}
#[test]
fn history_clone_and_prune_preserve_envelope_identity() {
let schedule = schedule();
let contract = exact_payload_contract();
let mut history = RealtimePayloadHistory::with_contract(contract.clone());
for position in 0..=6 {
history
.insert(
&schedule,
coordinate(position, RealtimeFrameSlot::Text),
position,
)
.unwrap();
}
let mut branch = history.clone();
assert_eq!(branch.contract(), Some(&contract));
assert!(branch
.envelopes()
.all(|envelope| envelope.contract() == &contract));
assert_eq!(branch.prune_for_next_frontier(&schedule, 7), Ok(2));
assert_eq!(branch.len(), 5);
assert_eq!(history.len(), 7);
assert!(branch.envelopes().all(|envelope| {
envelope.coordinate().position() >= 2 && envelope.contract() == &contract
}));
}
#[test]
fn exact_insert_get_and_required_resolution_preserve_coordinate_order() {
let schedule = schedule();
let mut history = payload_history();
let text = coordinate(2, RealtimeFrameSlot::Text);
let audio_zero = coordinate(2, RealtimeFrameSlot::Audio(0));
let audio_one = coordinate(2, RealtimeFrameSlot::Audio(1));
history
.insert_many(
&schedule,
[(audio_one, "a1"), (text, "text"), (audio_zero, "a0")],
)
.unwrap();
assert_eq!(history.get(&schedule, text).unwrap(), Some(&"text"));
assert_eq!(history.required(&schedule, audio_one).unwrap(), &"a1");
assert_eq!(
history
.resolve_required(&schedule, [audio_one, text, audio_one])
.unwrap(),
vec![&"a1", &"text", &"a1"]
);
assert_eq!(
history.retained_values().copied().collect::<Vec<_>>(),
vec!["text", "a0", "a1"]
);
assert_eq!(history.schedule(), &schedule);
}
#[test]
fn invalid_or_duplicate_publications_are_atomic() {
let schedule = schedule();
let mut history = payload_history();
let retained = coordinate(0, RealtimeFrameSlot::Text);
history.insert(&schedule, retained, 7).unwrap();
let invalid = coordinate(1, RealtimeFrameSlot::Audio(2));
assert_eq!(
history.insert_many(
&schedule,
[
(coordinate(1, RealtimeFrameSlot::Audio(0)), 8),
(invalid, 9)
]
),
Err(RealtimePayloadHistoryError::InvalidSlot {
slot: RealtimeFrameSlot::Audio(2),
total_audio_codebooks: 2,
})
);
assert_eq!(history.len(), 1);
assert_eq!(
history.insert_many(&schedule, [(invalid, 8), (invalid, 9)]),
Err(RealtimePayloadHistoryError::InvalidSlot {
slot: RealtimeFrameSlot::Audio(2),
total_audio_codebooks: 2,
})
);
let pending = coordinate(1, RealtimeFrameSlot::Audio(0));
assert_eq!(
history.insert_many(&schedule, [(pending, 8), (pending, 9)]),
Err(RealtimePayloadHistoryError::DuplicatePayload {
coordinate: pending,
})
);
assert_eq!(history.len(), 1);
assert_eq!(
history.insert(&schedule, retained, 10),
Err(RealtimePayloadHistoryError::DuplicatePayload {
coordinate: retained,
})
);
assert_eq!(history.required(&schedule, retained), Ok(&7));
}
#[test]
fn mismatch_missing_and_arithmetic_fail_without_mutation() {
let schedule = schedule();
let other = RealtimeSpeechConfig::new(
2,
1,
1,
1,
0,
1,
RealtimeFrameConvention::AbsoluteDelayedSlots,
vec![2, 0, 3],
)
.unwrap();
let mut history = payload_history();
let retained = coordinate(0, RealtimeFrameSlot::Text);
history.insert(&schedule, retained, 7).unwrap();
assert_eq!(
history.insert(&other, coordinate(1, RealtimeFrameSlot::Text), 8),
Err(RealtimePayloadHistoryError::ScheduleMismatch)
);
let missing = coordinate(1, RealtimeFrameSlot::Audio(0));
assert_eq!(
history.required(&schedule, missing),
Err(RealtimePayloadHistoryError::MissingPayload {
coordinate: missing,
})
);
assert_eq!(
history.insert_delayed(&schedule, usize::MAX, RealtimeFrameSlot::Text, 9,),
Err(RealtimePayloadHistoryError::CoordinateOverflow {
base_position: usize::MAX,
delay: 2,
})
);
assert_eq!(history.len(), 1);
assert_eq!(history.required(&schedule, retained), Ok(&7));
}
#[test]
fn pruning_uses_next_frontier_and_maximum_delay_deterministically() {
let schedule = schedule();
let mut history = payload_history();
for position in 0..=5 {
history
.insert(
&schedule,
coordinate(position, RealtimeFrameSlot::Text),
position,
)
.unwrap();
}
assert_eq!(history.prune_for_next_frontier(&schedule, 3), Ok(0));
assert_eq!(history.prune_for_next_frontier(&schedule, 6), Ok(1));
assert_eq!(
history
.entries()
.map(|(coordinate, payload)| (coordinate.position(), *payload))
.collect::<Vec<_>>(),
vec![(1, 1), (2, 2), (3, 3), (4, 4), (5, 5)]
);
}
}