corevm-engine 0.1.28

CoreVM engine that drives program execution either on the builder or CoreVM service side
Documentation
use super::{AccumulateEngine, AccumulateOps, AccumulationError};
use alloc::{borrow::Cow, collections::VecDeque, vec::Vec};
use corevm_host::{
	CoreVmInstruction, MessageQueue, OutgoingServiceMessage, ServiceMessage, StorageKey,
};
use jam_types::{ServiceId, Slot};

/// The number of time-slots an outgoing message should be retained in the storage.
///
/// Normally two slots should be enough: one slot to initiate transfer by the service-sender and
/// another slot to process the transfer by service-receiver. We add two more slots to account for
/// potential blockchain reorganizations.
const MAX_OUTGOING_MESSAGE_AGE: Slot = 4;

impl<A: AccumulateOps> AccumulateEngine<A> {
	/// Removes processed incoming inter-service messages from the storage and updated
	/// message queue under [`StorageKey::IncomingServiceMessages`] key.
	pub(crate) fn remove_processed_incoming_service_messages(
		&mut self,
		processed_service_messages: &[ServiceMessage],
	) -> Result<(), AccumulationError<A::Error>> {
		use AccumulationError::*;
		let mut queue: VecDeque<ServiceMessage> =
			self.get_typed(&StorageKey::IncomingServiceMessages)?.unwrap_or_default();
		let old_len = queue.len();
		// Messages are processed in order.
		for actual in processed_service_messages.iter() {
			let expected = queue.pop_front();
			if Some(actual) != expected.as_ref() {
				return Err(MessageMismatch);
			}
			self.ops.remove(&StorageKey::IncomingMessage(*actual));
		}
		let new_len = queue.len();
		if new_len != old_len {
			self.set_typed(StorageKey::IncomingServiceMessages, &queue)?;
		}
		Ok(())
	}

	/// Removes all incoming inter-service messages and the corresponding message queue from the
	/// storage.
	pub(crate) fn remove_all_incoming_service_messages(
		&mut self,
	) -> Result<(), AccumulationError<A::Error>> {
		let queue: VecDeque<ServiceMessage> =
			self.get_typed(&StorageKey::IncomingServiceMessages)?.unwrap_or_default();
		for message in queue.into_iter() {
			self.ops.remove(&StorageKey::IncomingMessage(message));
		}
		self.ops.remove(&StorageKey::IncomingServiceMessages);
		Ok(())
	}

	/// Removes all outgoing inter-service messages and the corresponding message queue from the
	/// storage.
	pub(crate) fn remove_all_outgoing_service_messages(
		&mut self,
	) -> Result<(), AccumulationError<A::Error>> {
		let queue: MessageQueue<OutgoingServiceMessage> =
			self.get_typed(&StorageKey::OutgoingServiceMessages)?.unwrap_or_default();
		for message in queue.messages.into_iter() {
			self.ops.remove(&StorageKey::OutgoingMessage(message.inner.index));
		}
		self.ops.remove(&StorageKey::OutgoingServiceMessages);
		Ok(())
	}

	/// Push the incoming inter-service message into the queue under
	/// [`StorageKey::IncomingServiceMessages`] key and save it in the storage under
	/// [`StorageKey::IncomingMessage`] key.
	pub(crate) fn push_incoming_service_message(
		&mut self,
		message: ServiceMessage,
	) -> Result<(), AccumulationError<A::Error>> {
		use AccumulationError::*;
		log::debug!("Incoming inter-service message {message:?}");
		// Get outgoing message from the storage of the sender.
		let data = self
			.ops
			.get(message.source, &StorageKey::OutgoingMessage(message.index))
			.ok_or(InvalidMessage)?
			.into_owned();
		// Insert it as the incoming message in the storage of the receiver (this service).
		self.ops.set(StorageKey::IncomingMessage(message), data.into()).map_err(Api)?;
		let mut queue: VecDeque<ServiceMessage> =
			self.get_typed(&StorageKey::IncomingServiceMessages)?.unwrap_or_default();
		queue.push_back(message);
		self.set_typed(StorageKey::IncomingServiceMessages, &queue)?;
		Ok(())
	}

	/// Transfers outgoing inter-service messsage to the storage of the destination service via
	/// `transfer` host-call and removes outgoing messages older than two slots from the storage of
	/// the current service.
	pub(crate) fn send_new_outgoing_service_messages_and_remove_old(
		&mut self,
		outgoing_messages: Vec<(ServiceId, Vec<u8>)>,
	) -> Result<(), AccumulationError<A::Error>> {
		use AccumulationError::*;
		let source = self.service_id;
		let mut queue: MessageQueue<OutgoingServiceMessage> =
			self.get_typed(&StorageKey::OutgoingServiceMessages)?.unwrap_or_default();
		queue.messages.retain(|message| {
			let retain = self.slot < message.slot + MAX_OUTGOING_MESSAGE_AGE;
			if !retain {
				log::trace!("Removing {:?}", StorageKey::OutgoingMessage(message.inner.index));
				self.ops.remove(&StorageKey::OutgoingMessage(message.inner.index));
			}
			retain
		});
		for (destination, data) in outgoing_messages.into_iter() {
			let inner = ServiceMessage { index: queue.counter, source };
			let message = OutgoingServiceMessage { slot: self.slot, inner };
			log::debug!("Outgoing inter-service message {message:?}");
			queue.messages.push_back(message.clone());
			queue.counter += 1;
			self.ops
				.set(StorageKey::OutgoingMessage(message.inner.index), Cow::Owned(data))
				.map_err(Api)?;
			let min_memo_gas =
				self.ops.min_memo_gas(destination).ok_or(UnknownService(destination))?;
			let instruction = CoreVmInstruction::PushServiceMessage(inner);
			self.ops
				.transfer_typed(destination, 0, min_memo_gas, &instruction)
				.map_err(Api)?;
		}
		Ok(())
	}
}