mod authority;
mod history;
pub use authority::{
AttemptAuthorityRecord, AttemptBudgetRecord, MAX_OPERATION_ATTEMPTS, OperationBindingRecord,
OperationBindingRequest,
};
use crate::model::artifacts::{ArtifactChecksumRecord, ChecksumError, canonical_hash};
use history::{AttemptEventRecord, Projection};
use serde::{Deserialize, Deserializer, Serialize, de};
use std::fmt;
use thiserror::Error;
pub const MAX_ATTEMPT_EVENTS: usize = 2048;
pub const MAX_ATTEMPT_JOURNAL_BYTES: u64 = 1024 * 1024;
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum MutationOutcomeRecord {
Applied,
NotApplied,
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum ObservationOutcomeRecord {
Applied,
NotApplied,
Uncertain,
}
#[derive(Clone, Debug)]
pub struct MutationReceiptRequest {
pub attempt: u32,
pub request: String,
pub outcome: MutationOutcomeRecord,
pub evidence: String,
}
#[derive(Clone, Debug)]
pub struct ObservationReceiptRequest {
pub attempt: u32,
pub request: String,
pub outcome: ObservationOutcomeRecord,
pub evidence: String,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
pub struct AttemptJournalView {
pub authority: ArtifactChecksumRecord,
pub mutations_used: u32,
pub observations_used: u32,
pub mutations_remaining: u32,
pub observations_remaining: u32,
pub pending_mutation: Option<u32>,
pub pending_observation: Option<u32>,
pub applied: bool,
pub history_len: usize,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(try_from = "JournalFields")]
pub struct AttemptJournalRecord {
version: u16,
authority: AttemptAuthorityRecord,
events: Vec<AttemptEventRecord>,
#[serde(skip)]
projection: Projection,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct JournalFields {
version: u16,
authority: AttemptAuthorityRecord,
#[serde(deserialize_with = "bounded_events")]
events: Vec<AttemptEventRecord>,
}
impl TryFrom<JournalFields> for AttemptJournalRecord {
type Error = AttemptJournalRecordError;
fn try_from(mut fields: JournalFields) -> Result<Self, Self::Error> {
if fields.version != 1 {
return Err(AttemptJournalRecordError::UnsupportedVersion(
fields.version,
));
}
let mut projection = Projection::empty();
for event in &mut fields.events {
event.normalize()?;
projection.apply(event, fields.authority.budget())?;
}
Ok(Self {
version: 1,
authority: fields.authority,
events: fields.events,
projection,
})
}
}
impl AttemptJournalRecord {
#[must_use]
pub const fn new(authority: AttemptAuthorityRecord) -> Self {
Self {
version: 1,
authority,
events: Vec::new(),
projection: Projection::empty(),
}
}
#[must_use]
pub const fn authority(&self) -> &AttemptAuthorityRecord {
&self.authority
}
#[must_use]
pub fn pending_observation_request(&self) -> Option<&str> {
self.projection
.pending_observation
.as_ref()
.map(|pending| pending.request.as_str())
}
#[must_use]
pub fn view(&self) -> AttemptJournalView {
AttemptJournalView {
authority: self.authority.digest(),
mutations_used: self.projection.mutations_used,
observations_used: self.projection.observations_used,
mutations_remaining: self.authority.budget().mutations()
- self.projection.mutations_used,
observations_remaining: self.authority.budget().observations()
- self.projection.observations_used,
pending_mutation: self.projection.pending_mutation,
pending_observation: self
.projection
.pending_observation
.as_ref()
.map(|pending| pending.attempt),
applied: self.projection.applied,
history_len: self.events.len(),
}
}
#[must_use]
pub fn digest(&self) -> ArtifactChecksumRecord {
let mut bytes = b"ic-backup/attempt-journal-history/v1\0".to_vec();
bytes.extend_from_slice(self.authority.digest().hash().as_bytes());
bytes.extend_from_slice(&(self.events.len() as u64).to_be_bytes());
for event in &self.events {
event.append_digest_bytes(&mut bytes);
}
ArtifactChecksumRecord::from_bytes(&bytes)
}
pub(crate) fn reserve_mutation(&mut self) -> Result<u32, AttemptJournalRecordError> {
let attempt = self.projection.next_attempt;
self.append(AttemptEventRecord::MutationReserved { attempt })?;
Ok(attempt)
}
pub(crate) fn reserve_observation(
&mut self,
mutation: u32,
request: &str,
) -> Result<u32, AttemptJournalRecordError> {
let attempt = self.projection.next_attempt;
self.append(AttemptEventRecord::ObservationReserved {
attempt,
mutation,
request: request.to_owned(),
})?;
Ok(attempt)
}
pub(crate) fn record_mutation(
&mut self,
receipt: MutationReceiptRequest,
) -> Result<(), AttemptJournalRecordError> {
if canonical_hash(&receipt.request)? != self.authority.binding().request() {
return Err(AttemptJournalRecordError::RequestMismatch);
}
self.append(AttemptEventRecord::MutationResolved {
mutation: receipt.attempt,
outcome: receipt.outcome,
evidence: receipt.evidence,
})
}
pub(crate) fn record_observation(
&mut self,
receipt: ObservationReceiptRequest,
) -> Result<(), AttemptJournalRecordError> {
self.append(AttemptEventRecord::ObservationRecorded {
observation: receipt.attempt,
request: receipt.request,
outcome: receipt.outcome,
evidence: receipt.evidence,
})
}
fn append(&mut self, mut event: AttemptEventRecord) -> Result<(), AttemptJournalRecordError> {
if self.events.len() == MAX_ATTEMPT_EVENTS {
return Err(AttemptJournalRecordError::HistoryTooLarge);
}
event.normalize()?;
let mut projection = self.projection.clone();
projection.apply(&event, self.authority.budget())?;
self.events.push(event);
self.projection = projection;
Ok(())
}
}
fn bounded_events<'de, D: Deserializer<'de>>(
deserializer: D,
) -> Result<Vec<AttemptEventRecord>, D::Error> {
struct EventsVisitor;
impl<'de> de::Visitor<'de> for EventsVisitor {
type Value = Vec<AttemptEventRecord>;
fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("a bounded chronological attempt event list")
}
fn visit_seq<A: de::SeqAccess<'de>>(
self,
mut sequence: A,
) -> Result<Self::Value, A::Error> {
let mut events = Vec::new();
while events.len() < MAX_ATTEMPT_EVENTS {
match sequence.next_element()? {
Some(event) => events.push(event),
None => return Ok(events),
}
}
if sequence.next_element::<de::IgnoredAny>()?.is_some() {
return Err(de::Error::custom(
AttemptJournalRecordError::HistoryTooLarge,
));
}
Ok(events)
}
}
deserializer.deserialize_seq(EventsVisitor)
}
#[derive(Debug, Error)]
pub enum AttemptJournalRecordError {
#[error("unsupported attempt journal version {0}")]
UnsupportedVersion(u16),
#[error("attempt allowance exceeds {MAX_OPERATION_ATTEMPTS}")]
BudgetTooLarge,
#[error("attempt history exceeds {MAX_ATTEMPT_EVENTS} events")]
HistoryTooLarge,
#[error("invalid attempt operation principal")]
InvalidPrincipal,
#[error(transparent)]
Checksum(#[from] ChecksumError),
#[error("mutation attempt {attempt} remains unresolved")]
MutationPending {
attempt: u32,
},
#[error("observation attempt {attempt} remains unresolved")]
ObservationPending {
attempt: u32,
},
#[error("no pending mutation attempt")]
NoPendingMutation,
#[error("no pending observation attempt")]
NoPendingObservation,
#[error("attempt identity mismatch: expected {expected}, actual {actual}")]
AttemptMismatch {
expected: u32,
actual: u32,
},
#[error("attempt request digest mismatch")]
RequestMismatch,
#[error("mutation attempt allowance exhausted")]
MutationBudgetExhausted,
#[error("observation attempt allowance exhausted")]
ObservationBudgetExhausted,
#[error("operation application already retained")]
AlreadyApplied,
}
#[cfg(test)]
mod tests;