use crate::{
hash_encoded, hash_raw, AddressKind, Engine, HostCallError, InnerVm, InputStats, OuterVm,
OutputLimitReached, PrettyBytes, TooLargeWorkOutput, TouchError,
};
use alloc::{vec, vec::Vec};
use corevm_codec::video;
use corevm_host::{CoreVmOutput, Outcome, OutputStream, PageNum, PageSegmentOps, PAGE_SIZE};
use jam_pvm_common::ApiError;
use jam_types::SEGMENT_LEN;
use log::{debug, trace};
impl<O: OuterVm> Engine<O> {
pub(crate) fn map_pages(&mut self, size: u32) -> Result<Option<(u32, u32)>, HostCallError> {
let Some(size) = size.checked_next_multiple_of(PAGE_SIZE) else {
return Ok(None);
};
let num_pages = size / PAGE_SIZE;
Ok(self.work_output.map_pages(num_pages)?.map(|page_range| {
let address = page_range.start * PAGE_SIZE;
let size = (page_range.end - page_range.start) * PAGE_SIZE;
(address, size)
}))
}
pub(crate) fn unmap_pages(&mut self, address: u32, size: u32) -> Result<(), ApiError> {
let start_page = PageNum::from_address(address);
let end_page = address.saturating_add(size).div_ceil(PAGE_SIZE);
self.work_output.unmap_pages(start_page.0, end_page);
self.inner_vm.void(start_page, end_page - start_page.0)?;
trace!("Freed page(s) {:?}", self.memory_map.pretty_page_range(start_page.0..end_page));
Ok(())
}
pub(crate) fn touch_once_range(&mut self, address: u32, length: u32) -> Result<(), TouchError> {
let pages = address_and_length_to_page_range(address, length)
.ok_or(TouchError::InvalidMemoryAccess)?;
for page in pages {
let page = PageNum(page);
if self.work_output.updated_pages.contains_key(&page) {
continue;
}
self.touch(page)?;
}
Ok(())
}
pub(crate) fn touch(&mut self, page: PageNum) -> Result<(), TouchError> {
let address = page.address();
let Some(class) = self.memory_map.classify_address(address) else {
trace!("Page fault at {:?}", self.memory_map.pretty_address(address));
return Err(TouchError::InvalidMemoryAccess);
};
if class == AddressKind::RoData {
trace!("Page fault at {:?}", self.memory_map.pretty_address(address));
self.inner_vm.poke_ro_rw_data_page(
&self.ro_data[..],
self.memory_map.ro_data_range(),
address,
)?;
return Ok(());
}
if self.num_exported_pages().saturating_add(self.output.segment_count()) >=
usize::from(self.export_count)
{
trace!(
"Page fault at {:?}: The page will not fit into work package exports",
self.memory_map.pretty_address(address)
);
return Err(TouchError::PageFault { page, num_pages: 1 });
}
let mut imported_page_hash = None;
if self.work_output.vm_state.resident_pages.contains_index(page.0) {
let Some(segment) = self.input.import_page(page)? else {
trace!(
"Page fault at {:?} (missing page)",
self.memory_map.pretty_address(address)
);
return Err(TouchError::PageFault { page, num_pages: 1 });
};
trace!("Page fault at {:?} (resident page)", self.memory_map.pretty_address(address));
self.inner_vm.zero_poke(segment.as_array().page(), address)?;
imported_page_hash = Some(hash_raw(segment.as_ref()));
} else {
match class {
AddressKind::RoData => {
unreachable!()
},
AddressKind::RwData => {
trace!("Page fault at {:?}", self.memory_map.pretty_address(address));
self.inner_vm.poke_ro_rw_data_page(
&self.rw_data[..],
self.memory_map.rw_data_range(),
address,
)?;
},
AddressKind::Stack => {
trace!("Page fault at {:?}", self.memory_map.pretty_address(address));
self.inner_vm.zero(page, 1)?;
},
AddressKind::Heap => {
if !self.work_output.heap_page_mapper().is_mapped(address / PAGE_SIZE) {
trace!(
"Page fault at {:?} (unmapped)",
self.memory_map.pretty_address(address)
);
return Err(TouchError::InvalidMemoryAccess);
}
trace!("Page fault at {:?}", self.memory_map.pretty_address(address));
self.inner_vm.zero(page, 1)?;
},
}
}
self.work_output.insert_updated_page(page, imported_page_hash).map_err(
|TooLargeWorkOutput| {
trace!(
"Page fault at {:?}: The page will not fit into work output",
self.memory_map.pretty_address(address)
);
TouchError::PageFault { page, num_pages: 1 }
},
)?;
Ok(())
}
fn check_exports_after(&self, additional_exports: usize) -> Result<(), usize> {
let num_exports = self.num_exported_pages() + self.output.segment_count();
if num_exports.saturating_add(additional_exports) > usize::from(self.export_count) {
return Err(num_exports);
}
Ok(())
}
fn num_exported_pages(&self) -> usize {
self.work_output
.updated_pages
.keys()
.filter(|page| {
self.work_output.vm_state.resident_pages.contains_index(page.0)
})
.count()
}
fn check_can_append(&self, stream: OutputStream, len: usize) -> Result<(), HostCallError> {
match stream {
OutputStream::Video =>
if self.work_output.vm_state.video.is_none() {
trace!("Video output mode is not set");
return Err(Outcome::Panic.into());
},
OutputStream::Audio =>
if self.work_output.vm_state.audio.is_none() {
trace!("Audio output mode is not set");
return Err(Outcome::Panic.into());
},
OutputStream::Stdout | OutputStream::Stderr => {},
}
let additional_segments = self
.output
.new_segments_after(self.timestamp, stream, len)
.ok_or(OutputLimitReached)?;
Ok(self.check_exports_after(additional_segments).map_err(|_| OutputLimitReached)?)
}
pub(crate) fn append_output(
&mut self,
stream: OutputStream,
inner_src: u32,
len: u32,
) -> Result<(), HostCallError> {
if len == 0 {
return Ok(());
}
self.touch_once_range(inner_src, len)?;
match stream {
OutputStream::Video => {
let Some(mode) = self.work_output.vm_state.video else {
return Err(Outcome::Panic.into());
};
let mut video_stream = Vec::new();
let encoder = match self.video_encoder {
Some(ref mut encoder) => encoder,
ref mut ret_encoder @ None => {
let mut config = video::Config::default();
config.quantization_level = mode.options.quantization_level();
config.chroma_subsampling = mode.options.chroma_subsampling();
config.raw = mode.options.raw();
let mut encoder = video::Encoder::new(mode.width, mode.height, config);
encoder.start(&mut video_stream);
ret_encoder.insert(encoder)
},
};
let expected_frame_len =
u64::from(mode.width.get()) * u64::from(mode.height.get()) + 3 * 256;
if u64::from(len) != expected_frame_len {
return Err(Outcome::Panic.into());
}
let mut buf = vec![0_u8; len as usize];
self.inner_vm.peek_into(&mut buf, inner_src)?;
encoder.write_rgb888_indexed8_frame(&buf[..], &mut video_stream);
self.check_can_append(stream, video_stream.len())?;
self.output.append(self.timestamp, stream, &video_stream);
trace!(
"Appended video frame ({} B) from {:?}",
video_stream.len(),
self.memory_map.pretty_address(inner_src)
);
return Ok(());
},
OutputStream::Audio =>
if self.work_output.vm_state.audio.is_none() {
return Err(Outcome::Panic.into());
},
OutputStream::Stdout | OutputStream::Stderr => {},
}
self.check_can_append(stream, len as usize)?;
let buf = self.output.pre_allocate(self.timestamp, stream, len as usize);
self.inner_vm.peek_into(buf, inner_src)?;
trace!("Appended {len} byte(s) to {stream:?} from buffer {:?}", PrettyBytes(buf));
Ok(())
}
pub(crate) fn append_output_slice(
&mut self,
stream: OutputStream,
slice: &[u8],
) -> Result<(), HostCallError> {
if slice.is_empty() {
return Ok(());
}
let len = slice.len();
self.check_can_append(stream, len).inspect_err(|_| {
trace!("Output limit reached while appending {len} byte(s) to {stream:?}")
})?;
self.output.append(self.timestamp, stream, slice);
trace!("Appended {len} bytes to {stream:?} from buffer {:?}", PrettyBytes(slice));
Ok(())
}
pub(crate) fn suspend(
mut self,
outcome: Outcome,
) -> Result<(CoreVmOutput, O, InputStats), ApiError> {
let stream_len = self.output.stream_len();
let output_len = self.output.total_len() as u32;
let mut work_output = self.work_output.into_inner();
let mut num_exported_pages = 0_u32;
let mut updated_pages = core::mem::take(&mut work_output.updated_pages).into_vec();
for (page, hash) in updated_pages.iter_mut() {
let address = page.address();
if work_output.vm_state.resident_pages.contains_index(page.0) {
let mut segment = [0_u8; SEGMENT_LEN];
segment.set_page_number(*page);
self.inner_vm.peek_into(segment.page_mut(), address)?;
self.outer_vm.export(&segment)?;
num_exported_pages += 1;
*hash = hash_raw(&segment);
} else {
trace!("Not exporting freed page {:?}", self.memory_map.pretty_page(page.0));
}
}
work_output.updated_pages = updated_pages.into();
self.output.export_segments(|segment| self.outer_vm.export(segment))?;
let num_output_segments = output_len.div_ceil(SEGMENT_LEN as u32);
let total_segment_count = num_exported_pages + num_output_segments;
for _ in total_segment_count..u32::from(self.export_count) {
let segment = [0_u8; SEGMENT_LEN];
self.outer_vm.export(&segment)?;
}
debug!(
"Exported page(s): {:?}",
self.memory_map
.pretty_page_ranges(work_output.updated_pages.iter().map(|(page, _)| page.0))
);
debug!(
"Exported {} segments: {} memory page(s), {} output segment(s), {} null segment(s)",
self.export_count,
num_exported_pages,
output_len.div_ceil(SEGMENT_LEN as u32),
u32::from(self.export_count).saturating_sub(total_segment_count)
);
let program_counter = self.inner_vm.expunge().expect("Failed to close VM handle");
work_output.vm_state.regs = self.args.regs;
work_output.vm_state.program_counter = program_counter;
work_output.vm_state.step = work_output.vm_state.step.saturating_add(1);
let new_hash = hash_encoded(&work_output.vm_state);
work_output.vm_output.remaining_gas = self.args.gas;
work_output.vm_output.outcome = outcome;
work_output.vm_output.num_memory_pages = num_exported_pages;
work_output.vm_output.stream_len = stream_len;
work_output.new_hash = new_hash;
Ok((work_output, self.outer_vm, self.input.into_stats()))
}
}
fn address_and_length_to_page_range(address: u32, length: u32) -> Option<core::ops::Range<u32>> {
let first_page = address / PAGE_SIZE;
let last = address.checked_add(length.saturating_sub(1))?;
let last_page = last / PAGE_SIZE;
Some(first_page..last_page + 1)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
testing::default_corevm_output, BoundedInput, BoundedWorkOutput, FatalError, InvokeArgs,
MemoryMap,
};
use alloc::{vec, vec::Vec};
use bytes::Bytes;
use corevm_host::{fs, ExecEnv, Range};
use jam_pvm_common::InvokeOutcome;
use jam_types::{max_exports, SegmentBytes, ServiceId, SignedGas, VecSet};
use polkavm::MemoryMapBuilder;
use rand::{seq::IndexedRandom, Rng, RngCore};
#[test]
fn touch_and_append_output_works() {
let _ = env_logger::try_init();
let mut rng = rand::rng();
let mut engine = Engine::with_data(
{
let len = rng.random_range(0..=2 * PAGE_SIZE);
let mut data = vec![0_u8; len as usize];
rng.fill_bytes(&mut data[..]);
data
},
{
let len = rng.random_range(0..=2 * PAGE_SIZE);
let mut data = vec![0_u8; len as usize];
rng.fill_bytes(&mut data[..]);
data
},
);
let ro_page_range = {
let r = engine.memory_map.ro_data_range();
let ro_start_page = r.start / PAGE_SIZE;
let ro_end_page = r.end / PAGE_SIZE;
ro_start_page..ro_end_page
};
let ro_data_page_range = {
let r = engine.memory_map.ro_data_range();
let data_len = engine.ro_data.len() as u32;
let ro_start_page = r.start / PAGE_SIZE;
let ro_end_page = (r.start + data_len).div_ceil(PAGE_SIZE);
ro_start_page..ro_end_page
};
let rw_page_range = {
let r = engine.memory_map.rw_data_range();
let rw_start_page = r.start / PAGE_SIZE;
let rw_end_page = r.end / PAGE_SIZE;
rw_start_page..rw_end_page
};
let rw_data_page_range = {
let r = engine.memory_map.rw_data_range();
let data_len = engine.rw_data.len() as u32;
let rw_start_page = r.start / PAGE_SIZE;
let rw_end_page = (r.start + data_len).div_ceil(PAGE_SIZE);
rw_start_page..rw_end_page
};
let stack_page_range = {
let r = engine.memory_map.stack_range();
let stack_start_page = r.start / PAGE_SIZE;
let stack_end_page = r.end / PAGE_SIZE;
stack_start_page..stack_end_page
};
let heap_page_range = {
let r = engine.memory_map.heap_range();
let heap_start_page = r.start / PAGE_SIZE;
let heap_end_page = r.end / PAGE_SIZE;
heap_start_page..heap_end_page
};
assert_eq!(None, engine.memory_map.classify_address(0));
let touch_result = engine.touch(PageNum(0));
assert!(
matches!(touch_result, Err(TouchError::InvalidMemoryAccess)),
"Result = {touch_result:?}"
);
macro_rules! check_touch {
($range: expr, readonly = $read_only: expr, alloc = $alloc: expr, data = $data: expr$(,)*) => {{
let range = $range;
let data = $data;
let (page, address) = if $alloc {
let (address, _size) = engine.map_pages(PAGE_SIZE).unwrap().unwrap();
let page = PageNum::from_address(address);
(page, address)
} else {
let page = PageNum(rng.random_range(range.clone()));
let address = page.address();
(page, address)
};
engine.work_output.reset_updated_pages();
engine.inner_vm.calls.clear();
assert!(engine.touch(page).is_ok());
if data.is_empty() {
assert_eq!(
[Call::Zero { page, num_pages: 1 }].as_slice(),
engine.inner_vm.calls.as_slice(),
"Wrong host-calls for empty data section"
);
} else {
let start = (page.0 - range.start) as usize * PAGE_SIZE as usize;
let end = (start + PAGE_SIZE as usize).min(data.len());
let outer_src = data[start..end].to_vec();
assert_eq!(
[
Call::Zero { page, num_pages: 1 },
Call::Poke { outer_src, inner_dst: address }
]
.as_slice(),
engine.inner_vm.calls.as_slice(),
"Wrong host-calls for non-empty data section"
);
}
if !$read_only {
assert_eq!(
[Range::new(page.0, page.0 + 1)].as_slice(),
engine.work_output.vm_state.resident_pages.as_slice(),
"Wrong resident pages"
);
assert_eq!(
VecSet::from_iter([page]),
engine
.work_output
.updated_pages
.keys()
.copied()
.collect::<VecSet<PageNum>>(),
"Wrong touched pages"
);
}
}};
}
for _ in 0..1000 {
check_touch!(
ro_data_page_range.end..ro_page_range.end,
readonly = true,
alloc = false,
data = [0_u8; 0],
);
check_touch!(
ro_data_page_range.start..ro_data_page_range.end,
readonly = true,
alloc = false,
data = engine.ro_data.clone(),
);
check_touch!(
rw_data_page_range.end..rw_page_range.end,
readonly = false,
alloc = false,
data = [0_u8; 0],
);
check_touch!(
rw_data_page_range.start..rw_data_page_range.end,
readonly = false,
alloc = false,
data = engine.rw_data.clone(),
);
check_touch!(
stack_page_range.start..stack_page_range.end,
readonly = false,
alloc = false,
data = [0_u8; 0],
);
check_touch!(
heap_page_range.start..heap_page_range.end,
readonly = false,
alloc = true,
data = [0_u8; 0],
);
{
engine.work_output.reset_updated_pages();
engine.inner_vm.calls.clear();
let stream = *BYTE_ORIENTED_OUTPUT_STREAMS.choose(&mut rng).unwrap();
let len = rng.random_range(1..=PAGE_SIZE);
let (address, _size) = engine.map_pages(len).unwrap().unwrap();
let page_aligned_address = address - address % PAGE_SIZE;
let start_page = page_aligned_address / PAGE_SIZE;
let end_page = (address + len).div_ceil(PAGE_SIZE);
let result = engine.append_output(stream, address, len);
assert!(result.is_ok(), "Result = {result:?}");
let mut calls = Vec::new();
if len != 0 {
for page in start_page..end_page {
calls.push(Call::Zero { page: PageNum(page), num_pages: 1 });
}
calls.push(Call::PeekInto { inner_src: address, len });
}
assert_eq!(calls.as_slice(), engine.inner_vm.calls.as_slice());
}
}
}
#[test]
fn page_limit_on_touch_works() {
let mut engine = Engine::new_for_tests();
let mut num_touched = 0;
for _ in 0..max_exports() {
let (address, _size) = engine.map_pages(PAGE_SIZE).unwrap().unwrap();
let page = PageNum::from_address(address);
let result = engine.touch(page);
assert!(engine.work_output.updated_pages.len() <= max_exports() as usize);
assert_eq!(Ok(()), engine.work_output.check_len());
if result.is_err() {
break;
}
num_touched += 1;
}
assert_ne!(0, num_touched);
}
const BYTE_ORIENTED_OUTPUT_STREAMS: [OutputStream; 2] =
[OutputStream::Stdout, OutputStream::Stderr];
#[test]
fn page_limit_on_output_works() {
let _ = env_logger::try_init();
let mut engine = Engine::new_for_tests();
engine.inner_vm.randomize_memory = true;
assert!(
engine.memory_map.heap_range().end - engine.memory_map.heap_range().start >
max_exports()
);
let mut rng = rand::rng();
let (address, _len) = engine.map_pages(PAGE_SIZE).unwrap().unwrap();
loop {
let stream = *BYTE_ORIENTED_OUTPUT_STREAMS.choose(&mut rng).unwrap();
let len = rng.random_range(0..=PAGE_SIZE);
let result = engine.append_output(stream, address, len);
assert_eq!(Ok(()), engine.work_output.check_len());
match result {
Ok(()) => {},
Err(HostCallError::Outcome(Outcome::OutputLimitReached)) => {
break;
},
Err(e) => {
panic!("Unexpected error {e:?}");
},
}
}
for _ in 0..10 {
let stream = *BYTE_ORIENTED_OUTPUT_STREAMS.choose(&mut rng).unwrap();
assert!(engine.append_output(stream, 0, 0).is_ok());
}
assert!(engine.work_output.updated_pages.len() <= max_exports() as usize);
assert_eq!(Ok(()), engine.work_output.check_len());
let (output, outer_vm, ..) = engine.suspend(Outcome::Halt).unwrap();
assert_eq!(
(output.vm_output.num_memory_pages + output.vm_output.num_output_segments()) as usize,
outer_vm.real_export_count() as usize
);
assert!(outer_vm.exports.len() <= max_exports() as usize);
}
#[test]
fn append_output_edge_cases() {
let mut engine = Engine::new_for_tests();
let mut rng = rand::rng();
let stream = *BYTE_ORIENTED_OUTPUT_STREAMS.choose(&mut rng).unwrap();
let result = engine.append_output(stream, engine.memory_map.heap_range().start, u32::MAX);
assert!(
matches!(result, Err(HostCallError::Fatal(FatalError::InvalidMemoryAccess))),
"Result = {result:?}"
);
for kind in AddressKind::ALL {
let address = engine.memory_map.address_range(kind).end;
if engine.memory_map.classify_address(address).is_some() {
continue;
}
let result = engine.append_output(stream, address, 1);
assert!(
matches!(result, Err(HostCallError::Fatal(FatalError::InvalidMemoryAccess))),
"Result = {result:?}, address kind = {kind:?}",
);
}
}
impl Engine<TestOuterVm> {
fn new_for_tests() -> Self {
Self::with_data(Default::default(), Default::default())
}
fn with_data(ro_data: Vec<u8>, rw_data: Vec<u8>) -> Self {
let inner_vm = TestInnerVm::default();
let outer_vm = TestOuterVm::default();
let memory_map = MemoryMapBuilder::new(PAGE_SIZE)
.ro_data_size(PAGE_SIZE * 10)
.rw_data_size(PAGE_SIZE * 10)
.stack_size(PAGE_SIZE * 10)
.build()
.expect("Failed to build memory map");
let memory_map = MemoryMap::new(&memory_map);
let heap_address_range = memory_map.heap_range();
let heap_page_range = Range::new(
heap_address_range.start / PAGE_SIZE,
heap_address_range.end / PAGE_SIZE,
);
Self {
work_output: BoundedWorkOutput::try_from(
default_corevm_output(),
Default::default(),
0,
heap_page_range,
)
.unwrap(),
ro_data: ro_data.into(),
rw_data: rw_data.into(),
output: Default::default(),
args: InvokeArgs::new(0),
input: BoundedInput::new(0, Default::default()),
host_call_handlers: Default::default(),
host_call_names: Default::default(),
num_video_frames: 0,
num_audio_bytes: 0,
timestamp: 0,
outer_vm,
inner_vm,
export_count: max_exports() as u16,
memory_map,
exec: ExecEnv {
program: fs::BlockRef { service_id: ServiceId::MAX, hash: fs::Hash::default() },
root_dir: fs::BlockRef {
service_id: ServiceId::MAX,
hash: fs::Hash::default(),
},
args: Default::default(),
env: Default::default(),
audio_input: None,
video_input: None,
input_key: None,
},
video_encoder: None,
}
}
}
#[derive(Default)]
struct TestInnerVm {
calls: Vec<Call>,
randomize_memory: bool,
}
impl InnerVm for TestInnerVm {
fn void(&mut self, page: PageNum, num_pages: u32) -> Result<(), ApiError> {
self.calls.push(Call::Void { page, num_pages });
Ok(())
}
fn zero(&mut self, page: PageNum, num_pages: u32) -> Result<(), ApiError> {
self.calls.push(Call::Zero { page, num_pages });
Ok(())
}
fn poke(&mut self, outer_src: &[u8], inner_dst: u32) -> Result<(), ApiError> {
self.calls.push(Call::Poke { outer_src: outer_src.to_vec(), inner_dst });
Ok(())
}
fn peek_into(&mut self, outer_dst: &mut [u8], inner_src: u32) -> Result<(), ApiError> {
self.calls.push(Call::PeekInto { len: outer_dst.len() as u32, inner_src });
if self.randomize_memory {
rand::rng().fill_bytes(outer_dst);
}
Ok(())
}
fn expunge(self) -> Result<u64, ApiError> {
Ok(0)
}
fn invoke(
&mut self,
_gas: SignedGas,
_regs: [u64; 13],
) -> Result<(InvokeOutcome, SignedGas, [u64; 13]), ApiError> {
unreachable!()
}
}
#[derive(Debug, PartialEq, Eq)]
enum Call {
Void { page: PageNum, num_pages: u32 },
Zero { page: PageNum, num_pages: u32 },
Poke { outer_src: Vec<u8>, inner_dst: u32 },
PeekInto { len: u32, inner_src: u32 },
}
#[derive(Default)]
struct TestOuterVm {
exports: Vec<SegmentBytes>,
}
impl TestOuterVm {
fn real_export_count(&self) -> u32 {
let num_null_segments = self
.exports
.iter()
.rev()
.take_while(|segment| segment.as_slice().iter().all(|&b| b == 0))
.count();
(self.exports.len() - num_null_segments) as u32
}
}
impl OuterVm for TestOuterVm {
type InnerVm = TestInnerVm;
fn export(&mut self, segment: &[u8; SEGMENT_LEN]) -> Result<(), ApiError> {
self.exports.push(segment.into());
Ok(())
}
fn read_file_block(&mut self, _block_ref: &fs::BlockRef) -> Option<Bytes> {
unreachable!()
}
fn machine(
&mut self,
_code: &[u8],
_program_counter: u64,
) -> Result<TestInnerVm, ApiError> {
unreachable!()
}
}
}