Coverage Report

Created: 2026-09-11 07:47

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/home/runner/work/lyquor/lyquor/platform/vm/src/mem.rs
Line
Count
Source
1
use std::marker::PhantomData;
2
use std::sync::{Arc, atomic};
3
4
use crate::guest::VmGuest;
5
use crate::{RwLock as SyncRwLock, RwLockWriteGuard as SyncRwLockWriteGuard};
6
use lyquid::mem::Wasm32;
7
use lyquor_api::store::LazyBytes;
8
use lyquor_primitives::StateCategory;
9
use lyquor_state::{Key, State, StateR, Value};
10
use tokio::sync::{RwLock, RwLockWriteGuard};
11
use userspace_pagefault::{
12
    self, AccessType, AsyncPageFaultFuture, AsyncPageStore, PageChunks, PagedSegment, SharedMemory,
13
};
14
15
#[cfg(test)] use lyquor_test::test;
16
17
const LYTEPAGE_SIZE: usize = 1024; // the logical page size for LyteMemory
18
const BEGIN_GUARD_SIZE: usize = 2 << 30;
19
const END_GUARD_SIZE: usize = 2 << 30;
20
21
60
fn lytememory_setup() {
22
    static INITIALIZED: std::sync::Once = std::sync::Once::new();
23
60
    INITIALIZED.call_once(|| 
{31
24
        // this is a non-so-bad hack to wasmtime: creating a dummy wasmtime engine will trigger the
25
        // on-time initialization of the trap handler by wasmtime. Then we can override this trap
26
        // handler and only hand back when the signals are not triggered within any managed
27
        // LyteMemory.
28
31
        let _ = wasmtime::Engine::new(wasmtime::Config::new().macos_use_mach_ports(false));
29
31
    });
30
60
}
31
32
struct Bitmask(Box<[atomic::AtomicU64]>);
33
34
impl Bitmask {
35
339
    fn new(size: usize) -> Self {
36
        Self(
37
1.37M
            
std::iter::repeat_with339
(|| atomic::AtomicU64::new(0))
38
339
                .take((size + 63) >> 6)
39
339
                .collect::<Vec<_>>()
40
339
                .into(),
41
        )
42
339
    }
43
44
    #[inline(always)]
45
747
    fn mark(&self, idx: usize) {
46
747
        self.0[idx >> 6].fetch_or(1 << (idx & 63), atomic::Ordering::Relaxed);
47
747
    }
48
49
    #[inline(always)]
50
494
    fn is_marked(&self, idx: usize) -> bool {
51
494
        (self.0[idx >> 6].fetch_or(0, atomic::Ordering::Relaxed) >> (idx & 63)) & 1 == 1
52
494
    }
53
54
    const TABLE: [u64; 64] = [
55
        0x00, 0x3a, 0x01, 0x3b, 0x2f, 0x35, 0x02, 0x3c, 0x27, 0x30, 0x1b, 0x36, 0x21, 0x2a, 0x03, 0x3d, 0x33, 0x25,
56
        0x28, 0x31, 0x12, 0x1c, 0x14, 0x37, 0x1e, 0x22, 0x0b, 0x2b, 0x0e, 0x16, 0x04, 0x3e, 0x39, 0x2e, 0x34, 0x26,
57
        0x1a, 0x20, 0x29, 0x32, 0x24, 0x11, 0x13, 0x1d, 0x0a, 0x0d, 0x15, 0x38, 0x2d, 0x19, 0x1f, 0x23, 0x10, 0x09,
58
        0x0c, 0x2c, 0x18, 0x0f, 0x08, 0x17, 0x07, 0x06, 0x05, 0x3f,
59
    ];
60
61
    /// Turn a power of 2 (>0) to the exponent.
62
    #[inline(always)]
63
383
    fn fast_log(mut x: u64) -> u64 {
64
383
        x |= x >> 1;
65
383
        x |= x >> 2;
66
383
        x |= x >> 4;
67
383
        x |= x >> 8;
68
383
        x |= x >> 16;
69
383
        x |= x >> 32;
70
383
        Self::TABLE[((x.wrapping_mul(0x03f6eaf2cd271461)) >> 58) as usize]
71
383
    }
72
73
    #[inline]
74
157
    fn get_marked_and_clear(&self) -> impl Iterator<Item = usize> + '_ {
75
630k
        
self.0157
.
iter157
().
zip157
((
0..157
).
step_by157
(64)).
flat_map157
(|(bits, i)| {
76
630k
            let mut v = bits.swap(0, atomic::Ordering::Relaxed);
77
631k
            
std::iter::from_fn630k
(move || {
78
                // this lowbit finding algorithm is suitable for a sparse bitmask, which is the
79
                // case typically for LyteMemory dirty pages
80
631k
                let lowbit = v;
81
631k
                if lowbit == 0 {
82
630k
                    None
83
                } else {
84
383
                    let lowbit = lowbit.isolate_lowest_one();
85
383
                    v ^= lowbit; // clear the lowest bit
86
383
                    Some(i + Self::fast_log(lowbit) as usize)
87
                }
88
631k
            })
89
630k
        })
90
157
    }
91
}
92
93
#[test]
94
fn test_bitmask() {
95
    let all_bits: Vec<_> = (0..64).collect();
96
    let cases = [
97
        (64, &all_bits[..]),
98
        (64, &[0, 2, 4, 7, 11, 33, 63][..]),
99
        (1234, &[0, 32, 63, 64, 100, 256, 1024, 1233][..]),
100
    ];
101
    for (size, pos) in cases {
102
        let bm = Bitmask::new(size);
103
        for i in pos {
104
            bm.mark(*i);
105
        }
106
        for (i, j) in bm.get_marked_and_clear().zip(pos.iter()) {
107
            assert_eq!(i, *j);
108
        }
109
    }
110
}
111
112
2.34k
fn litepage_state_key(litepage_id: u32) -> Key {
113
    static PREFIX: std::sync::OnceLock<Key> = std::sync::OnceLock::new();
114
2.34k
    let key = PREFIX.get_or_init(|| lyquid::LYTEMEM_PAGE_PREFIX.
into29
());
115
    // using big-endian here so the state debug will be easier to read
116
2.34k
    let id: LazyBytes = litepage_id.to_be_bytes().into();
117
2.34k
    id.prepend(key)
118
2.34k
}
119
120
struct Segment {
121
    /// Dirty page bitmask. One position (bit) corresponds to one OS page.
122
    dirty: Bitmask,
123
    /// Loaded page bitmask. One position (bit) corresponds to one OS page.
124
    loaded: Bitmask,
125
    /// Size of an OS page on the host.
126
    page_size: usize,
127
    /// Offset of this segment from the beginning of the LytePage area.
128
    litepage_base: usize,
129
    /// Offest of this segment in the memory (PagedSegment).
130
    base: usize,
131
    /// Length of the segment.
132
    length: usize,
133
}
134
135
impl Segment {
136
168
    fn new(base: usize, length: usize, litepage_base: usize, page_size: usize) -> Self {
137
        // align to the closest page
138
168
        let bitmask_size = length.div_ceil(page_size);
139
168
        Self {
140
168
            dirty: Bitmask::new(bitmask_size),
141
168
            loaded: Bitmask::new(bitmask_size),
142
168
            page_size,
143
168
            litepage_base,
144
168
            base,
145
168
            length,
146
168
        }
147
168
    }
148
149
    #[inline(always)]
150
1.16k
    fn addr_to_idx(&self, mem_addr: usize) -> usize {
151
1.16k
        (mem_addr - self.base) / self.page_size
152
1.16k
    }
153
154
    #[inline(always)]
155
386
    fn mark_dirty(&self, mem_addr: usize) {
156
386
        self.dirty.mark(self.addr_to_idx(mem_addr));
157
386
    }
158
159
    #[inline(always)]
160
494
    fn is_loaded(&self, mem_addr: usize) -> bool {
161
494
        self.loaded.is_marked(self.addr_to_idx(mem_addr))
162
494
    }
163
164
    #[inline(always)]
165
282
    fn mark_loaded(&self, mem_addr: usize) {
166
282
        self.loaded.mark(self.addr_to_idx(mem_addr));
167
282
    }
168
169
    /// Write the dirty pages to the key-value state and clear the dirty bitmask.
170
152
    async fn writeback_changes<S: State>(&self, state: &SyncRwLock<S>, mem: &[u8]) {
171
152
        let pending = {
172
152
            let state = state.read();
173
152
            let mut pending = Vec::new();
174
302
            for ospage_id in 
self.dirty152
.
get_marked_and_clear152
() {
175
302
                let segment_off = ospage_id * self.page_size;
176
302
                let litepage_start = (self.litepage_base + segment_off) / LYTEPAGE_SIZE;
177
302
                let mem_addr = self.base + segment_off;
178
302
                let ospage = &mem[mem_addr..mem_addr + self.page_size];
179
1.20k
                for (i, litepage) in 
ospage302
.
chunks302
(LYTEPAGE_SIZE).
enumerate302
() {
180
1.20k
                    let key = litepage_state_key((litepage_start + i) as u32);
181
1.20k
                    let litepage: Value = litepage.to_vec().into();
182
1.20k
                    pending.push((key.clone(), litepage, state.get(key)));
183
1.20k
                }
184
            }
185
152
            pending
186
        };
187
188
152
        let mut changes = Vec::new();
189
1.20k
        for (key, litepage, old) in 
pending152
{
190
1.20k
            let changed = match old.await {
191
71
                Some(old) => litepage.as_ref() != old.as_ref(), // if the page was previous persisted, compare
192
1.04M
                None => !
litepage.as_ref().iter()1.13k
.
all1.13k
(|&b| b == 0), // otherwise check if it's empty
193
            };
194
1.20k
            if changed {
195
294
                changes.push((key, litepage));
196
914
            }
197
        }
198
199
152
        let mut state = state.write();
200
294
        for (key, litepage) in 
changes152
{
201
294
            state.set(key, Some(litepage));
202
294
        }
203
152
    }
204
205
    /// Revert the dirty pages to their original data and clear the dirty bitmask.
206
2
    async fn drop_changes<S: State>(&self, state: &SyncRwLock<S>) -> Vec<(usize, Option<Value>)> {
207
2
        let pending = {
208
2
            let state = state.read();
209
2
            let mut pending = Vec::new();
210
2
            for ospage_id in self.dirty.get_marked_and_clear() {
211
2
                let segment_off = ospage_id * self.page_size;
212
2
                let litepage_start = (self.litepage_base + segment_off) / LYTEPAGE_SIZE;
213
2
                let mem_addr = self.base + segment_off;
214
8
                for i in 
0..self.page_size / LYTEPAGE_SIZE2
{
215
8
                    let key = litepage_state_key((litepage_start + i) as u32);
216
8
                    pending.push((mem_addr + i * LYTEPAGE_SIZE, state.get(key)));
217
8
                }
218
            }
219
2
            pending
220
        };
221
222
2
        let mut restores = Vec::with_capacity(pending.len());
223
8
        for (mem_addr, old) in 
pending2
{
224
8
            restores.push((mem_addr, old.await));
225
        }
226
2
        restores
227
2
    }
228
}
229
230
#[derive(Clone)]
231
struct InstanceSegment {
232
    shm: SharedMemory,
233
    seg: Arc<Segment>,
234
    parking: Arc<crate::sync::ParkingSpot>,
235
}
236
237
impl InstanceSegment {
238
60
    fn new(
239
60
        base: usize, size: usize, litepage_base: usize, page_size: usize,
240
60
    ) -> Result<Self, userspace_pagefault::Error> {
241
        Ok(Self {
242
60
            shm: SharedMemory::new(size)
?0
,
243
60
            seg: Arc::new(Segment::new(base, size, litepage_base, page_size)),
244
60
            parking: Arc::new(crate::sync::ParkingSpot::new()),
245
        })
246
60
    }
247
}
248
249
#[derive(Clone)]
250
struct Context {
251
    state: crate::instance::LyquidState,
252
    network: Arc<Segment>,
253
    instance: InstanceSegment,
254
    litepage: std::ops::Range<usize>,
255
}
256
257
impl AsyncPageStore for Context {
258
1.42k
    fn page_fault_async(&mut self, offset: usize, length: usize, access: AccessType) -> AsyncPageFaultFuture<'_> {
259
1.42k
        Box::pin(async move {
260
1.42k
            if !self.litepage.contains(&offset) {
261
                // the accessed page is outside LytePage (network/instance) address range. Returns None
262
                // because no loading is needed.
263
929
                return None;
264
494
            }
265
266
            /*
267
            println!(
268
                "{:?} pagefault {:?} @ {:x}",
269
                self as *mut Self,
270
                access,
271
                offset - BEGIN_GUARD_SIZE
272
            );
273
            */
274
275
            // the loaded OS page may span across multiple LytePages.
276
494
            let litepage_start = (offset - self.litepage.start) / LYTEPAGE_SIZE;
277
282
            let (seg, reads) = {
278
                // check which segment this address accesses.
279
494
                let (state, seg) = match offset {
280
494
                    
off297
if off < self.instance.seg.base =>
(self.state.network.read(), &self.network)297
,
281
197
                    _ => (self.state.instance.read(), &self.instance.seg),
282
                };
283
494
                if let AccessType::Write = access {
284
386
                    // mark the page as dirty
285
386
                    seg.mark_dirty(offset);
286
386
                
}108
287
494
                if seg.is_loaded(offset) {
288
                    // the page was already loaded
289
212
                    return None;
290
282
                }
291
292
                //println!("loading @ {:x}", offset - BEGIN_GUARD_SIZE);
293
294
282
                let reads = (0..length / LYTEPAGE_SIZE)
295
1.12k
                    .
map282
(|i| state.get(litepage_state_key((litepage_start + i) as u32)))
296
282
                    .collect::<Vec<_>>();
297
282
                (seg, reads)
298
            };
299
300
            // collect the contents of all LytePages read from state store.
301
282
            let mut chunks = Vec::with_capacity(reads.len());
302
1.12k
            for read in 
reads282
{
303
1.12k
                chunks.push(Box::new(read.await.unwrap_or_else(|| 
{1.07k
304
1.07k
                    std::iter::repeat_with(|| 0)
305
1.07k
                        .take(LYTEPAGE_SIZE)
306
1.07k
                        .collect::<Vec<u8>>()
307
1.07k
                        .into()
308
1.12k
                
}1.07k
)) as Box<dyn AsRef<[u8]>>);
309
            }
310
282
            seg.mark_loaded(offset);
311
282
            Some(Box::new(chunks.into_iter()) as PageChunks<'_>)
312
1.42k
        })
313
1.42k
    }
314
}
315
316
/// Factory and shared backing resources for Lyquid virtual memory objects.
317
pub struct LyteMemoryManager {
318
    instance: InstanceSegment,
319
    instance_mapped: userspace_pagefault::Segment,
320
    network_range: std::ops::Range<usize>,
321
    litepage: std::ops::Range<usize>,
322
    page_size: usize,
323
    total: usize,
324
}
325
326
impl LyteMemoryManager {
327
    /// Create the shared instance segment and reserve the LyteMemory virtual address range.
328
60
    pub fn new() -> Result<Self, userspace_pagefault::Error> {
329
60
        lytememory_setup();
330
60
        let wasm_usable = lyquid::LYTEMEM_SIZE_IN_MB << 20;
331
60
        let network = lyquid::NETWORK_MEMSIZE_IN_MB << 20;
332
60
        let instance = lyquid::INSTANCE_MEMSIZE_IN_MB << 20;
333
        // WASM 32-bit starting address for LytePage area
334
60
        let start = BEGIN_GUARD_SIZE + lyquid::LYTEMEM_BASE;
335
        // WASM 32-bit ending address for LytePage area
336
60
        let end = start + instance + network;
337
        // total virtual size, including guards
338
60
        let total = BEGIN_GUARD_SIZE + wasm_usable + END_GUARD_SIZE;
339
60
        assert!(end <= total);
340
60
        let page_size = userspace_pagefault::get_page_size()
?0
;
341
60
        let network_range = start..start + network;
342
60
        let instance_range = start + network..end;
343
        // all LyteMemories generated from this LyteMemoryManager will share the same
344
        // InstanceSegment
345
60
        let instance = InstanceSegment::new(
346
60
            instance_range.start,
347
60
            instance_range.len(),
348
60
            network_range.len(),
349
60
            page_size,
350
0
        )?;
351
352
60
        let instance_mapped =
353
60
            userspace_pagefault::Segment::new(None, total, page_size, userspace_pagefault::ProtFlags::PROT_NONE)
?0
;
354
60
        instance_mapped.make_shared(
355
60
            instance.seg.base,
356
60
            &instance.shm,
357
60
            userspace_pagefault::ProtFlags::PROT_READ | userspace_pagefault::ProtFlags::PROT_WRITE,
358
0
        )?;
359
360
60
        Ok(Self {
361
60
            instance,
362
60
            instance_mapped,
363
60
            network_range,
364
60
            litepage: start..end,
365
60
            page_size,
366
60
            total,
367
60
        })
368
60
    }
369
370
    /// Flush the dirty pages in the instance segement to the key-value state and clear the dirty
371
    /// bitmask.
372
58
    pub(crate) async fn flush_instance_state<'a, T, S: State>(
373
58
        &self, lock: &'a RwLock<T>, state: &'a SyncRwLock<S>,
374
58
    ) -> (RwLockWriteGuard<'a, T>, SyncRwLockWriteGuard<'a, S>) {
375
58
        let guard = lock.write().await;
376
58
        self.instance
377
58
            .seg
378
58
            .writeback_changes(state, self.instance_mapped.as_slice())
379
58
            .await;
380
58
        let state = state.write();
381
58
        (guard, state)
382
58
    }
383
384
    /// Create a new LyteMemory object. The object can be clone-shared but the underlying resource
385
    /// is the same.
386
108
    pub fn new_lytememory<G: VmGuest>(
387
108
        &self, state: crate::instance::LyquidState, category: StateCategory,
388
108
    ) -> Result<LyteMemory<G>, userspace_pagefault::Error> {
389
108
        let ctx = Context {
390
108
            state,
391
108
            network: Arc::new(Segment::new(
392
108
                self.network_range.start,
393
108
                self.network_range.len(),
394
108
                0,
395
108
                self.page_size,
396
108
            )),
397
108
            instance: self.instance.clone(),
398
108
            litepage: self.litepage.clone(),
399
108
        };
400
401
108
        let mem = Arc::new(PagedSegment::new_async(self.total, ctx.clone(), None)
?0
);
402
108
        mem.make_shared(self.instance.seg.base, &self.instance.shm)
?0
;
403
404
108
        Ok(LyteMemory {
405
108
            mem,
406
108
            ctx,
407
108
            category,
408
108
            guest: PhantomData,
409
108
        })
410
108
    }
411
}
412
413
/// A sharable object that represents a userspace-paged virtual memory used by LVM.
414
#[derive(Clone)]
415
pub struct LyteMemory<G: VmGuest = Wasm32> {
416
    mem: Arc<PagedSegment<'static>>,
417
    ctx: Context,
418
    category: StateCategory,
419
    guest: PhantomData<G>,
420
}
421
422
impl<G: VmGuest> LyteMemory<G> {
423
    /// Reset the dirty-page detection of the instance segment. This is required after a flush
424
    /// write-back of the instance segment, otherwise future changes will not be detected.
425
0
    pub fn reset_instance_write_detection(&self) {
426
0
        self.mem
427
0
            .reset_write_detection(self.ctx.instance.seg.base, self.ctx.instance.seg.length)
428
0
            .expect("Failed to reset write detection."); // FIXME: handle this error
429
0
    }
430
431
    /// Flush the dirty pages of the network segment to the key-value state carried by LyteMemory.
432
    /// The dirty-page detection is also reset properly after flushing the pages.
433
94
    pub async fn flush_network_state(&self) {
434
94
        self.ctx
435
94
            .network
436
94
            .writeback_changes(self.ctx.state.network.as_ref(), self.mem.as_slice())
437
94
            .await;
438
94
        self.mem
439
94
            .reset_write_detection(self.ctx.network.base, self.ctx.network.length)
440
94
            .expect("Failed to reset write detection."); // FIXME: handle this error
441
        // the lock is released at the end
442
94
    }
443
444
    /// Revert all dirty pages of the network segment.
445
    /// The dirty-page detection is also reset properly after flushing the pages.
446
2
    pub async fn revert_network_state(&self) {
447
2
        let restores = self.ctx.network.drop_changes(self.ctx.state.network.as_ref()).await;
448
2
        let (base, len) = self.mem.as_raw_parts();
449
2
        let mem = unsafe { std::slice::from_raw_parts_mut(base, len) };
450
8
        for (mem_addr, old) in 
restores2
{
451
8
            let litepage = &mut mem[mem_addr..mem_addr + LYTEPAGE_SIZE];
452
8
            match old {
453
2
                Some(old) => litepage.copy_from_slice(&old),
454
6
                None => litepage.fill(0),
455
            }
456
        }
457
2
        self.mem
458
2
            .reset_write_detection(self.ctx.network.base, self.ctx.network.length)
459
2
            .expect("Failed to reset write detection."); // FIXME: handle this error
460
        // the lock is released at the end
461
2
    }
462
463
    /// Get the raw pointer to the host memory at a given WASM address.
464
    #[inline(always)]
465
8.72k
    pub fn host_ptr_by_wasm_addr(&self, addr: G::Usize) -> *mut u8 {
466
8.72k
        let (base, _) = self.mem.as_raw_parts();
467
8.72k
        base.wrapping_add(G::usize_to_host(addr) + BEGIN_GUARD_SIZE)
468
8.72k
    }
469
470
    /// Get a read-only slice of memory by WASM address.
471
3.34k
    pub fn slice_by_wasm_addr(&self, offset: G::Usize, length: G::Usize) -> &[u8] {
472
3.34k
        unsafe { std::slice::from_raw_parts(self.host_ptr_by_wasm_addr(offset), G::usize_to_host(length)) }
473
3.34k
    }
474
475
    /// Get a mutable slice of the memory with WASM address.
476
    ///
477
    /// # Safety
478
    ///
479
    /// The caller must ensure no other live references alias the returned range for the duration
480
    /// of the returned slice.
481
    #[allow(clippy::mut_from_ref)]
482
5.25k
    pub unsafe fn slice_by_wasm_addr_mut(&self, offset: G::Usize, length: G::Usize) -> &mut [u8] {
483
5.25k
        unsafe { std::slice::from_raw_parts_mut(self.host_ptr_by_wasm_addr(offset), G::usize_to_host(length)) }
484
5.25k
    }
485
486
    /// Get the function category of this LyteMemory's owner.
487
642
    pub fn category(&self) -> StateCategory {
488
642
        self.category
489
642
    }
490
491
    /// Return the wait/notify queue associated with the shared instance segment.
492
216
    pub(crate) fn instance_parking_spot(&self) -> Arc<crate::sync::ParkingSpot> {
493
216
        self.ctx.instance.parking.clone()
494
216
    }
495
}
496
497
struct WasmRuntimeMemory<G: VmGuest> {
498
    inner: LyteMemory<G>,
499
    size: usize,
500
}
501
502
unsafe impl<G: VmGuest> wasmtime::LinearMemory for WasmRuntimeMemory<G> {
503
216
    fn byte_size(&self) -> usize {
504
216
        self.size
505
216
    }
506
507
0
    fn byte_capacity(&self) -> usize {
508
0
        self.size
509
0
    }
510
511
0
    fn grow_to(&mut self, size: usize) -> wasmtime::Result<()> {
512
0
        if size > self.size {
513
0
            return Err(wasmtime::Error::msg("LyteMemory is not growable."));
514
0
        }
515
0
        Ok(())
516
0
    }
517
518
108
    fn as_ptr(&self) -> *mut u8 {
519
108
        let (ptr, _) = self.inner.mem.as_raw_parts();
520
108
        unsafe { ptr.add(BEGIN_GUARD_SIZE) }
521
108
    }
522
}
523
524
/// Memory factory that makes wasmtime memory for the given LyteMemory.
525
pub struct WasmMemoryCreator<G: VmGuest = Wasm32>(pub LyteMemory<G>);
526
527
unsafe impl<G: VmGuest> wasmtime::MemoryCreator for WasmMemoryCreator<G> {
528
108
    fn new_memory(
529
108
        &self, _mt: wasmtime::MemoryType, _min: usize, _max: Option<usize>, _reserved: Option<usize>, _guard: usize,
530
108
    ) -> Result<Box<dyn wasmtime::LinearMemory>, String> {
531
108
        Ok(Box::new(WasmRuntimeMemory {
532
108
            inner: self.0.clone(),
533
108
            size: lyquid::LYTEMEM_SIZE_IN_MB << 20,
534
108
        }))
535
108
    }
536
}