Coverage Report

Created: 2026-07-26 02:19

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
44
fn lytememory_setup() {
22
    static INITIALIZED: std::sync::Once = std::sync::Once::new();
23
44
    INITIALIZED.call_once(|| 
{22
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
22
        let _ = wasmtime::Engine::new(wasmtime::Config::new().macos_use_mach_ports(false));
29
22
    });
30
44
}
31
32
struct Bitmask(Box<[atomic::AtomicU64]>);
33
34
impl Bitmask {
35
243
    fn new(size: usize) -> Self {
36
        Self(
37
983k
            
std::iter::repeat_with243
(|| atomic::AtomicU64::new(0))
38
243
                .take((size + 63) >> 6)
39
243
                .collect::<Vec<_>>()
40
243
                .into(),
41
        )
42
243
    }
43
44
    #[inline(always)]
45
595
    fn mark(&self, idx: usize) {
46
595
        self.0[idx >> 6].fetch_or(1 << (idx & 63), atomic::Ordering::Relaxed);
47
595
    }
48
49
    #[inline(always)]
50
387
    fn is_marked(&self, idx: usize) -> bool {
51
387
        (self.0[idx >> 6].fetch_or(0, atomic::Ordering::Relaxed) >> (idx & 63)) & 1 == 1
52
387
    }
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
323
    fn fast_log(mut x: u64) -> u64 {
64
323
        x |= x >> 1;
65
323
        x |= x >> 2;
66
323
        x |= x >> 4;
67
323
        x |= x >> 8;
68
323
        x |= x >> 16;
69
323
        x |= x >> 32;
70
323
        Self::TABLE[((x.wrapping_mul(0x03f6eaf2cd271461)) >> 58) as usize]
71
323
    }
72
73
    #[inline]
74
126
    fn get_marked_and_clear(&self) -> impl Iterator<Item = usize> + '_ {
75
503k
        
self.0126
.
iter126
().
zip126
((
0..126
).
step_by126
(64)).
flat_map126
(|(bits, i)| {
76
503k
            let mut v = bits.swap(0, atomic::Ordering::Relaxed);
77
504k
            
std::iter::from_fn503k
(move || {
78
                // this lowbit finding algorithm is suitable for a sparse bitmask, which is the
79
                // case typically for LyteMemory dirty pages
80
504k
                let mut lowbit = v;
81
504k
                if lowbit == 0 {
82
503k
                    None
83
                } else {
84
323
                    lowbit = lowbit & lowbit.wrapping_neg();
85
323
                    v ^= lowbit; // clear the lowest bit
86
323
                    Some(i + Self::fast_log(lowbit) as usize)
87
                }
88
504k
            })
89
503k
        })
90
126
    }
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
1.79k
fn litepage_state_key(litepage_id: u32) -> Key {
113
    static PREFIX: std::sync::OnceLock<Key> = std::sync::OnceLock::new();
114
1.79k
    let key = PREFIX.get_or_init(|| lyquid::LYTEMEM_PAGE_PREFIX.
into21
());
115
    // using big-endian here so the state debug will be easier to read
116
1.79k
    let id: LazyBytes = litepage_id.to_be_bytes().into();
117
1.79k
    id.prepend(key)
118
1.79k
}
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
120
    fn new(base: usize, length: usize, litepage_base: usize, page_size: usize) -> Self {
137
        // align to the closest page
138
120
        let bitmask_size = length.div_ceil(page_size);
139
120
        Self {
140
120
            dirty: Bitmask::new(bitmask_size),
141
120
            loaded: Bitmask::new(bitmask_size),
142
120
            page_size,
143
120
            litepage_base,
144
120
            base,
145
120
            length,
146
120
        }
147
120
    }
148
149
    #[inline(always)]
150
903
    fn addr_to_idx(&self, mem_addr: usize) -> usize {
151
903
        (mem_addr - self.base) / self.page_size
152
903
    }
153
154
    #[inline(always)]
155
311
    fn mark_dirty(&self, mem_addr: usize) {
156
311
        self.dirty.mark(self.addr_to_idx(mem_addr));
157
311
    }
158
159
    #[inline(always)]
160
387
    fn is_loaded(&self, mem_addr: usize) -> bool {
161
387
        self.loaded.is_marked(self.addr_to_idx(mem_addr))
162
387
    }
163
164
    #[inline(always)]
165
205
    fn mark_loaded(&self, mem_addr: usize) {
166
205
        self.loaded.mark(self.addr_to_idx(mem_addr));
167
205
    }
168
169
    /// Write the dirty pages to the key-value state and clear the dirty bitmask.
170
121
    async fn writeback_changes<S: State>(&self, state: &SyncRwLock<S>, mem: &[u8]) {
171
121
        let pending = {
172
121
            let state = state.read();
173
121
            let mut pending = Vec::new();
174
242
            for ospage_id in 
self.dirty121
.
get_marked_and_clear121
() {
175
242
                let segment_off = ospage_id * self.page_size;
176
242
                let litepage_start = (self.litepage_base + segment_off) / LYTEPAGE_SIZE;
177
242
                let mem_addr = self.base + segment_off;
178
242
                let ospage = &mem[mem_addr..mem_addr + self.page_size];
179
968
                for (i, litepage) in 
ospage242
.
chunks242
(LYTEPAGE_SIZE).
enumerate242
() {
180
968
                    let key = litepage_state_key((litepage_start + i) as u32);
181
968
                    let litepage: Value = litepage.to_vec().into();
182
968
                    pending.push((key.clone(), litepage, state.get(key)));
183
968
                }
184
            }
185
121
            pending
186
        };
187
188
121
        let mut changes = Vec::new();
189
968
        for (key, litepage, old) in 
pending121
{
190
968
            let changed = match old.await {
191
71
                Some(old) => litepage.as_ref() != old.as_ref(), // if the page was previous persisted, compare
192
825k
                None => !
litepage.as_ref().iter()897
.
all897
(|&b| b == 0), // otherwise check if it's empty
193
            };
194
968
            if changed {
195
234
                changes.push((key, litepage));
196
734
            }
197
        }
198
199
121
        let mut state = state.write();
200
234
        for (key, litepage) in 
changes121
{
201
234
            state.set(key, Some(litepage));
202
234
        }
203
121
    }
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
44
    fn new(
239
44
        base: usize, size: usize, litepage_base: usize, page_size: usize,
240
44
    ) -> Result<Self, userspace_pagefault::Error> {
241
        Ok(Self {
242
44
            shm: SharedMemory::new(size)
?0
,
243
44
            seg: Arc::new(Segment::new(base, size, litepage_base, page_size)),
244
44
            parking: Arc::new(crate::sync::ParkingSpot::new()),
245
        })
246
44
    }
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
941
    fn page_fault_async(&mut self, offset: usize, length: usize, access: AccessType) -> AsyncPageFaultFuture<'_> {
259
941
        Box::pin(async move {
260
941
            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
554
                return None;
264
387
            }
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
387
            let litepage_start = (offset - self.litepage.start) / LYTEPAGE_SIZE;
277
205
            let (seg, reads) = {
278
                // check which segment this address accesses.
279
387
                let (state, seg) = match offset {
280
387
                    
off235
if off < self.instance.seg.base =>
(self.state.network.read(), &self.network)235
,
281
152
                    _ => (self.state.instance.read(), &self.instance.seg),
282
                };
283
387
                if let AccessType::Write = access {
284
311
                    // mark the page as dirty
285
311
                    seg.mark_dirty(offset);
286
311
                
}76
287
387
                if seg.is_loaded(offset) {
288
                    // the page was already loaded
289
182
                    return None;
290
205
                }
291
292
                //println!("loading @ {:x}", offset - BEGIN_GUARD_SIZE);
293
294
205
                let reads = (0..length / LYTEPAGE_SIZE)
295
820
                    .
map205
(|i| state.get(litepage_state_key((litepage_start + i) as u32)))
296
205
                    .collect::<Vec<_>>();
297
205
                (seg, reads)
298
            };
299
300
            // collect the contents of all LytePages read from state store.
301
205
            let mut chunks = Vec::with_capacity(reads.len());
302
820
            for read in 
reads205
{
303
820
                chunks.push(Box::new(read.await.unwrap_or_else(|| 
{781
304
781
                    std::iter::repeat_with(|| 0)
305
781
                        .take(LYTEPAGE_SIZE)
306
781
                        .collect::<Vec<u8>>()
307
781
                        .into()
308
820
                
}781
)) as Box<dyn AsRef<[u8]>>);
309
            }
310
205
            seg.mark_loaded(offset);
311
205
            Some(Box::new(chunks.into_iter()) as PageChunks<'_>)
312
941
        })
313
941
    }
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
44
    pub fn new() -> Result<Self, userspace_pagefault::Error> {
329
44
        lytememory_setup();
330
44
        let wasm_usable = lyquid::LYTEMEM_SIZE_IN_MB << 20;
331
44
        let network = lyquid::NETWORK_MEMSIZE_IN_MB << 20;
332
44
        let instance = lyquid::INSTANCE_MEMSIZE_IN_MB << 20;
333
        // WASM 32-bit starting address for LytePage area
334
44
        let start = BEGIN_GUARD_SIZE + lyquid::LYTEMEM_BASE;
335
        // WASM 32-bit ending address for LytePage area
336
44
        let end = start + instance + network;
337
        // total virtual size, including guards
338
44
        let total = BEGIN_GUARD_SIZE + wasm_usable + END_GUARD_SIZE;
339
44
        assert!(end <= total);
340
44
        let page_size = userspace_pagefault::get_page_size()
?0
;
341
44
        let network_range = start..start + network;
342
44
        let instance_range = start + network..end;
343
        // all LyteMemories generated from this LyteMemoryManager will share the same
344
        // InstanceSegment
345
44
        let instance = InstanceSegment::new(
346
44
            instance_range.start,
347
44
            instance_range.len(),
348
44
            network_range.len(),
349
44
            page_size,
350
0
        )?;
351
352
44
        let instance_mapped =
353
44
            userspace_pagefault::Segment::new(None, total, page_size, userspace_pagefault::ProtFlags::PROT_NONE)
?0
;
354
44
        instance_mapped.make_shared(
355
44
            instance.seg.base,
356
44
            &instance.shm,
357
44
            userspace_pagefault::ProtFlags::PROT_READ | userspace_pagefault::ProtFlags::PROT_WRITE,
358
0
        )?;
359
360
44
        Ok(Self {
361
44
            instance,
362
44
            instance_mapped,
363
44
            network_range,
364
44
            litepage: start..end,
365
44
            page_size,
366
44
            total,
367
44
        })
368
44
    }
369
370
    /// Flush the dirty pages in the instance segement to the key-value state and clear the dirty
371
    /// bitmask.
372
43
    pub(crate) async fn flush_instance_state<'a, T, S: State>(
373
43
        &self, lock: &'a RwLock<T>, state: &'a SyncRwLock<S>,
374
43
    ) -> (RwLockWriteGuard<'a, T>, SyncRwLockWriteGuard<'a, S>) {
375
43
        let guard = lock.write().await;
376
43
        self.instance
377
43
            .seg
378
43
            .writeback_changes(state, self.instance_mapped.as_slice())
379
43
            .await;
380
43
        let state = state.write();
381
43
        (guard, state)
382
43
    }
383
384
    /// Create a new LyteMemory object. The object can be clone-shared but the underlying resource
385
    /// is the same.
386
76
    pub fn new_lytememory<G: VmGuest>(
387
76
        &self, state: crate::instance::LyquidState, category: StateCategory,
388
76
    ) -> Result<LyteMemory<G>, userspace_pagefault::Error> {
389
76
        let ctx = Context {
390
76
            state,
391
76
            network: Arc::new(Segment::new(
392
76
                self.network_range.start,
393
76
                self.network_range.len(),
394
76
                0,
395
76
                self.page_size,
396
76
            )),
397
76
            instance: self.instance.clone(),
398
76
            litepage: self.litepage.clone(),
399
76
        };
400
401
76
        let mem = Arc::new(PagedSegment::new_async(self.total, ctx.clone(), None)
?0
);
402
76
        mem.make_shared(self.instance.seg.base, &self.instance.shm)
?0
;
403
404
76
        Ok(LyteMemory {
405
76
            mem,
406
76
            ctx,
407
76
            category,
408
76
            guest: PhantomData,
409
76
        })
410
76
    }
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
78
    pub async fn flush_network_state(&self) {
434
78
        self.ctx
435
78
            .network
436
78
            .writeback_changes(self.ctx.state.network.as_ref(), self.mem.as_slice())
437
78
            .await;
438
78
        self.mem
439
78
            .reset_write_detection(self.ctx.network.base, self.ctx.network.length)
440
78
            .expect("Failed to reset write detection."); // FIXME: handle this error
441
        // the lock is released at the end
442
78
    }
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
5.70k
    pub fn host_ptr_by_wasm_addr(&self, addr: G::Usize) -> *mut u8 {
466
5.70k
        let (base, _) = self.mem.as_raw_parts();
467
5.70k
        base.wrapping_add(G::usize_to_host(addr) + BEGIN_GUARD_SIZE)
468
5.70k
    }
469
470
    /// Get a read-only slice of memory by WASM address.
471
2.15k
    pub fn slice_by_wasm_addr(&self, offset: G::Usize, length: G::Usize) -> &[u8] {
472
2.15k
        unsafe { std::slice::from_raw_parts(self.host_ptr_by_wasm_addr(offset), G::usize_to_host(length)) }
473
2.15k
    }
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
3.54k
    pub unsafe fn slice_by_wasm_addr_mut(&self, offset: G::Usize, length: G::Usize) -> &mut [u8] {
483
3.54k
        unsafe { std::slice::from_raw_parts_mut(self.host_ptr_by_wasm_addr(offset), G::usize_to_host(length)) }
484
3.54k
    }
485
486
    /// Get the function category of this LyteMemory's owner.
487
383
    pub fn category(&self) -> StateCategory {
488
383
        self.category
489
383
    }
490
491
    /// Return the wait/notify queue associated with the shared instance segment.
492
152
    pub(crate) fn instance_parking_spot(&self) -> Arc<crate::sync::ParkingSpot> {
493
152
        self.ctx.instance.parking.clone()
494
152
    }
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
152
    fn byte_size(&self) -> usize {
504
152
        self.size
505
152
    }
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
76
    fn as_ptr(&self) -> *mut u8 {
519
76
        let (ptr, _) = self.inner.mem.as_raw_parts();
520
76
        unsafe { ptr.add(BEGIN_GUARD_SIZE) }
521
76
    }
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
76
    fn new_memory(
529
76
        &self, _mt: wasmtime::MemoryType, _min: usize, _max: Option<usize>, _reserved: Option<usize>, _guard: usize,
530
76
    ) -> Result<Box<dyn wasmtime::LinearMemory>, String> {
531
76
        Ok(Box::new(WasmRuntimeMemory {
532
76
            inner: self.0.clone(),
533
76
            size: lyquid::LYTEMEM_SIZE_IN_MB << 20,
534
76
        }))
535
76
    }
536
}