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};
const MAX_OUTGOING_MESSAGE_AGE: Slot = 4;
impl<A: AccumulateOps> AccumulateEngine<A> {
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();
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(())
}
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(())
}
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(())
}
pub(crate) fn push_incoming_service_message(
&mut self,
message: ServiceMessage,
) -> Result<(), AccumulationError<A::Error>> {
use AccumulationError::*;
log::debug!("Incoming inter-service message {message:?}");
let data = self
.ops
.get(message.source, &StorageKey::OutgoingMessage(message.index))
.ok_or(InvalidMessage)?
.into_owned();
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(())
}
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(())
}
}