Skip to content

vllm.distributed.device_communicators.shm_object_storage

Classes:

ObjectSerde

Bases: ABC

Methods:

Source code in vllm/distributed/device_communicators/shm_object_storage.py
class ObjectSerde(ABC):
    @abstractmethod
    def serialize(self, value: Any) -> tuple[Any, int, bytes, int]:
        """Serialize an object to bytes."""
        raise NotImplementedError

    @abstractmethod
    def deserialize(self, data: memoryview) -> Any:
        """Deserialize bytes back to an object."""
        raise NotImplementedError

deserialize(data) abstractmethod

Deserialize bytes back to an object.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
@abstractmethod
def deserialize(self, data: memoryview) -> Any:
    """Deserialize bytes back to an object."""
    raise NotImplementedError

serialize(value) abstractmethod

Serialize an object to bytes.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
@abstractmethod
def serialize(self, value: Any) -> tuple[Any, int, bytes, int]:
    """Serialize an object to bytes."""
    raise NotImplementedError

SingleWriterShmObjectStorage

A single-writer, multiple-reader object storage system built on top of a shared memory ring buffer. Provides key-value storage with automatic memory management and cross-process serialization support.

This storage system follows a FIFO (First-In-First-Out) eviction policy where the oldest objects are automatically freed when memory runs low. Memory is reclaimed based on reader reference counting - objects are only freed when all readers have finished accessing them.

Architecture: - Single writer process can put(key, value) objects - Multiple reader processes can get(address, monotonic_id, signature, key) objects - Built on SingleWriterShmRingBuffer for efficient shared memory management - Thread-safe operations with reader synchronization via locks

Key Features: - FIFO Eviction: Oldest objects are evicted first when memory is full - Reference Counting: Objects are only freed when no readers are accessing them - Duplicate Key Handling: Existing keys are not overwritten, just re-referenced - Customized Serialization: By default uses Msgpack for efficient serialization of Python objects, but can be extended for custom types - Cross-Process Safety: Uses shared memory with proper synchronization - Automatic Cleanup: Garbage collection happens transparently during allocation

Handle Signatures: - Handles are signed with an HMAC of (key, address, monotonic_id), keyed by a random secret stored at the start of the shared memory - The writer signs the handles it issues and rejects stale or forged handles supplied by requests - Readers verify the signature in get/touch before reading shared memory

Memory Layout per Object: [4-byte reference_count][metadata_size][serialized_object_data]

Thread Safety: - Writer operations (put, clear) are single-threaded by design - Reader operations (get) are thread-safe with lock-based reference counting - Memory reclamation is handled exclusively by the writer process

Methods:

  • __init__ –

    Initialize the object storage.

  • clear –

    Clear the object storage.

  • close –

    Close the shared memory.

  • default_is_free_check –

    Default is_free function that checks if the first 4 bytes are zero.

  • free_unused –

    Free unused buffers in the ring buffer.

  • get_cached –

    Get the cached object by key if it exists.

  • get_signature –

    Sign the handle of a cached object so readers can verify it.

  • handle –

    Get handle for sharing across processes.

  • increment_reader_flag –

    Set the in-use flag for the reader.

  • increment_writer_flag –

    Set the in-use flag for the writer.

  • is_cached –

    Check if the object with the given key is cached.

  • put –

    Store a key-value pair in the object storage.

  • release_touches –

    Stop protecting touched items that were not looked up.

  • touch –

    Touch an existing cached item.

  • verify_signature –

    Verify that a handle was issued for the given key.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
class SingleWriterShmObjectStorage:
    """A single-writer, multiple-reader object storage system built on top of a
    shared memory ring buffer. Provides key-value storage with automatic memory
    management and cross-process serialization support.

    This storage system follows a FIFO (First-In-First-Out) eviction policy
    where the oldest objects are automatically freed when memory runs low.
    Memory is reclaimed based on reader reference counting - objects are only
    freed when all readers have finished accessing them.

    Architecture:
    - Single writer process can put(key, value) objects
    - Multiple reader processes can get(address, monotonic_id, signature, key)
      objects
    - Built on SingleWriterShmRingBuffer for efficient shared memory management
    - Thread-safe operations with reader synchronization via locks

    Key Features:
    - FIFO Eviction: Oldest objects are evicted first when memory is full
    - Reference Counting: Objects are only freed when no readers are
      accessing them
    - Duplicate Key Handling: Existing keys are not overwritten, just
      re-referenced
    - Customized Serialization: By default uses Msgpack for efficient
      serialization of Python objects, but can be extended for custom types
    - Cross-Process Safety: Uses shared memory with proper synchronization
    - Automatic Cleanup: Garbage collection happens transparently during
      allocation

    Handle Signatures:
    - Handles are signed with an HMAC of (key, address, monotonic_id), keyed by
      a random secret stored at the start of the shared memory
    - The writer signs the handles it issues and rejects stale or forged
      handles supplied by requests
    - Readers verify the signature in get/touch before reading shared memory

    Memory Layout per Object:
    `[4-byte reference_count][metadata_size][serialized_object_data]`

    Thread Safety:
    - Writer operations (put, clear) are single-threaded by design
    - Reader operations (get) are thread-safe with lock-based reference
      counting
    - Memory reclamation is handled exclusively by the writer process
    """

    def __init__(
        self,
        max_object_size: int,
        n_readers: int,
        ring_buffer: SingleWriterShmRingBuffer,
        serde_class: type[ObjectSerde] = MsgpackSerde,
        reader_lock: LockType | None = None,
    ):
        """Initialize the object storage.

        Args:
            max_object_size: Maximum size for a single object in bytes.
            n_readers: Number of reader processes that can access the storage.
            ring_buffer: The shared memory ring buffer for storing objects.
            serde_class: Serializer/deserializer for objects.
            reader_lock: Optional lock for synchronizing reader access.

        Raises:
            ValueError: If reader_lock is None for readers.

        """
        self.max_object_size = max_object_size
        self.n_readers = n_readers
        self.serde_class = serde_class
        self.ser_de = serde_class()
        self.ring_buffer = ring_buffer
        self.is_writer = self.ring_buffer.is_writer

        self.flag_bytes = 4  # for in-use flag

        if self.is_writer:
            # Key-value mapping: key -> (address, monotonic_id)
            self.key_index: dict[str, tuple[int, int]] = {}
            # Reverse mapping: monotonic_id -> key
            self.id_index: dict[int, str] = {}
            # Writer flag to track in-use status: monotonic_id -> count
            self.writer_flag: dict[int, int] = {}
            # Items touch() keeps from being freed until they are looked up
            # or release_touches() is called
            self._touched: set[int] = set()
        else:
            if reader_lock is None:
                raise ValueError("Lock must be provided for readers.")

        self._reader_lock = reader_lock

    def clear(self) -> None:
        """Clear the object storage."""
        if self.is_writer:
            self.ring_buffer.clear()
            self.key_index.clear()
            self.id_index.clear()
            self.writer_flag.clear()
            self._touched.clear()
            logger.debug("Object storage cleared and reinitialized.")

    def copy_to_buffer(
        self,
        data: bytes | list[bytes],
        data_bytes: int,
        metadata: bytes,
        md_bytes: int,
        data_view: memoryview,
    ) -> None:
        data_view[self.flag_bytes : self.flag_bytes + md_bytes] = metadata
        if isinstance(data, bytes):
            data_view[-data_bytes:] = data
        elif isinstance(data, list):
            start_idx = self.flag_bytes + md_bytes
            for item_bytes in data:
                item_size = len(item_bytes)
                data_view[start_idx : start_idx + item_size] = item_bytes
                start_idx += item_size
        else:
            raise ValueError(f"Unsupported data type for serialization: {type(data)}")

    def _validate_monotonic_id(
        self,
        address: int,
        monotonic_id: int,
        buf_metadata: tuple[int, int],
    ) -> None:
        if buf_metadata[0] != monotonic_id:
            raise ValueError(
                f"Data for address:id '{address}:{monotonic_id}'"
                " has been modified or is invalid."
            )

    def _verify_reader_signature(
        self,
        address: int,
        monotonic_id: int,
        signature: list[int] | None,
        key: str | None,
    ) -> None:
        if signature is None or key is None:
            raise ValueError("Missing SHM handle signature for cache key.")
        if (
            not isinstance(signature, list)
            or len(signature) != hashlib.sha256().digest_size
            or any(type(byte) is not int or not 0 <= byte <= 255 for byte in signature)
            or not isinstance(key, str)
        ):
            raise ValueError("Invalid SHM handle signature for cache key.")

        expected_signature = self._make_signature(key, address, monotonic_id)
        if not hmac.compare_digest(bytes(signature), expected_signature):
            raise ValueError("Invalid SHM handle signature for cache key.")

    def verify_signature(
        self,
        key: str | None,
        address: int,
        monotonic_id: int,
        signature: list[int] | None,
    ) -> None:
        """Verify that a handle was issued for the given key.

        For writers: the handle must also be the key's current entry
        For readers: only the signature is checked

        Args:
            key: String key the handle belongs to
            address: Address of the object
            monotonic_id: Monotonic ID of the object
            signature: Signature issued with the handle

        """
        if self.is_writer and (
            key is None or self.key_index.get(key) != (address, monotonic_id)
        ):
            raise ValueError("Invalid SHM handle signature for cache key.")
        self._verify_reader_signature(address, monotonic_id, signature, key)

    def _make_signature(self, key: str, address: int, monotonic_id: int) -> bytes:
        payload = (
            hashlib.sha256(key.encode("utf-8")).digest()
            + address.to_bytes(8, "little", signed=False)
            + monotonic_id.to_bytes(4, "little", signed=False)
        )
        return hmac.new(
            self.ring_buffer.get_auth_secret(),
            payload,
            hashlib.sha256,
        ).digest()

    def increment_writer_flag(self, id: int) -> None:
        """Set the in-use flag for the writer."""
        self.writer_flag[id] = self.writer_flag.get(id, 0) + 1

    def increment_reader_flag(self, data_view: memoryview) -> None:
        """Set the in-use flag for the reader."""
        # >0 for in-use flag
        reader_count = self.ring_buffer.byte2int(data_view)
        data_view[:] = self.ring_buffer.int2byte(reader_count + 1)

    def free_unused(self) -> None:
        """Free unused buffers in the ring buffer."""
        # try to free up 2*max_object_size bytes of space in the ring buffer,
        # since the buffer might be fragmented
        freed_ids = self.ring_buffer.free_buf(
            self.default_is_free_check, 2 * self.max_object_size
        )
        # update the metadata after freeing up space
        for freed_id in freed_ids:
            key_to_free = self.id_index[freed_id]
            del self.key_index[key_to_free]
            del self.id_index[freed_id]
            del self.writer_flag[freed_id]

    def is_cached(self, key: str) -> bool:
        """Check if the object with the given key is cached."""
        return key in self.key_index

    def get_cached(self, key: str) -> tuple[int, int]:
        """Get the cached object by key if it exists."""
        address, monotonic_id = self.key_index[key]
        self.increment_writer_flag(monotonic_id)
        self._touched.discard(monotonic_id)
        return address, monotonic_id

    def release_touches(self) -> None:
        """Stop protecting touched items that were not looked up."""
        self._touched.clear()

    def get_signature(self, key: str) -> list[int]:
        """Sign the handle of a cached object so readers can verify it."""
        address, monotonic_id = self.key_index[key]
        return list(self._make_signature(key, address, monotonic_id))

    def put(self, key: str, value: Any) -> tuple[int, int]:
        """Store a key-value pair in the object storage.
        Attempts to free max_object_size bytes using FIFO order
        when the ring buffer runs out of space during a put() operation.

        Args:
            key: String key to identify the object
            value: Any serializable Python object

        Raises:
            MemoryError: If there's not enough space in the buffer
            ValueError: If the serialized object is too large
            ValueError: If the key already exists in the storage

        """
        if key in self.key_index:
            raise ValueError(f"Key '{key}' already exists in the storage.")

        object_data, data_bytes, object_metadata, md_bytes = self.ser_de.serialize(
            value
        )
        buffer_size = self.flag_bytes + data_bytes + md_bytes
        # Sanity checks
        if buffer_size > self.max_object_size:
            raise ValueError(
                f"Serialized object size ({buffer_size} bytes) exceeds "
                f"max object size ({self.max_object_size} bytes)"
            )

        # Allocate new buffer
        try:
            address, monotonic_id = self.ring_buffer.allocate_buf(buffer_size)
        except MemoryError:
            self.free_unused()
            # try again after freeing up space
            address, monotonic_id = self.ring_buffer.allocate_buf(buffer_size)

        # Write data to buffer
        with self.ring_buffer.access_buf(address) as (data_view, metadata):
            data_view[: self.flag_bytes] = self.ring_buffer.int2byte(0)
            self.copy_to_buffer(
                object_data, data_bytes, object_metadata, md_bytes, data_view
            )
        self.increment_writer_flag(monotonic_id)

        # Update key index
        self.key_index[key] = (address, monotonic_id)
        self.id_index[monotonic_id] = key
        return address, monotonic_id

    def get(
        self,
        address: int,
        monotonic_id: int,
        signature: list[int] | None = None,
        key: str | None = None,
    ) -> Any:
        if not self.is_writer:
            # Reject forged handles before dereferencing the supplied address.
            self.verify_signature(key, address, monotonic_id, signature)

        # Read data from buffer
        with self.ring_buffer.access_buf(address) as (data_view, buf_metadata):
            self._validate_monotonic_id(address, monotonic_id, buf_metadata)

            obj = self.ser_de.deserialize(data_view[self.flag_bytes :])

            # decrease the in-use flag for reader reads
            if self._reader_lock is not None:
                with self._reader_lock:
                    self.increment_reader_flag(data_view[: self.flag_bytes])
            else:
                # if self._reader_lock is None, it means we are the writer
                # in this case, we do not need to decrease the reader count
                assert self.is_writer

        return obj

    def touch(
        self,
        key: str,
        address: int = 0,
        monotonic_id: int = 0,
        signature: list[int] | None = None,
    ) -> None:
        """Touch an existing cached item.

        For writers (ShmObjectStoreSenderCache): Protect the item from
        eviction until get_cached or release_touches
        For readers (ShmObjectStoreReceiverCache): Validate the handle

        Args:
            key: String key of the object to touch
            address: Address of the object (only for readers)
            monotonic_id: Monotonic ID of the object (only for readers)
            signature: Server-issued handle signature (only for readers)

        """
        if self._reader_lock is None:
            if key not in self.key_index:
                return None
            self._touched.add(self.key_index[key][1])
        else:
            # Reject forged handles before dereferencing the supplied address.
            self.verify_signature(key, address, monotonic_id, signature)
            with self.ring_buffer.access_buf(address) as (_, buf_metadata):
                self._validate_monotonic_id(address, monotonic_id, buf_metadata)

    def close(self) -> None:
        """Close the shared memory."""
        self.ring_buffer.close()

    def handle(self):
        """Get handle for sharing across processes."""
        return ShmObjectStorageHandle(
            max_object_size=self.max_object_size,
            n_readers=self.n_readers,
            ring_buffer_handle=self.ring_buffer.handle(),
            serde_class=self.serde_class,
            reader_lock=self._reader_lock,
        )

    @staticmethod
    def create_from_handle(
        handle: ShmObjectStorageHandle,
    ) -> "SingleWriterShmObjectStorage":
        logger.debug("Creating storage from handle: %s", handle)
        ring_buffer = SingleWriterShmRingBuffer(*handle.ring_buffer_handle)
        return SingleWriterShmObjectStorage(
            max_object_size=handle.max_object_size,
            n_readers=handle.n_readers,
            ring_buffer=ring_buffer,
            serde_class=handle.serde_class,
            reader_lock=handle.reader_lock,
        )

    def default_is_free_check(self, id: int, buf: memoryview) -> bool:
        """Default is_free function that checks if the first 4 bytes are zero.
        This indicates that the buffer is free.
        """
        if id in self._touched:
            return False
        reader_count = int.from_bytes(buf[0:4], "little", signed=True)
        writer_count = self.writer_flag[id]
        return reader_count >= writer_count * self.n_readers

__init__(max_object_size, n_readers, ring_buffer, serde_class=MsgpackSerde, reader_lock=None)

Initialize the object storage.

Parameters:

  • max_object_size

    (int) –

    Maximum size for a single object in bytes.

  • n_readers

    (int) –

    Number of reader processes that can access the storage.

  • ring_buffer

    (SingleWriterShmRingBuffer) –

    The shared memory ring buffer for storing objects.

  • serde_class

    (type[ObjectSerde], default: MsgpackSerde ) –

    Serializer/deserializer for objects.

  • reader_lock

    (Lock | None, default: None ) –

    Optional lock for synchronizing reader access.

Raises:

  • ValueError –

    If reader_lock is None for readers.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def __init__(
    self,
    max_object_size: int,
    n_readers: int,
    ring_buffer: SingleWriterShmRingBuffer,
    serde_class: type[ObjectSerde] = MsgpackSerde,
    reader_lock: LockType | None = None,
):
    """Initialize the object storage.

    Args:
        max_object_size: Maximum size for a single object in bytes.
        n_readers: Number of reader processes that can access the storage.
        ring_buffer: The shared memory ring buffer for storing objects.
        serde_class: Serializer/deserializer for objects.
        reader_lock: Optional lock for synchronizing reader access.

    Raises:
        ValueError: If reader_lock is None for readers.

    """
    self.max_object_size = max_object_size
    self.n_readers = n_readers
    self.serde_class = serde_class
    self.ser_de = serde_class()
    self.ring_buffer = ring_buffer
    self.is_writer = self.ring_buffer.is_writer

    self.flag_bytes = 4  # for in-use flag

    if self.is_writer:
        # Key-value mapping: key -> (address, monotonic_id)
        self.key_index: dict[str, tuple[int, int]] = {}
        # Reverse mapping: monotonic_id -> key
        self.id_index: dict[int, str] = {}
        # Writer flag to track in-use status: monotonic_id -> count
        self.writer_flag: dict[int, int] = {}
        # Items touch() keeps from being freed until they are looked up
        # or release_touches() is called
        self._touched: set[int] = set()
    else:
        if reader_lock is None:
            raise ValueError("Lock must be provided for readers.")

    self._reader_lock = reader_lock

clear()

Clear the object storage.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def clear(self) -> None:
    """Clear the object storage."""
    if self.is_writer:
        self.ring_buffer.clear()
        self.key_index.clear()
        self.id_index.clear()
        self.writer_flag.clear()
        self._touched.clear()
        logger.debug("Object storage cleared and reinitialized.")

close()

Close the shared memory.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def close(self) -> None:
    """Close the shared memory."""
    self.ring_buffer.close()

default_is_free_check(id, buf)

Default is_free function that checks if the first 4 bytes are zero. This indicates that the buffer is free.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def default_is_free_check(self, id: int, buf: memoryview) -> bool:
    """Default is_free function that checks if the first 4 bytes are zero.
    This indicates that the buffer is free.
    """
    if id in self._touched:
        return False
    reader_count = int.from_bytes(buf[0:4], "little", signed=True)
    writer_count = self.writer_flag[id]
    return reader_count >= writer_count * self.n_readers

free_unused()

Free unused buffers in the ring buffer.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def free_unused(self) -> None:
    """Free unused buffers in the ring buffer."""
    # try to free up 2*max_object_size bytes of space in the ring buffer,
    # since the buffer might be fragmented
    freed_ids = self.ring_buffer.free_buf(
        self.default_is_free_check, 2 * self.max_object_size
    )
    # update the metadata after freeing up space
    for freed_id in freed_ids:
        key_to_free = self.id_index[freed_id]
        del self.key_index[key_to_free]
        del self.id_index[freed_id]
        del self.writer_flag[freed_id]

get_cached(key)

Get the cached object by key if it exists.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def get_cached(self, key: str) -> tuple[int, int]:
    """Get the cached object by key if it exists."""
    address, monotonic_id = self.key_index[key]
    self.increment_writer_flag(monotonic_id)
    self._touched.discard(monotonic_id)
    return address, monotonic_id

get_signature(key)

Sign the handle of a cached object so readers can verify it.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def get_signature(self, key: str) -> list[int]:
    """Sign the handle of a cached object so readers can verify it."""
    address, monotonic_id = self.key_index[key]
    return list(self._make_signature(key, address, monotonic_id))

handle()

Get handle for sharing across processes.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def handle(self):
    """Get handle for sharing across processes."""
    return ShmObjectStorageHandle(
        max_object_size=self.max_object_size,
        n_readers=self.n_readers,
        ring_buffer_handle=self.ring_buffer.handle(),
        serde_class=self.serde_class,
        reader_lock=self._reader_lock,
    )

increment_reader_flag(data_view)

Set the in-use flag for the reader.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def increment_reader_flag(self, data_view: memoryview) -> None:
    """Set the in-use flag for the reader."""
    # >0 for in-use flag
    reader_count = self.ring_buffer.byte2int(data_view)
    data_view[:] = self.ring_buffer.int2byte(reader_count + 1)

increment_writer_flag(id)

Set the in-use flag for the writer.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def increment_writer_flag(self, id: int) -> None:
    """Set the in-use flag for the writer."""
    self.writer_flag[id] = self.writer_flag.get(id, 0) + 1

is_cached(key)

Check if the object with the given key is cached.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def is_cached(self, key: str) -> bool:
    """Check if the object with the given key is cached."""
    return key in self.key_index

put(key, value)

Store a key-value pair in the object storage. Attempts to free max_object_size bytes using FIFO order when the ring buffer runs out of space during a put() operation.

Parameters:

  • key

    (str) –

    String key to identify the object

  • value

    (Any) –

    Any serializable Python object

Raises:

  • MemoryError –

    If there's not enough space in the buffer

  • ValueError –

    If the serialized object is too large

  • ValueError –

    If the key already exists in the storage

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def put(self, key: str, value: Any) -> tuple[int, int]:
    """Store a key-value pair in the object storage.
    Attempts to free max_object_size bytes using FIFO order
    when the ring buffer runs out of space during a put() operation.

    Args:
        key: String key to identify the object
        value: Any serializable Python object

    Raises:
        MemoryError: If there's not enough space in the buffer
        ValueError: If the serialized object is too large
        ValueError: If the key already exists in the storage

    """
    if key in self.key_index:
        raise ValueError(f"Key '{key}' already exists in the storage.")

    object_data, data_bytes, object_metadata, md_bytes = self.ser_de.serialize(
        value
    )
    buffer_size = self.flag_bytes + data_bytes + md_bytes
    # Sanity checks
    if buffer_size > self.max_object_size:
        raise ValueError(
            f"Serialized object size ({buffer_size} bytes) exceeds "
            f"max object size ({self.max_object_size} bytes)"
        )

    # Allocate new buffer
    try:
        address, monotonic_id = self.ring_buffer.allocate_buf(buffer_size)
    except MemoryError:
        self.free_unused()
        # try again after freeing up space
        address, monotonic_id = self.ring_buffer.allocate_buf(buffer_size)

    # Write data to buffer
    with self.ring_buffer.access_buf(address) as (data_view, metadata):
        data_view[: self.flag_bytes] = self.ring_buffer.int2byte(0)
        self.copy_to_buffer(
            object_data, data_bytes, object_metadata, md_bytes, data_view
        )
    self.increment_writer_flag(monotonic_id)

    # Update key index
    self.key_index[key] = (address, monotonic_id)
    self.id_index[monotonic_id] = key
    return address, monotonic_id

release_touches()

Stop protecting touched items that were not looked up.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def release_touches(self) -> None:
    """Stop protecting touched items that were not looked up."""
    self._touched.clear()

touch(key, address=0, monotonic_id=0, signature=None)

Touch an existing cached item.

For writers (ShmObjectStoreSenderCache): Protect the item from eviction until get_cached or release_touches For readers (ShmObjectStoreReceiverCache): Validate the handle

Parameters:

  • key

    (str) –

    String key of the object to touch

  • address

    (int, default: 0 ) –

    Address of the object (only for readers)

  • monotonic_id

    (int, default: 0 ) –

    Monotonic ID of the object (only for readers)

  • signature

    (list[int] | None, default: None ) –

    Server-issued handle signature (only for readers)

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def touch(
    self,
    key: str,
    address: int = 0,
    monotonic_id: int = 0,
    signature: list[int] | None = None,
) -> None:
    """Touch an existing cached item.

    For writers (ShmObjectStoreSenderCache): Protect the item from
    eviction until get_cached or release_touches
    For readers (ShmObjectStoreReceiverCache): Validate the handle

    Args:
        key: String key of the object to touch
        address: Address of the object (only for readers)
        monotonic_id: Monotonic ID of the object (only for readers)
        signature: Server-issued handle signature (only for readers)

    """
    if self._reader_lock is None:
        if key not in self.key_index:
            return None
        self._touched.add(self.key_index[key][1])
    else:
        # Reject forged handles before dereferencing the supplied address.
        self.verify_signature(key, address, monotonic_id, signature)
        with self.ring_buffer.access_buf(address) as (_, buf_metadata):
            self._validate_monotonic_id(address, monotonic_id, buf_metadata)

verify_signature(key, address, monotonic_id, signature)

Verify that a handle was issued for the given key.

For writers: the handle must also be the key's current entry For readers: only the signature is checked

Parameters:

  • key

    (str | None) –

    String key the handle belongs to

  • address

    (int) –

    Address of the object

  • monotonic_id

    (int) –

    Monotonic ID of the object

  • signature

    (list[int] | None) –

    Signature issued with the handle

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def verify_signature(
    self,
    key: str | None,
    address: int,
    monotonic_id: int,
    signature: list[int] | None,
) -> None:
    """Verify that a handle was issued for the given key.

    For writers: the handle must also be the key's current entry
    For readers: only the signature is checked

    Args:
        key: String key the handle belongs to
        address: Address of the object
        monotonic_id: Monotonic ID of the object
        signature: Signature issued with the handle

    """
    if self.is_writer and (
        key is None or self.key_index.get(key) != (address, monotonic_id)
    ):
        raise ValueError("Invalid SHM handle signature for cache key.")
    self._verify_reader_signature(address, monotonic_id, signature, key)

SingleWriterShmRingBuffer

A single-writer, multiple-reader ring buffer implementation using shared memory. This class provides a thread-safe ring buffer where one process can write data while multiple processes/threads can read from it.

Architecture: - Uses shared memory for cross-process communication - Maintains metadata for each allocated buffer chunk in the writer process - Supports custom "is_free_fn" functions to determine when buffers can be reused - Each buffer chunk contains: [4-byte id][4-byte size][actual_data]

Key Concepts: - monotonic_id_start/end: Track the range of active buffer IDs - data_buffer_start/end: Track the physical memory range in use - Automatic wraparound when reaching buffer end - Lazy garbage collection based on is_free_fn checks

Example Usage Scenarios:

Scenario 1: Simple Linear Allocation

Buffer size: 100 bytes
Initial state: [................................................. ]
               ^start=end(0)

After allocating 20 bytes (id=0):
[id:0|size:20|data........][...................................]
^start(0)                  ^end(28)

After allocating 30 bytes (id=1):
[id:0|size:20|data........][id:1|size:30|data..............][..]
^start(0)                                                   ^end(66)

Scenario 2: Memory Reclamation

Before freeing (both buffers still in use):
[id:0|size:20|data........][id:1|size:30|data..............][..]
^start(0)                                                   ^end(66)

After id:0 is marked free by readers:
[FREED.................... ][id:1|size:30|data..............][..]
                            ^start(28)                       ^end(66)

After both are freed:
[FREED..............................................][..]
                                                     ^start=end(66)

Scenario 3: Wraparound Allocation (continuing from Scenario 2)

Starting from after memory reclamation in Scenario 2:
[FREED..............................................][..]
                                                     ^start=end(66)

Allocate 40 bytes (id=2) - only 34 bytes available at end, so wraparound:
[id:2|size:40|data........................][FREED.............][..]
                                          ^end(148)            ^start(66)

Scenario 4: Error Handling - Out of Space

Starting from after wraparound allocation in Scenario 3:
[id:2|size:40|data........................][FREED.............][..]
                                          ^end(148)            ^start(66)

Trying to allocate 20 more bytes:
occupied_size_new = end + size - start = 148 + 28 - 66 > buffer_size(100)
-> Raises MemoryError: "Not enough space in the data buffer"

Thread Safety: - Single writer: Only one process/thread should write (allocate_buf) - Multiple readers: Multiple processes/threads can read (access_buf) - Reader synchronization handled by is_free_fn callback - Writer handles garbage collection (free_buf) based on reader feedback

Shared Memory Layout: [32-byte handle-signing secret][ring buffer chunks...]

Memory Layout per Buffer Chunk: [4-byte monotonic_id][4-byte chunk_size][actual_data...] ^metadata_start ^data_start

The monotonic_id ensures data integrity - readers can verify they're accessing the correct data even after buffer wraparound or reuse.

Methods:

  • allocate_buf –

    Allocate a buffer MD_SIZE + size bytes in the shared memory.

  • byte2int –

    Convert bytes back to an integer.

  • clear –

    Clear the ring buffer.

  • close –

    Close the shared memory.

  • free_buf –

    Free a buffer of the given size. This is a no-op in shared memory,

  • get_auth_secret –

    Return the current handle-signing secret.

  • int2byte –

    Convert an integer to bytes.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
class SingleWriterShmRingBuffer:
    """A single-writer, multiple-reader ring buffer implementation using shared
    memory. This class provides a thread-safe ring buffer where one process
    can write data while multiple processes/threads can read from it.

    Architecture:
    - Uses shared memory for cross-process communication
    - Maintains metadata for each allocated buffer chunk in the writer process
    - Supports custom "is_free_fn" functions to determine when buffers can be
      reused
    - Each buffer chunk contains: `[4-byte id][4-byte size][actual_data]`

    Key Concepts:
    - monotonic_id_start/end: Track the range of active buffer IDs
    - data_buffer_start/end: Track the physical memory range in use
    - Automatic wraparound when reaching buffer end
    - Lazy garbage collection based on is_free_fn checks

    Example Usage Scenarios:

    Scenario 1: Simple Linear Allocation
    ```
    Buffer size: 100 bytes
    Initial state: [................................................. ]
                   ^start=end(0)

    After allocating 20 bytes (id=0):
    [id:0|size:20|data........][...................................]
    ^start(0)                  ^end(28)

    After allocating 30 bytes (id=1):
    [id:0|size:20|data........][id:1|size:30|data..............][..]
    ^start(0)                                                   ^end(66)
    ```

    Scenario 2: Memory Reclamation
    ```
    Before freeing (both buffers still in use):
    [id:0|size:20|data........][id:1|size:30|data..............][..]
    ^start(0)                                                   ^end(66)

    After id:0 is marked free by readers:
    [FREED.................... ][id:1|size:30|data..............][..]
                                ^start(28)                       ^end(66)

    After both are freed:
    [FREED..............................................][..]
                                                         ^start=end(66)
    ```

    Scenario 3: Wraparound Allocation (continuing from Scenario 2)
    ```
    Starting from after memory reclamation in Scenario 2:
    [FREED..............................................][..]
                                                         ^start=end(66)

    Allocate 40 bytes (id=2) - only 34 bytes available at end, so wraparound:
    [id:2|size:40|data........................][FREED.............][..]
                                              ^end(148)            ^start(66)
    ```

    Scenario 4: Error Handling - Out of Space
    ```
    Starting from after wraparound allocation in Scenario 3:
    [id:2|size:40|data........................][FREED.............][..]
                                              ^end(148)            ^start(66)

    Trying to allocate 20 more bytes:
    occupied_size_new = end + size - start = 148 + 28 - 66 > buffer_size(100)
    -> Raises MemoryError: "Not enough space in the data buffer"
    ```

    Thread Safety:
    - Single writer: Only one process/thread should write (allocate_buf)
    - Multiple readers: Multiple processes/threads can read (access_buf)
    - Reader synchronization handled by is_free_fn callback
    - Writer handles garbage collection (free_buf) based on reader feedback

    Shared Memory Layout:
    `[32-byte handle-signing secret][ring buffer chunks...]`

    Memory Layout per Buffer Chunk:
    `[4-byte monotonic_id][4-byte chunk_size][actual_data...]`
    ^metadata_start                         ^data_start

    The monotonic_id ensures data integrity - readers can verify they're
    accessing the correct data even after buffer wraparound or reuse.
    """

    def __init__(
        self,
        data_buffer_size: int,
        name: str | None = None,
        create: bool = False,
    ):
        self.data_buffer_size = data_buffer_size
        self.is_writer = create

        self.ID_NBYTES = 4
        self.ID_MAX = 2**31  # exclusive, so 2**31 - 1 is the max value
        self.SIZE_NBYTES = 4
        # 4 bytes for id, 4 bytes for buffer size
        self.MD_SIZE = self.ID_NBYTES + self.SIZE_NBYTES
        self.AUTH_SECRET_BYTES = 32
        self.DATA_OFFSET = self.AUTH_SECRET_BYTES
        self.monotonic_id_end = 0
        self.monotonic_id_start = 0
        self.data_buffer_start = 0
        self.data_buffer_end = 0

        if create:
            logger.debug("Creating new shared memory buffer: %s", name)
            # we are creating a buffer
            self.metadata: dict[int, int] = {}  # monotonic_id -> start address
            self.shared_memory = shared_memory.SharedMemory(
                create=True,
                size=self.data_buffer_size + self.DATA_OFFSET,
                name=name,
            )
            self._init_auth_secret()
        else:
            # we are opening an existing buffer
            # fix to https://stackoverflow.com/q/62748654/9191338
            # Python incorrectly tracks shared memory even if it is not
            # created by the process. The following patch is a workaround.
            with patch(
                "multiprocessing.resource_tracker.register",
                lambda *args, **kwargs: None,
            ):
                self.shared_memory = shared_memory.SharedMemory(name=name)
                # See https://docs.python.org/3/library/multiprocessing.shared_memory.html # noqa
                # Some platforms allocate memory based on page size,
                # so the shared memory block size may be larger or equal
                # to the requested size. The size parameter is ignored
                # when attaching to an existing block.
                assert (
                    self.shared_memory.size >= self.data_buffer_size + self.DATA_OFFSET
                )

        logger.debug(
            "Shared memory created/opened with name: %s, size: %d",
            self.shared_memory.name,
            self.data_buffer_size,
        )

    def handle(self):
        return (
            self.data_buffer_size,
            self.shared_memory.name,
        )

    def clear(self) -> None:
        """Clear the ring buffer."""
        assert self.is_writer, "Only the writer can clear the buffer."
        self.metadata.clear()
        # Keep IDs increasing so handles issued before the clear (e.g. for
        # requests still draining) fail the ID check once their chunk is
        # reused, instead of resolving to a different object.
        self.monotonic_id_start = self.monotonic_id_end
        self.data_buffer_start = 0
        self.data_buffer_end = 0

    def get_auth_secret(self) -> bytes:
        """Return the current handle-signing secret."""
        assert self.shared_memory.buf is not None, "Buffer has been closed"
        return bytes(self.shared_memory.buf[: self.AUTH_SECRET_BYTES])

    def _init_auth_secret(self) -> None:
        assert self.is_writer, "Only the writer can set the auth secret."
        assert self.shared_memory.buf is not None, "Buffer has been closed"
        self.shared_memory.buf[: self.AUTH_SECRET_BYTES] = secrets.token_bytes(
            self.AUTH_SECRET_BYTES
        )

    def close(self) -> None:
        """Close the shared memory."""
        if hasattr(self, "shared_memory"):
            self.shared_memory.close()
            if self.is_writer:
                with suppress(FileNotFoundError):
                    self.shared_memory.unlink()

    def __del__(self):
        self.close()

    def int2byte(self, integer: int) -> bytes:
        """Convert an integer to bytes."""
        return integer.to_bytes(self.ID_NBYTES, "little", signed=True)

    def byte2int(self, byte_data: bytes) -> int:
        """Convert bytes back to an integer."""
        return int.from_bytes(byte_data, "little", signed=True)

    def allocate_buf(self, size: int) -> tuple[int, int]:
        """Allocate a buffer `MD_SIZE` + `size` bytes in the shared memory.
        Memory layout:
        `[4-byte monotonic_id][4-byte size][buffer data...]`
        """
        assert self.is_writer, "Only the writer can allocate buffers."
        assert size > 0, "Size must be greater than 0"
        assert self.shared_memory.buf is not None, "Buffer has been closed"
        size += self.MD_SIZE  # add metadata size to the buffer size
        # reset to beginning if the buffer does have enough contiguous space
        buffer_end_reset = self.data_buffer_end % self.data_buffer_size
        if buffer_end_reset + size > self.data_buffer_size:
            buffer_end_reset = (
                self.data_buffer_end // self.data_buffer_size + 1
            ) * self.data_buffer_size
        else:  # no reset needed
            buffer_end_reset = self.data_buffer_end

        # check if we have enough space in the data buffer
        # i.e. if the new end (self.data_buffer_end + size)
        # exceeds the start of the data buffer
        occupied_size_new = buffer_end_reset + size - self.data_buffer_start
        if occupied_size_new > self.data_buffer_size:
            raise MemoryError(
                "Not enough space in the data buffer, "
                "try calling free_buf() to free up space"
            )
        self.data_buffer_end = buffer_end_reset

        # first 4 bytes as the monotonic id
        buf_idx = self.data_buffer_end % self.data_buffer_size
        physical_idx = self.DATA_OFFSET + buf_idx
        self.shared_memory.buf[physical_idx : physical_idx + self.ID_NBYTES] = (
            self.int2byte(self.monotonic_id_end)
        )
        # next 4 bytes as the size of the data buffer
        self.shared_memory.buf[
            physical_idx + self.ID_NBYTES : physical_idx + self.MD_SIZE
        ] = self.int2byte(size)

        # record metadata
        self.metadata[self.monotonic_id_end % self.ID_MAX] = self.data_buffer_end
        # update buffer and monotonic id indices
        current_buffer_end = self.data_buffer_end
        current_id_end = self.monotonic_id_end
        self.data_buffer_end += size
        self.monotonic_id_end = (self.monotonic_id_end + 1) % self.ID_MAX
        return current_buffer_end, current_id_end

    @contextmanager
    def access_buf(self, address: int):
        assert self.shared_memory.buf is not None, "Buffer has been closed"
        buf_idx = address % self.data_buffer_size
        physical_idx = self.DATA_OFFSET + buf_idx

        # read metadata
        metadata_buff = self.shared_memory.buf[
            physical_idx : physical_idx + self.MD_SIZE
        ]
        id = self.byte2int(metadata_buff[: self.ID_NBYTES])
        size = self.byte2int(metadata_buff[self.ID_NBYTES : self.MD_SIZE])

        # yield the data buffer and metadata
        data_buff = self.shared_memory.buf[
            physical_idx + self.MD_SIZE : physical_idx + size
        ]
        with (
            memoryview(data_buff) as data_view,
        ):
            yield data_view, (id, size)

    def free_buf(
        self,
        is_free_fn: Callable[[int, memoryview], bool],
        nbytes: int | None = None,
    ) -> Iterable[int]:
        """Free a buffer of the given size. This is a no-op in shared memory,
        but we need to keep track of the metadata.

        If freed memory spreads across the end and start of the ring buffer,
        the actual freed memory will be in two segments. In this case there
        still might not be a contiguous space of `nbytes` available.

        Args:
            is_free_fn (Callable[[int, memoryview], bool]): Predicate called with
                a monotonic id and the buffer, returning True when that buffer
                can be reclaimed.
            nbytes (int, optional): The size of the buffer to free. If None,
                frees the maximum size of the ring buffer.

        """
        assert self.is_writer, "Only the writer can free buffers."
        logger.debug(
            "Freeing up space in the ring buffer, "
            "monotonic_id_start: %d, monotonic_id_end: %d",
            self.monotonic_id_start,
            self.monotonic_id_end,
        )
        monotonic_id_before = self.monotonic_id_start
        # if nbytes is None, free up the maximum size of the ring buffer
        if nbytes is None:
            nbytes = self.data_buffer_size
        freed_bytes = 0
        while self.monotonic_id_start in self.metadata and freed_bytes < nbytes:
            address = self.metadata[self.monotonic_id_start]
            with self.access_buf(address) as (data_buff, metadata):
                if is_free_fn(self.monotonic_id_start, data_buff):
                    # check passed, we can free the buffer
                    del self.metadata[self.monotonic_id_start]
                    self.monotonic_id_start = (
                        self.monotonic_id_start + 1
                    ) % self.ID_MAX
                    if self.monotonic_id_start in self.metadata:
                        # pointing to the start addr of next allocation
                        self.data_buffer_start += (
                            self.metadata[self.monotonic_id_start]
                            - self.data_buffer_start
                        ) % self.data_buffer_size
                    else:
                        # no remaining allocation, reset to zero
                        self.data_buffer_start = self.data_buffer_end = 0
                    freed_bytes += metadata[1]
                else:
                    # there are still readers, we cannot free the buffer
                    break

        logger.debug(
            "Freed %d bytes from the ring buffer, "
            "monotonic_id_start: %d, monotonic_id_end: %d",
            freed_bytes,
            self.monotonic_id_start,
            self.monotonic_id_end,
        )

        # buffer wrap around
        if self.data_buffer_start >= self.data_buffer_size:
            self.data_buffer_start -= self.data_buffer_size
            self.data_buffer_end -= self.data_buffer_size

        monotonic_id_after = self.monotonic_id_start
        # id wrap around
        if monotonic_id_after >= monotonic_id_before:
            return range(monotonic_id_before, monotonic_id_after)
        else:
            return chain(
                range(monotonic_id_before, self.ID_MAX), range(0, monotonic_id_after)
            )

allocate_buf(size)

Allocate a buffer MD_SIZE + size bytes in the shared memory. Memory layout: [4-byte monotonic_id][4-byte size][buffer data...]

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def allocate_buf(self, size: int) -> tuple[int, int]:
    """Allocate a buffer `MD_SIZE` + `size` bytes in the shared memory.
    Memory layout:
    `[4-byte monotonic_id][4-byte size][buffer data...]`
    """
    assert self.is_writer, "Only the writer can allocate buffers."
    assert size > 0, "Size must be greater than 0"
    assert self.shared_memory.buf is not None, "Buffer has been closed"
    size += self.MD_SIZE  # add metadata size to the buffer size
    # reset to beginning if the buffer does have enough contiguous space
    buffer_end_reset = self.data_buffer_end % self.data_buffer_size
    if buffer_end_reset + size > self.data_buffer_size:
        buffer_end_reset = (
            self.data_buffer_end // self.data_buffer_size + 1
        ) * self.data_buffer_size
    else:  # no reset needed
        buffer_end_reset = self.data_buffer_end

    # check if we have enough space in the data buffer
    # i.e. if the new end (self.data_buffer_end + size)
    # exceeds the start of the data buffer
    occupied_size_new = buffer_end_reset + size - self.data_buffer_start
    if occupied_size_new > self.data_buffer_size:
        raise MemoryError(
            "Not enough space in the data buffer, "
            "try calling free_buf() to free up space"
        )
    self.data_buffer_end = buffer_end_reset

    # first 4 bytes as the monotonic id
    buf_idx = self.data_buffer_end % self.data_buffer_size
    physical_idx = self.DATA_OFFSET + buf_idx
    self.shared_memory.buf[physical_idx : physical_idx + self.ID_NBYTES] = (
        self.int2byte(self.monotonic_id_end)
    )
    # next 4 bytes as the size of the data buffer
    self.shared_memory.buf[
        physical_idx + self.ID_NBYTES : physical_idx + self.MD_SIZE
    ] = self.int2byte(size)

    # record metadata
    self.metadata[self.monotonic_id_end % self.ID_MAX] = self.data_buffer_end
    # update buffer and monotonic id indices
    current_buffer_end = self.data_buffer_end
    current_id_end = self.monotonic_id_end
    self.data_buffer_end += size
    self.monotonic_id_end = (self.monotonic_id_end + 1) % self.ID_MAX
    return current_buffer_end, current_id_end

byte2int(byte_data)

Convert bytes back to an integer.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def byte2int(self, byte_data: bytes) -> int:
    """Convert bytes back to an integer."""
    return int.from_bytes(byte_data, "little", signed=True)

clear()

Clear the ring buffer.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def clear(self) -> None:
    """Clear the ring buffer."""
    assert self.is_writer, "Only the writer can clear the buffer."
    self.metadata.clear()
    # Keep IDs increasing so handles issued before the clear (e.g. for
    # requests still draining) fail the ID check once their chunk is
    # reused, instead of resolving to a different object.
    self.monotonic_id_start = self.monotonic_id_end
    self.data_buffer_start = 0
    self.data_buffer_end = 0

close()

Close the shared memory.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def close(self) -> None:
    """Close the shared memory."""
    if hasattr(self, "shared_memory"):
        self.shared_memory.close()
        if self.is_writer:
            with suppress(FileNotFoundError):
                self.shared_memory.unlink()

free_buf(is_free_fn, nbytes=None)

Free a buffer of the given size. This is a no-op in shared memory, but we need to keep track of the metadata.

If freed memory spreads across the end and start of the ring buffer, the actual freed memory will be in two segments. In this case there still might not be a contiguous space of nbytes available.

Parameters:

  • is_free_fn

    (Callable[[int, memoryview], bool]) –

    Predicate called with a monotonic id and the buffer, returning True when that buffer can be reclaimed.

  • nbytes

    (int, default: None ) –

    The size of the buffer to free. If None, frees the maximum size of the ring buffer.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def free_buf(
    self,
    is_free_fn: Callable[[int, memoryview], bool],
    nbytes: int | None = None,
) -> Iterable[int]:
    """Free a buffer of the given size. This is a no-op in shared memory,
    but we need to keep track of the metadata.

    If freed memory spreads across the end and start of the ring buffer,
    the actual freed memory will be in two segments. In this case there
    still might not be a contiguous space of `nbytes` available.

    Args:
        is_free_fn (Callable[[int, memoryview], bool]): Predicate called with
            a monotonic id and the buffer, returning True when that buffer
            can be reclaimed.
        nbytes (int, optional): The size of the buffer to free. If None,
            frees the maximum size of the ring buffer.

    """
    assert self.is_writer, "Only the writer can free buffers."
    logger.debug(
        "Freeing up space in the ring buffer, "
        "monotonic_id_start: %d, monotonic_id_end: %d",
        self.monotonic_id_start,
        self.monotonic_id_end,
    )
    monotonic_id_before = self.monotonic_id_start
    # if nbytes is None, free up the maximum size of the ring buffer
    if nbytes is None:
        nbytes = self.data_buffer_size
    freed_bytes = 0
    while self.monotonic_id_start in self.metadata and freed_bytes < nbytes:
        address = self.metadata[self.monotonic_id_start]
        with self.access_buf(address) as (data_buff, metadata):
            if is_free_fn(self.monotonic_id_start, data_buff):
                # check passed, we can free the buffer
                del self.metadata[self.monotonic_id_start]
                self.monotonic_id_start = (
                    self.monotonic_id_start + 1
                ) % self.ID_MAX
                if self.monotonic_id_start in self.metadata:
                    # pointing to the start addr of next allocation
                    self.data_buffer_start += (
                        self.metadata[self.monotonic_id_start]
                        - self.data_buffer_start
                    ) % self.data_buffer_size
                else:
                    # no remaining allocation, reset to zero
                    self.data_buffer_start = self.data_buffer_end = 0
                freed_bytes += metadata[1]
            else:
                # there are still readers, we cannot free the buffer
                break

    logger.debug(
        "Freed %d bytes from the ring buffer, "
        "monotonic_id_start: %d, monotonic_id_end: %d",
        freed_bytes,
        self.monotonic_id_start,
        self.monotonic_id_end,
    )

    # buffer wrap around
    if self.data_buffer_start >= self.data_buffer_size:
        self.data_buffer_start -= self.data_buffer_size
        self.data_buffer_end -= self.data_buffer_size

    monotonic_id_after = self.monotonic_id_start
    # id wrap around
    if monotonic_id_after >= monotonic_id_before:
        return range(monotonic_id_before, monotonic_id_after)
    else:
        return chain(
            range(monotonic_id_before, self.ID_MAX), range(0, monotonic_id_after)
        )

get_auth_secret()

Return the current handle-signing secret.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def get_auth_secret(self) -> bytes:
    """Return the current handle-signing secret."""
    assert self.shared_memory.buf is not None, "Buffer has been closed"
    return bytes(self.shared_memory.buf[: self.AUTH_SECRET_BYTES])

int2byte(integer)

Convert an integer to bytes.

Source code in vllm/distributed/device_communicators/shm_object_storage.py
def int2byte(self, integer: int) -> bytes:
    """Convert an integer to bytes."""
    return integer.to_bytes(self.ID_NBYTES, "little", signed=True)