Skip to main content

wasmtime/runtime/vm/instance/allocator/
pooling.rs

1//! Implements the pooling instance allocator.
2//!
3//! The pooling instance allocator maps memory in advance and allocates
4//! instances, memories, tables, and stacks from a pool of available resources.
5//! Using the pooling instance allocator can speed up module instantiation when
6//! modules can be constrained based on configurable limits
7//! ([`InstanceLimits`]). Each new instance is stored in a "slot"; as instances
8//! are allocated and freed, these slots are either filled or emptied:
9//!
10//! ```text
11//! ┌──────┬──────┬──────┬──────┬──────┐
12//! │Slot 0│Slot 1│Slot 2│Slot 3│......│
13//! └──────┴──────┴──────┴──────┴──────┘
14//! ```
15//!
16//! Each slot has a "slot ID"--an index into the pool. Slot IDs are handed out
17//! by the [`index_allocator`] module. Note that each kind of pool-allocated
18//! item is stored in its own separate pool: [`memory_pool`], [`table_pool`],
19//! [`stack_pool`]. See those modules for more details.
20
21mod decommit_queue;
22mod index_allocator;
23mod memory_pool;
24mod metrics;
25mod table_pool;
26
27#[cfg(feature = "gc")]
28mod gc_heap_pool;
29
30#[cfg(all(feature = "async"))]
31mod generic_stack_pool;
32#[cfg(all(feature = "async", unix, not(miri)))]
33mod unix_stack_pool;
34
35#[cfg(all(feature = "async"))]
36cfg_select! {
37    all(unix, not(miri), not(asan)) => {
38        use unix_stack_pool as stack_pool;
39    }
40    _ => {
41        use generic_stack_pool as stack_pool;
42    }
43}
44
45use self::decommit_queue::DecommitQueue;
46use self::memory_pool::MemoryPool;
47pub use self::metrics::PoolingAllocatorMetrics;
48use self::table_pool::TablePool;
49use super::{
50    InstanceAllocationRequest, InstanceAllocator, MemoryAllocationIndex, TableAllocationIndex,
51};
52use crate::Enabled;
53use crate::config::PoolingAllocationConfig;
54use crate::prelude::*;
55use crate::runtime::vm::{
56    CompiledModuleId, Memory, Table,
57    instance::Instance,
58    mpk::{self, ProtectionKey, ProtectionMask},
59    sys::vm::PageMap,
60};
61use core::future::Future;
62use core::pin::Pin;
63use core::sync::atomic::AtomicUsize;
64use std::borrow::Cow;
65use std::fmt::Display;
66use std::sync::{Mutex, MutexGuard};
67use std::{
68    mem,
69    sync::atomic::{AtomicU64, Ordering},
70};
71use wasmtime_environ::{
72    DefinedMemoryIndex, DefinedTableIndex, HostPtr, MemoryKind, Module, Tunables, VMOffsets,
73};
74
75#[cfg(feature = "gc")]
76use super::GcHeapAllocationIndex;
77#[cfg(feature = "gc")]
78use crate::runtime::vm::{GcHeap, GcRuntime};
79#[cfg(feature = "gc")]
80use gc_heap_pool::GcHeapPool;
81
82/// Pad a value out to a full cache line (or two, on aarch64 prefetch
83/// granularity) so neighboring shards don't false-share.
84#[repr(align(128))]
85#[derive(Debug)]
86struct CachePadded<T>(T);
87
88/// Identifier of one shard of the pooling allocator's sharded data
89/// structures (the decommit queues and each pool's index allocator).
90#[derive(Copy, Clone, Debug, PartialEq, Eq)]
91pub(crate) struct ShardId(u32);
92
93impl ShardId {
94    pub(crate) fn from_index(index: usize) -> ShardId {
95        ShardId(u32::try_from(index).unwrap())
96    }
97
98    pub(crate) fn index(self) -> usize {
99        usize::try_from(self.0).unwrap()
100    }
101}
102
103/// The number of shards used for the pooling allocator's sharded data
104/// structures: one per available CPU, capped to 16.
105///
106/// The cap bounds worst-case probing when pools run near-full, the
107/// dilution of per-shard warm-slot budgets, and per-shard memory
108/// overhead, while still being enough shards to make lock collisions
109/// rare given the very short critical sections involved.
110pub(crate) fn default_shard_count() -> u32 {
111    let n = std::thread::available_parallelism()
112        .map(|n| n.get())
113        .unwrap_or(1)
114        .min(16);
115    u32::try_from(n).unwrap()
116}
117
118/// Pick this thread's shard (used for both the sharded decommit queue and
119/// the sharded index allocators): assigned round-robin at first use per
120/// thread, cached in a thread-local.
121pub(crate) fn thread_shard(nshards: usize) -> ShardId {
122    static NEXT_SHARD: AtomicUsize = AtomicUsize::new(0);
123    std::thread_local! {
124        static SHARD: usize = NEXT_SHARD.fetch_add(1, Ordering::Relaxed);
125    }
126    ShardId::from_index(SHARD.with(|s| *s) % nshards)
127}
128
129/// Enumerate all shard ids for a sharded structure with `nshards` shards,
130/// starting with the current thread's home shard and wrapping around.
131pub(crate) fn shard_ids_from_home(nshards: usize) -> impl Iterator<Item = ShardId> {
132    let home = thread_shard(nshards).index();
133    (0..nshards).map(move |i| ShardId::from_index((home + i) % nshards))
134}
135
136#[cfg(feature = "async")]
137use stack_pool::StackPool;
138
139#[cfg(feature = "component-model")]
140use wasmtime_environ::{
141    StaticModuleIndex,
142    component::{Component, VMComponentOffsets},
143};
144
145fn round_up_to_pow2(n: usize, to: usize) -> usize {
146    debug_assert!(to > 0);
147    debug_assert!(to.is_power_of_two());
148    (n + to - 1) & !(to - 1)
149}
150
151impl PoolingAllocationConfig {
152    /// Tests whether [`Self::pagemap_scan`] is available or not on the host
153    /// system.
154    pub fn is_pagemap_scan_available() -> bool {
155        PageMap::new().is_some()
156    }
157}
158
159/// An error returned when the pooling allocator cannot allocate a table,
160/// memory, etc... because the maximum number of concurrent allocations for that
161/// entity has been reached.
162#[derive(Debug)]
163pub struct PoolConcurrencyLimitError {
164    limit: usize,
165    kind: Cow<'static, str>,
166}
167
168impl core::error::Error for PoolConcurrencyLimitError {}
169
170impl Display for PoolConcurrencyLimitError {
171    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
172        let limit = self.limit;
173        let kind = &self.kind;
174        write!(f, "maximum concurrent limit of {limit} for {kind} reached")
175    }
176}
177
178impl PoolConcurrencyLimitError {
179    fn new(limit: usize, kind: impl Into<Cow<'static, str>>) -> Self {
180        Self {
181            limit,
182            kind: kind.into(),
183        }
184    }
185}
186
187/// Implements the pooling instance allocator.
188///
189/// This allocator internally maintains pools of instances, memories, tables,
190/// and stacks.
191///
192/// Note: the resource pools are manually dropped so that the fault handler
193/// terminates correctly.
194#[derive(Debug)]
195pub struct PoolingInstanceAllocator {
196    // The number of live core module and component instances at any given
197    // time. Note that this can temporarily go over the configured limit. This
198    // doesn't mean we have actually overshot, but that we attempted to allocate
199    // a new instance and incremented the counter, we've seen (or are about to
200    // see) that the counter is beyond the configured threshold, and are going
201    // to decrement the counter and return an error but haven't done so yet. See
202    // the increment trait methods for more details.
203    live_core_instances: AtomicU64,
204    live_component_instances: AtomicU64,
205
206    /// Sharded to avoid a single global mutex on every deallocation when
207    /// decommit batching is enabled: each thread appends to its own shard
208    /// (assigned round-robin at first use) and flushes that shard when it
209    /// reaches the configured batch size. Slot-exhaustion paths flush all
210    /// shards.
211    decommit_queues: Box<[CachePadded<Mutex<DecommitQueue>>]>,
212
213    memories: MemoryPool,
214    live_memories: AtomicUsize,
215
216    tables: TablePool,
217    live_tables: AtomicUsize,
218
219    #[cfg(feature = "gc")]
220    gc_heaps: Option<GcHeapPool>,
221    #[cfg(feature = "gc")]
222    live_gc_heaps: AtomicUsize,
223
224    #[cfg(feature = "async")]
225    stacks: StackPool,
226    #[cfg(feature = "async")]
227    live_stacks: AtomicUsize,
228
229    pagemap: Option<PageMap>,
230    config: PoolingAllocationConfig,
231}
232
233impl Drop for PoolingInstanceAllocator {
234    fn drop(&mut self) {
235        if !cfg!(debug_assertions) {
236            return;
237        }
238
239        // NB: when cfg(not(debug_assertions)) it is okay that we don't flush
240        // the queue, as the sub-pools will unmap those ranges anyways, so
241        // there's no point in decommitting them. But we do need to flush the
242        // queue when debug assertions are enabled to make sure that all
243        // entities get returned to their associated sub-pools and we can
244        // differentiate between a leaking slot and an enqueued-for-decommit
245        // slot.
246        self.flush_all_decommit_queues();
247
248        debug_assert_eq!(self.live_component_instances.load(Ordering::Acquire), 0);
249        debug_assert_eq!(self.live_core_instances.load(Ordering::Acquire), 0);
250        debug_assert_eq!(self.live_memories.load(Ordering::Acquire), 0);
251        debug_assert_eq!(self.live_tables.load(Ordering::Acquire), 0);
252
253        debug_assert!(self.memories.is_empty());
254        debug_assert!(self.tables.is_empty());
255
256        #[cfg(feature = "gc")]
257        if let Some(gc_heaps) = &self.gc_heaps {
258            debug_assert!(gc_heaps.is_empty());
259            debug_assert_eq!(self.live_gc_heaps.load(Ordering::Acquire), 0);
260        }
261
262        #[cfg(feature = "async")]
263        {
264            debug_assert!(self.stacks.is_empty());
265            debug_assert_eq!(self.live_stacks.load(Ordering::Acquire), 0);
266        }
267    }
268}
269
270impl PoolingInstanceAllocator {
271    /// Releases memory kept resident for unused slots, returning the number
272    /// of bytes released. See
273    /// [`MemoryPool::release_resident_unused_memory`].
274    pub fn release_resident_unused_memory(&self) -> usize {
275        self.memories.release_resident_unused_memory()
276    }
277
278    /// Creates a new pooling instance allocator with the given strategy and limits.
279    pub fn new(config: &PoolingAllocationConfig, tunables: &Tunables) -> Result<Self> {
280        Ok(Self {
281            live_component_instances: AtomicU64::new(0),
282            live_core_instances: AtomicU64::new(0),
283            decommit_queues: (0..default_shard_count())
284                .map(|_| CachePadded(Mutex::new(DecommitQueue::default())))
285                .try_collect::<Box<[_]>, OutOfMemory>()?,
286            memories: MemoryPool::new(config, tunables)?,
287            live_memories: AtomicUsize::new(0),
288            tables: TablePool::new(config)?,
289            live_tables: AtomicUsize::new(0),
290            #[cfg(feature = "gc")]
291            gc_heaps: if tunables.collector.is_some() {
292                Some(GcHeapPool::new(config, tunables)?)
293            } else {
294                None
295            },
296            #[cfg(feature = "gc")]
297            live_gc_heaps: AtomicUsize::new(0),
298            #[cfg(feature = "async")]
299            stacks: StackPool::new(config)?,
300            #[cfg(feature = "async")]
301            live_stacks: AtomicUsize::new(0),
302            pagemap: match config.pagemap_scan {
303                Enabled::Auto => PageMap::new(),
304                Enabled::Yes => Some(PageMap::new().ok_or_else(|| {
305                    format_err!(
306                        "required to enable PAGEMAP_SCAN but this system \
307                         does not support it"
308                    )
309                })?),
310                Enabled::No => None,
311            },
312            config: config.clone(),
313        })
314    }
315
316    fn core_instance_size(&self) -> usize {
317        round_up_to_pow2(
318            self.config.limits.core_instance_size,
319            mem::align_of::<Instance>(),
320        )
321    }
322
323    fn validate_table_plans(&self, module: &Module) -> Result<()> {
324        self.tables.validate(module)
325    }
326
327    fn validate_memory_plans(&self, module: &Module) -> Result<()> {
328        self.memories.validate_memories(module)
329    }
330
331    fn validate_core_instance_size(&self, offsets: &VMOffsets<HostPtr>) -> Result<()> {
332        let layout = Instance::alloc_layout(offsets);
333        if layout.size() <= self.core_instance_size() {
334            return Ok(());
335        }
336
337        // If this `module` exceeds the allocation size allotted to it then an
338        // error will be reported here. The error of "required N bytes but
339        // cannot allocate that" is pretty opaque, however, because it's not
340        // clear what the breakdown of the N bytes are and what to optimize
341        // next. To help provide a better error message here some fancy-ish
342        // logic is done here to report the breakdown of the byte request into
343        // the largest portions and where it's coming from.
344        let mut message = format!(
345            "instance allocation for this module \
346             requires {} bytes which exceeds the configured maximum \
347             of {} bytes; breakdown of allocation requirement:\n\n",
348            layout.size(),
349            self.core_instance_size(),
350        );
351
352        let mut remaining = layout.size();
353        let mut push = |name: &str, bytes: usize| {
354            assert!(remaining >= bytes);
355            remaining -= bytes;
356
357            // If the `name` region is more than 5% of the allocation request
358            // then report it here, otherwise ignore it. We have less than 20
359            // fields so we're guaranteed that something should be reported, and
360            // otherwise it's not particularly interesting to learn about 5
361            // different fields that are all 8 or 0 bytes. Only try to report
362            // the "major" sources of bytes here.
363            if bytes > layout.size() / 20 {
364                message.push_str(&format!(
365                    " * {:.02}% - {} bytes - {}\n",
366                    ((bytes as f32) / (layout.size() as f32)) * 100.0,
367                    bytes,
368                    name,
369                ));
370            }
371        };
372
373        // The `Instance` itself requires some size allocated to it.
374        push("instance state management", mem::size_of::<Instance>());
375
376        // Afterwards the `VMContext`'s regions are why we're requesting bytes,
377        // so ask it for descriptions on each region's byte size.
378        for (desc, size) in offsets.region_sizes() {
379            push(desc, size as usize);
380        }
381
382        // double-check we accounted for all the bytes
383        assert_eq!(remaining, 0);
384
385        bail!("{message}")
386    }
387
388    #[cfg(feature = "component-model")]
389    fn validate_component_instance_size(
390        &self,
391        offsets: &VMComponentOffsets<HostPtr>,
392        core_instances_aggregate_size: usize,
393    ) -> Result<()> {
394        let vmcomponentctx_size = usize::try_from(offsets.size_of_vmctx()).unwrap();
395        let total_instance_size = core_instances_aggregate_size.saturating_add(vmcomponentctx_size);
396        if total_instance_size <= self.config.limits.component_instance_size {
397            return Ok(());
398        }
399
400        // TODO: Add context with detailed accounting of what makes up all the
401        // `VMComponentContext`'s space like we do for module instances.
402        bail!(
403            "instance allocation for this component requires {total_instance_size} bytes of `VMComponentContext` \
404             and aggregated core instance runtime space which exceeds the configured maximum of {} bytes. \
405             `VMComponentContext` used {vmcomponentctx_size} bytes, `core module instances` used \
406             {core_instances_aggregate_size} bytes.",
407            self.config.limits.component_instance_size
408        )
409    }
410
411    /// Returns the decommit-queue shard for `shard`.
412    fn decommit_queue(&self, shard: ShardId) -> &Mutex<DecommitQueue> {
413        &self.decommit_queues[shard.index()].0
414    }
415
416    /// Enumerate all decommit-queue shard ids, starting with the current
417    /// thread's home shard.
418    fn decommit_shard_ids(&self) -> impl Iterator<Item = ShardId> {
419        shard_ids_from_home(self.decommit_queues.len())
420    }
421
422    fn flush_decommit_queue(&self, mut locked_queue: MutexGuard<'_, DecommitQueue>) -> bool {
423        // Take the queue out of the mutex and drop the lock, to minimize
424        // contention.
425        let queue = mem::take(&mut *locked_queue);
426        drop(locked_queue);
427        queue.flush(self)
428    }
429
430    /// Flush every shard of the decommit queue, e.g. on allocator drop.
431    /// Returns whether any slot was returned to any pool.
432    fn flush_all_decommit_queues(&self) -> bool {
433        let mut any = false;
434        for shard in self.decommit_shard_ids() {
435            let queue = self.decommit_queue(shard).lock().unwrap();
436            any |= self.flush_decommit_queue(queue);
437        }
438        any
439    }
440
441    /// Execute `f` and if it returns `Err(PoolConcurrencyLimitError)`, then try
442    /// flushing the decommit queue. If flushing the queue freed up slots, then
443    /// try running `f` again.
444    ///
445    /// Queue shards are flushed one at a time, retrying `f` after each flush
446    /// that returned slots to a pool, rather than eagerly flushing all
447    /// shards: one flushed shard is often enough to satisfy the allocation,
448    /// and this avoids acquiring every shard's lock (at the cost of raising
449    /// the chances that another thread steals the freshly-flushed slots
450    /// before we get a chance to grab one, in which case we keep flushing).
451    ///
452    /// Note that [`Self::flush_decommit_queue`] takes the shard's queue out
453    /// of its mutex and drops the lock immediately, so no queue lock is held
454    /// while decommitting or while `f` runs.
455    #[cfg(feature = "async")]
456    fn with_flush_and_retry<T>(&self, mut f: impl FnMut() -> Result<T>) -> Result<T> {
457        let mut result = f();
458        for shard in self.decommit_shard_ids() {
459            match &result {
460                Err(e) if e.is::<PoolConcurrencyLimitError>() => {}
461                _ => break,
462            }
463            let queue = self.decommit_queue(shard).lock().unwrap();
464            if self.flush_decommit_queue(queue) {
465                result = f();
466            }
467        }
468        result
469    }
470
471    fn merge_or_flush(&self, mut local_queue: DecommitQueue) {
472        match local_queue.raw_len() {
473            // If we didn't enqueue any regions for decommit, then we must have
474            // either memset the whole entity or eagerly remapped it to zero
475            // because we don't have linux's `madvise(DONTNEED)` semantics. In
476            // either case, the entity slot is ready for reuse immediately.
477            0 => {
478                local_queue.flush(self);
479            }
480
481            // We enqueued at least our batch size of regions for decommit, so
482            // flush the local queue immediately. Don't bother inspecting (or
483            // locking!) the shared queue.
484            n if n >= self.config.decommit_batch_size => {
485                local_queue.flush(self);
486            }
487
488            // If we enqueued some regions for decommit, but did not reach our
489            // batch size, so we don't want to flush it yet, then merge the
490            // local queue into this thread's shard of the shared queue.
491            n => {
492                debug_assert!(n < self.config.decommit_batch_size);
493                let shard = thread_shard(self.decommit_queues.len());
494                let mut shared_queue = self.decommit_queue(shard).lock().unwrap();
495                shared_queue.append(&mut local_queue);
496                // And if this shard now has at least as many regions enqueued
497                // for decommit as our batch size, then we can flush it.
498                if shared_queue.raw_len() >= self.config.decommit_batch_size {
499                    self.flush_decommit_queue(shared_queue);
500                }
501            }
502        }
503    }
504
505    pub fn config(&self) -> &PoolingAllocationConfig {
506        &self.config
507    }
508}
509
510unsafe impl InstanceAllocator for PoolingInstanceAllocator {
511    #[cfg(feature = "component-model")]
512    fn validate_component<'a>(
513        &self,
514        component: &Component,
515        offsets: &VMComponentOffsets<HostPtr>,
516        get_module: &'a dyn Fn(StaticModuleIndex) -> &'a Module,
517    ) -> Result<()> {
518        let mut num_core_instances = 0;
519        let mut num_memories = 0;
520        let mut num_tables = 0;
521        let mut core_instances_aggregate_size: usize = 0;
522        for init in &component.initializers {
523            use wasmtime_environ::component::GlobalInitializer::*;
524            use wasmtime_environ::component::InstantiateModule;
525            match init {
526                InstantiateModule(InstantiateModule::Import(_, _), _) => {
527                    num_core_instances += 1;
528                    // Can't statically account for the total vmctx size, number
529                    // of memories, and number of tables in this component.
530                }
531                InstantiateModule(InstantiateModule::Static(static_module_index, _), _) => {
532                    let module = get_module(*static_module_index);
533                    let offsets = VMOffsets::new(HostPtr, &module);
534                    let layout = Instance::alloc_layout(&offsets);
535                    self.validate_module(module, &offsets)?;
536                    num_core_instances += 1;
537                    num_memories += module.num_defined_memories();
538                    num_tables += module.num_defined_tables();
539                    core_instances_aggregate_size += layout.size();
540                }
541                LowerImport { .. }
542                | ExtractMemory(_)
543                | ExtractTable(_)
544                | ExtractRealloc(_)
545                | ExtractCallback(_)
546                | ExtractPostReturn(_)
547                | Resource(_) => {}
548            }
549        }
550
551        if num_core_instances
552            > usize::try_from(self.config.limits.max_core_instances_per_component).unwrap()
553        {
554            bail!(
555                "The component transitively contains {num_core_instances} core module instances, \
556                 which exceeds the configured maximum of {} in the pooling allocator",
557                self.config.limits.max_core_instances_per_component
558            );
559        }
560
561        if num_memories > usize::try_from(self.config.limits.max_memories_per_component).unwrap() {
562            bail!(
563                "The component transitively contains {num_memories} Wasm linear memories, which \
564                 exceeds the configured maximum of {} in the pooling allocator",
565                self.config.limits.max_memories_per_component
566            );
567        }
568
569        if num_tables > usize::try_from(self.config.limits.max_tables_per_component).unwrap() {
570            bail!(
571                "The component transitively contains {num_tables} tables, which exceeds the \
572                 configured maximum of {} in the pooling allocator",
573                self.config.limits.max_tables_per_component
574            );
575        }
576
577        self.validate_component_instance_size(offsets, core_instances_aggregate_size)
578            .context("component instance size does not fit in pooling allocator requirements")?;
579
580        Ok(())
581    }
582
583    fn validate_module(&self, module: &Module, offsets: &VMOffsets<HostPtr>) -> Result<()> {
584        self.validate_memory_plans(module)
585            .context("module memory does not fit in pooling allocator requirements")?;
586        self.validate_table_plans(module)
587            .context("module table does not fit in pooling allocator requirements")?;
588        self.validate_core_instance_size(offsets)
589            .context("module instance size does not fit in pooling allocator requirements")?;
590        Ok(())
591    }
592
593    #[cfg(feature = "gc")]
594    fn validate_memory(&self, memory: &wasmtime_environ::Memory) -> Result<()> {
595        self.memories.validate_memory(memory)
596    }
597
598    #[cfg(feature = "component-model")]
599    fn increment_component_instance_count(&self) -> Result<()> {
600        let old_count = self.live_component_instances.fetch_add(1, Ordering::AcqRel);
601        if old_count >= u64::from(self.config.limits.total_component_instances) {
602            self.decrement_component_instance_count();
603            return Err(PoolConcurrencyLimitError::new(
604                usize::try_from(self.config.limits.total_component_instances).unwrap(),
605                "component instances",
606            )
607            .into());
608        }
609        Ok(())
610    }
611
612    #[cfg(feature = "component-model")]
613    fn decrement_component_instance_count(&self) {
614        self.live_component_instances.fetch_sub(1, Ordering::AcqRel);
615    }
616
617    fn increment_core_instance_count(&self) -> Result<()> {
618        let old_count = self.live_core_instances.fetch_add(1, Ordering::AcqRel);
619        if old_count >= u64::from(self.config.limits.total_core_instances) {
620            self.decrement_core_instance_count();
621            return Err(PoolConcurrencyLimitError::new(
622                usize::try_from(self.config.limits.total_core_instances).unwrap(),
623                "core instances",
624            )
625            .into());
626        }
627        Ok(())
628    }
629
630    fn decrement_core_instance_count(&self) {
631        self.live_core_instances.fetch_sub(1, Ordering::AcqRel);
632    }
633
634    fn allocate_memory<'a, 'b: 'a, 'c: 'a>(
635        &'a self,
636        request: &'a mut InstanceAllocationRequest<'b, 'c>,
637        ty: &'a wasmtime_environ::Memory,
638        memory_index: Option<DefinedMemoryIndex>,
639        _memory_kind: MemoryKind,
640    ) -> Pin<Box<dyn Future<Output = Result<(MemoryAllocationIndex, Memory)>> + Send + 'a>> {
641        crate::runtime::box_future(async move {
642            async {
643                // FIXME(rust-lang/rust#145127) this should ideally use a version of
644                // `with_flush_and_retry` but adapted for async closures instead of only
645                // sync closures. Right now that won't compile though so this is the
646                // manually expanded version of the method.
647                let mut e = match self.memories.allocate(request, ty, memory_index).await {
648                    Ok(result) => return Ok(result),
649                    Err(e) => e,
650                };
651
652                for shard in self.decommit_shard_ids() {
653                    if !e.is::<PoolConcurrencyLimitError>() {
654                        break;
655                    }
656                    let queue = self.decommit_queue(shard).lock().unwrap();
657                    if self.flush_decommit_queue(queue) {
658                        match self.memories.allocate(request, ty, memory_index).await {
659                            Ok(result) => return Ok(result),
660                            Err(err) => e = err,
661                        }
662                    }
663                }
664
665                Err(e)
666            }
667            .await
668            .inspect(|_| {
669                self.live_memories.fetch_add(1, Ordering::Relaxed);
670            })
671        })
672    }
673
674    unsafe fn deallocate_memory(
675        &self,
676        _memory_index: Option<DefinedMemoryIndex>,
677        allocation_index: MemoryAllocationIndex,
678        memory: Memory,
679    ) {
680        let prev = self.live_memories.fetch_sub(1, Ordering::Relaxed);
681        debug_assert!(prev > 0);
682
683        // Reset the image slot. Depending on whether this is successful or not
684        // the `image` is preserved for future use. On success it's queued up to
685        // get deallocated later, and on failure the slot is deallocated
686        // immediately without preserving the image.
687        let mut image = memory.unwrap_static_image();
688        let mut queue = DecommitQueue::default();
689        let bytes_resident = image.clear_and_remain_ready(
690            self.pagemap.as_ref(),
691            self.memories.keep_resident,
692            |ptr, len| {
693                // SAFETY: the memory in `image` won't be used until this
694                // decommit queue is flushed, and by definition the memory is
695                // not in use when calling this function.
696                unsafe {
697                    queue.push_raw(ptr, len);
698                }
699            },
700        );
701
702        match bytes_resident {
703            Ok(bytes_resident) => {
704                // SAFETY: this image is not in use and its memory regions were enqueued
705                // with `push_raw` above.
706                unsafe {
707                    queue.push_memory(allocation_index, image, bytes_resident);
708                }
709                self.merge_or_flush(queue);
710            }
711            Err(e) => {
712                log::warn!("ignoring clear_and_remain_ready error {e}");
713                // SAFETY: `allocation_index` comes from this pool, as an unsafe
714                // contract of this function itself, and it's guaranteed to be no
715                // longer in use so safe to deallocate. The slot couldn't be
716                // preserved so it's dropped here.
717                //
718                // Note that at this point it's not clear how many bytes are
719                // resident in memory, so it's inevitably going to leave statistics
720                // a little off. Also note though that non-Linux platforms don't
721                // keep track of resident bytes anyway, and this path is only
722                // reachable on non-Linux platforms because Linux can't return an
723                // error.
724                unsafe {
725                    self.memories.deallocate(allocation_index, None, 0);
726                }
727            }
728        }
729    }
730
731    fn allocate_table<'a, 'b: 'a, 'c: 'a>(
732        &'a self,
733        request: &'a mut InstanceAllocationRequest<'b, 'c>,
734        ty: &'a wasmtime_environ::Table,
735        _table_index: DefinedTableIndex,
736    ) -> Pin<Box<dyn Future<Output = Result<(super::TableAllocationIndex, Table)>> + Send + 'a>>
737    {
738        crate::runtime::box_future(async move {
739            async {
740                // FIXME: see `allocate_memory` above for comments about duplication
741                // with `with_flush_and_retry`.
742                let mut e = match self.tables.allocate(request, ty).await {
743                    Ok(result) => return Ok(result),
744                    Err(e) => e,
745                };
746
747                for shard in self.decommit_shard_ids() {
748                    if !e.is::<PoolConcurrencyLimitError>() {
749                        break;
750                    }
751                    let queue = self.decommit_queue(shard).lock().unwrap();
752                    if self.flush_decommit_queue(queue) {
753                        match self.tables.allocate(request, ty).await {
754                            Ok(result) => return Ok(result),
755                            Err(err) => e = err,
756                        }
757                    }
758                }
759
760                Err(e)
761            }
762            .await
763            .inspect(|_| {
764                self.live_tables.fetch_add(1, Ordering::Relaxed);
765            })
766        })
767    }
768
769    unsafe fn deallocate_table(
770        &self,
771        _table_index: DefinedTableIndex,
772        allocation_index: TableAllocationIndex,
773        mut table: Table,
774    ) {
775        let prev = self.live_tables.fetch_sub(1, Ordering::Relaxed);
776        debug_assert!(prev > 0);
777
778        let mut queue = DecommitQueue::default();
779        // SAFETY: This table is no longer in use by the allocator when this
780        // method is called and additionally all image ranges are pushed with
781        // the understanding that the memory won't get used until the whole
782        // queue is flushed.
783        let bytes_resident = unsafe {
784            self.tables.reset_table_pages_to_zero(
785                self.pagemap.as_ref(),
786                allocation_index,
787                &mut table,
788                |ptr, len| {
789                    queue.push_raw(ptr, len);
790                },
791            )
792        };
793
794        // SAFETY: the table has had all its memory regions enqueued above.
795        unsafe {
796            queue.push_table(allocation_index, table, bytes_resident);
797        }
798        self.merge_or_flush(queue);
799    }
800
801    #[cfg(feature = "async")]
802    fn allocate_fiber_stack(&self) -> Result<wasmtime_fiber::FiberStack> {
803        let ret = self.with_flush_and_retry(|| self.stacks.allocate())?;
804        self.live_stacks.fetch_add(1, Ordering::Relaxed);
805        Ok(ret)
806    }
807
808    #[cfg(feature = "async")]
809    unsafe fn deallocate_fiber_stack(&self, mut stack: wasmtime_fiber::FiberStack) {
810        self.live_stacks.fetch_sub(1, Ordering::Relaxed);
811        let mut queue = DecommitQueue::default();
812        // SAFETY: the stack is no longer in use by definition when this
813        // function is called and memory ranges pushed here are otherwise no
814        // longer in use.
815        let bytes_resident = unsafe {
816            self.stacks
817                .zero_stack(&mut stack, |ptr, len| queue.push_raw(ptr, len))
818        };
819        // SAFETY: this stack's memory regions were enqueued above.
820        unsafe {
821            queue.push_stack(stack, bytes_resident);
822        }
823        self.merge_or_flush(queue);
824    }
825
826    fn purge_module(&self, module: CompiledModuleId) {
827        self.memories.purge_module(module);
828    }
829
830    fn next_available_pkey(&self) -> Option<ProtectionKey> {
831        self.memories.next_available_pkey()
832    }
833
834    fn restrict_to_pkey(&self, pkey: ProtectionKey) {
835        mpk::allow(ProtectionMask::zero().or(pkey));
836    }
837
838    fn allow_all_pkeys(&self) {
839        mpk::allow(ProtectionMask::all());
840    }
841
842    #[cfg(feature = "gc")]
843    fn allocate_gc_heap(
844        &self,
845        engine: &crate::Engine,
846        gc_runtime: &dyn GcRuntime,
847        memory_alloc_index: MemoryAllocationIndex,
848    ) -> Result<(GcHeapAllocationIndex, Box<dyn GcHeap>)> {
849        let ret =
850            self.gc_heaps
851                .as_ref()
852                .unwrap()
853                .allocate(engine, gc_runtime, memory_alloc_index)?;
854        self.live_gc_heaps.fetch_add(1, Ordering::Relaxed);
855        Ok(ret)
856    }
857
858    #[cfg(feature = "gc")]
859    fn deallocate_gc_heap(
860        &self,
861        allocation_index: GcHeapAllocationIndex,
862        gc_heap: Box<dyn GcHeap>,
863    ) -> MemoryAllocationIndex {
864        let gc_heaps = self.gc_heaps.as_ref().unwrap();
865        self.live_gc_heaps.fetch_sub(1, Ordering::Relaxed);
866        gc_heaps.deallocate(allocation_index, gc_heap)
867    }
868
869    fn as_pooling(&self) -> Option<&PoolingInstanceAllocator> {
870        Some(self)
871    }
872}
873
874#[cfg(test)]
875#[cfg(target_pointer_width = "64")]
876mod test {
877    use super::*;
878    use crate::config::InstanceLimits;
879
880    #[test]
881    fn test_pooling_allocator_with_memory_pages_exceeded() {
882        let config = PoolingAllocationConfig {
883            limits: InstanceLimits {
884                total_memories: 1,
885                max_memory_size: 0x100010000,
886                ..Default::default()
887            },
888            ..PoolingAllocationConfig::default()
889        };
890        assert_eq!(
891            PoolingInstanceAllocator::new(
892                &config,
893                &Tunables {
894                    memory_reservation: 0x10000,
895                    ..Tunables::default_host()
896                },
897            )
898            .map_err(|e| e.to_string())
899            .expect_err("expected a failure constructing instance allocator"),
900            "maximum memory size of 0x100010000 bytes exceeds the configured \
901             memory reservation of 0x10000 bytes"
902        );
903    }
904
905    #[cfg(all(
906        unix,
907        target_pointer_width = "64",
908        feature = "async",
909        not(miri),
910        not(asan)
911    ))]
912    #[test]
913    fn test_stack_zeroed() -> Result<()> {
914        let config = PoolingAllocationConfig {
915            max_unused_warm_slots: 0,
916            limits: InstanceLimits {
917                total_stacks: 1,
918                total_memories: 0,
919                total_tables: 0,
920                ..Default::default()
921            },
922            stack_size: 128,
923            async_stack_zeroing: true,
924            ..PoolingAllocationConfig::default()
925        };
926        let allocator = PoolingInstanceAllocator::new(&config, &Tunables::default_host())?;
927
928        unsafe {
929            for _ in 0..255 {
930                let stack = allocator.allocate_fiber_stack()?;
931
932                // The stack pointer is at the top, so decrement it first
933                let addr = stack.top().unwrap().sub(1);
934
935                assert_eq!(*addr, 0);
936                *addr = 1;
937
938                allocator.deallocate_fiber_stack(stack);
939            }
940        }
941
942        Ok(())
943    }
944
945    #[cfg(all(
946        unix,
947        target_pointer_width = "64",
948        feature = "async",
949        not(miri),
950        not(asan)
951    ))]
952    #[test]
953    fn test_stack_unzeroed() -> Result<()> {
954        let config = PoolingAllocationConfig {
955            max_unused_warm_slots: 0,
956            limits: InstanceLimits {
957                total_stacks: 1,
958                total_memories: 0,
959                total_tables: 0,
960                ..Default::default()
961            },
962            stack_size: 128,
963            async_stack_zeroing: false,
964            ..PoolingAllocationConfig::default()
965        };
966        let allocator = PoolingInstanceAllocator::new(&config, &Tunables::default_host())?;
967
968        unsafe {
969            for i in 0..255 {
970                let stack = allocator.allocate_fiber_stack()?;
971
972                // The stack pointer is at the top, so decrement it first
973                let addr = stack.top().unwrap().sub(1);
974
975                assert_eq!(*addr, i);
976                *addr = i + 1;
977
978                allocator.deallocate_fiber_stack(stack);
979            }
980        }
981
982        Ok(())
983    }
984}