use crate::model::{
artifacts::ArtifactChecksumRecord,
attempt_journal::{AttemptAuthorityRecord, AttemptJournalRecord},
ic_observation::{
IcObservationRequestError, IcObservationResponseError, IcObservationResponseInput,
ObservationReservation, fmt_response_input, validate_response_input,
},
ic_snapshot_data::{IcSnapshotDataRequest, MAX_IC_SNAPSHOT_DATA_REPLY_BYTES},
ic_snapshot_upload::{
IcSnapshotUploadAttemptError, IcSnapshotUploadKind, IcSnapshotUploadRequest,
original_authority,
},
operation_plan::OperationPlanRecord,
};
use std::fmt;
use thiserror::Error;
mod settlement;
pub use settlement::{IcSnapshotUploadDataAttribution, IcSnapshotUploadDataSettlement};
#[derive(Debug)]
pub struct IcSnapshotUploadDataObservationRequest<'request, 'source, 'metadata> {
plan: &'request OperationPlanRecord,
mutation: &'request IcSnapshotUploadRequest<'source>,
payload: &'request IcSnapshotDataRequest<'metadata>,
authority: AttemptAuthorityRecord,
mutation_attempt: u32,
observation_attempt: u32,
original_chunk_checksum: &'request ArtifactChecksumRecord,
}
impl<'request, 'source, 'metadata>
IcSnapshotUploadDataObservationRequest<'request, 'source, 'metadata>
{
pub fn new(
plan: &'request OperationPlanRecord,
operation_sequence: u64,
journal: &AttemptJournalRecord,
mutation: &'request IcSnapshotUploadRequest<'source>,
payload: &'request IcSnapshotDataRequest<'metadata>,
) -> Result<Self, IcSnapshotUploadDataObservationError> {
let IcSnapshotUploadKind::Data {
snapshot_id,
source_kind,
chunk_checksum,
..
} = mutation.kind()
else {
return Err(IcSnapshotUploadDataObservationError::DataUploadRequired);
};
let authority = original_authority(plan, operation_sequence, journal, mutation)?;
if payload.target() != mutation.target()
|| payload.snapshot_id() != snapshot_id
|| !payload.matches_kind(source_kind)
{
return Err(IcSnapshotUploadDataObservationError::ReadbackMismatch);
}
let actual = payload.metadata().metadata();
let original = mutation.source().metadata();
let actual_sizes = [
actual.wasm_module_size,
actual.wasm_memory_size,
actual.stable_memory_size,
];
let original_sizes = [
original.wasm_module_size,
original.wasm_memory_size,
original.stable_memory_size,
];
if !matches!(
actual.source,
Some(ic_management_canister_types::SnapshotSource::MetadataUpload(_))
) || actual_sizes != original_sizes
{
return Err(IcSnapshotUploadDataObservationError::DestinationMetadataMismatch);
}
let current = journal.view();
let request = Self {
plan,
mutation,
payload,
authority,
mutation_attempt: current
.pending_mutation
.ok_or(IcObservationRequestError::NoPendingMutation)?,
observation_attempt: current
.pending_observation
.ok_or(IcObservationRequestError::NoPendingObservation)?,
original_chunk_checksum: chunk_checksum,
};
request.validate_journal(journal)?;
Ok(request)
}
#[must_use]
pub const fn plan(&self) -> &OperationPlanRecord {
self.plan
}
#[must_use]
pub const fn mutation(&self) -> &'request IcSnapshotUploadRequest<'source> {
self.mutation
}
#[must_use]
pub const fn payload(&self) -> &'request IcSnapshotDataRequest<'metadata> {
self.payload
}
#[must_use]
pub const fn authority(&self) -> &AttemptAuthorityRecord {
&self.authority
}
#[must_use]
pub const fn mutation_attempt(&self) -> u32 {
self.mutation_attempt
}
#[must_use]
pub const fn observation_attempt(&self) -> u32 {
self.observation_attempt
}
#[must_use]
pub const fn original_chunk_checksum(&self) -> &ArtifactChecksumRecord {
self.original_chunk_checksum
}
pub fn validate_journal(
&self,
journal: &AttemptJournalRecord,
) -> Result<(), IcObservationRequestError> {
ObservationReservation {
authority: &self.authority,
mutation_attempt: self.mutation_attempt,
observation_attempt: self.observation_attempt,
request: self.payload.digest(),
}
.validate(journal)
}
}
#[derive(Debug, Error)]
pub enum IcSnapshotUploadDataObservationError {
#[error("data upload observation requires original data intent")]
DataUploadRequired,
#[error("data upload observation readback differs")]
ReadbackMismatch,
#[error("data upload observation destination metadata differs")]
DestinationMetadataMismatch,
#[error(transparent)]
Upload(#[from] IcSnapshotUploadAttemptError),
#[error(transparent)]
Observation(#[from] IcObservationRequestError),
}
#[derive(Clone)]
pub struct IcSnapshotUploadDataObservationResponse {
input: IcObservationResponseInput,
}
impl IcSnapshotUploadDataObservationResponse {
pub fn new(mut input: IcObservationResponseInput) -> Result<Self, IcObservationResponseError> {
validate_response_input(&mut input, MAX_IC_SNAPSHOT_DATA_REPLY_BYTES)?;
Ok(Self { input })
}
#[must_use]
pub const fn input(&self) -> &IcObservationResponseInput {
&self.input
}
}
impl fmt::Debug for IcSnapshotUploadDataObservationResponse {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt_response_input(
&self.input,
formatter,
"IcSnapshotUploadDataObservationResponse",
)
}
}
#[cfg(test)]
mod tests;