Skip to main content

PegaFlow Python Package

High-performance key-value storage engine with Python bindings, built with Rust and PyO3.

Features

  • PegaEngine: Fast Rust-based key-value storage with Python bindings
  • PegaKVConnector: vLLM KV connector for distributed inference with KV cache transfer

Installation

From Source

# Install maturin if you haven't already
pip install maturin

# Build and install in development mode
cd python
maturin develop

# Or build a wheel
maturin build --release

From PyPI (coming soon)

pip install pegaflow

Usage

Basic KV Storage

from pegaflow import PegaEngine

# Create a new engine
engine = PegaEngine()

# Store key-value pairs
engine.put("name", "PegaFlow")
engine.put("version", "0.1.0")

# Retrieve values
name = engine.get("name")  # Returns "PegaFlow"
missing = engine.get("nonexistent")  # Returns None

# Remove keys
removed = engine.remove("name")  # Returns "PegaFlow"

vLLM KV Connector

from vllm import LLM
from vllm.distributed.kv_transfer.kv_transfer_agent import KVTransferConfig

# Configure vLLM to use PegaKVConnector
kv_transfer_config = KVTransferConfig(
    kv_connector="PegaKVConnector",
    kv_role="kv_both",
    kv_connector_module_path="pegaflow.connector",
)

# Create LLM with KV transfer enabled
llm = LLM(
    model="gpt2",
    kv_transfer_config=kv_transfer_config,
)

Connector Modes

PegaKVConnector defaults to read_write: it queries PegaFlow for reusable KV blocks, loads matched blocks into vLLM, and saves newly computed full blocks back to PegaFlow.

Set pegaflow.mode to save_only when another vLLM connector is responsible for reads and PegaFlow should only persist KV blocks for later reuse. This is intended for MultiConnector decode-side setups where an upstream connector owns the external hit/load path, while PegaFlow records the resulting KV cache. In save_only mode, PegaFlow does not query or load KV blocks.

vllm serve Qwen/Qwen3-0.6B \
  --kv-transfer-config '{
    "kv_connector": "MultiConnector",
    "kv_role": "kv_both",
    "kv_connector_extra_config": {
      "connectors": [
        {
          "kv_connector": "<external-read-connector>",
          "kv_role": "kv_both"
        },
        {
          "kv_connector": "PegaKVConnector",
          "kv_role": "kv_both",
          "kv_connector_module_path": "pegaflow.connector",
          "kv_connector_extra_config": {
            "pegaflow.mode": "save_only"
          }
        }
      ]
    }
  }'

Valid values are read_write and save_only.

TP Shards Across Hosts

CUDA IPC is host-local. When one tensor-parallel replica spans multiple hosts, run one PegaFlow server on each host and configure the connector with every server endpoint in global TP-rank order:

{
  "kv_connector": "PegaKVConnector",
  "kv_role": "kv_both",
  "kv_connector_module_path": "pegaflow.connector",
  "kv_connector_extra_config": {
    "pegaflow.tp_shard_endpoints": [
      "http://host-a:50055",
      "http://host-b:50055"
    ]
  }
}

For TP8 and two endpoints, global ranks 0-3 register with the first server and ranks 4-7 register with the second. Each server sees a local TP4 topology and must manage the four GPUs on its own host. Every vLLM process must receive the same ordered endpoint list.

The scheduler queries every shard and only reuses the prefix available from all of them. Each worker loads with the lease issued by its local server. The connector gives every shard a distinct namespace, so deployments with a different host split cannot reuse an incompatible cache layout.

TP sharding currently requires equal contiguous shards and TP-only parallelism. Pipeline, decode-context, and prefill-context parallelism are rejected when more than one endpoint is configured.

P/D Partial Tail Blocks

vLLM normally exposes hashes only for complete KV blocks. In a P/D deployment, enable pegaflow.pd_tail_save on prefill and pegaflow.pd_tail_load on decode to reuse the final partial prompt block as well. Start both vLLM processes with the same explicit PYTHONHASHSEED and --prefix-caching-hash-algo xxhash_cbor.

Prefill: {"pegaflow.pd_tail_save": true}

Decode: {"pegaflow.pd_tail_load": true, "pegaflow.wait_for_full_prefix": true}

pegaflow.wait_for_full_prefix makes decode wait (up to 30s) until the full prompt prefix is fetchable from a remote node via MetaServer + RDMA. It only applies when prefill and decode run separate engines; it does not observe saves landing in a shared/local engine and has no effect when RDMA is not configured.

Development

See the examples directory for more usage examples.

Testing

Running Unit Tests

The test suite includes integration tests that verify the EngineRpcClient can correctly communicate with a running pegaflow-server instance.

Prerequisites

  1. Build the Rust extension:

    cd python
    maturin develop --release
    
  2. Build the server binary:

    cd ..
    cargo build --release --bin pegaflow-server
    
  3. Ensure CUDA is available (tests require GPU):

    python -c "import torch; assert torch.cuda.is_available()"
    

Running Tests

cd python

# Run all tests
pytest tests/ -v

# Run specific test file
pytest tests/test_engine_client.py -v

# Run with coverage
pytest tests/ --cov=pegaflow --cov-report=html

Test Structure

  • tests/conftest.py: Contains pytest fixtures for:

    • pega_server: Automatically starts/stops pegaflow-server for integration tests
    • engine_client: Creates an EngineRpcClient connected to the test server
    • client_context: Provides a ClientContext representing a vLLM instance with GPU KV cache tensors
    • registered_instance: Provides a registered instance ID for query tests
  • tests/test_engine_client.py: Integration tests for:

    • Server connectivity
    • Query operations with various inputs

Test Fixtures

The ClientContext class abstracts a vLLM instance and provides:

  • register_kv_caches(): Register GPU KV cache tensors with the server
  • query(block_hashes): Query available blocks
  • unregister_context(): Unregister context from server

Example test usage:

def test_query(client_context):
    """Test query operation."""
    result = client_context.query([])
    assert result is not None

License

MIT

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distributions

No source distribution files available for this release.See tutorial on generating distribution archives.

Built Distributions

If you're not sure about the file name format, learn more about wheel file names.

pegaflow_llm_cu13-0.23.9-cp314-cp314-manylinux_2_34_x86_64.whl (7.9 MB view details)

Uploaded CPython 3.14manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.9-cp314-cp314-manylinux_2_34_aarch64.whl (7.8 MB view details)

Uploaded CPython 3.14manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.9-cp313-cp313-manylinux_2_34_x86_64.whl (7.9 MB view details)

Uploaded CPython 3.13manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.9-cp313-cp313-manylinux_2_34_aarch64.whl (7.8 MB view details)

Uploaded CPython 3.13manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.9-cp312-cp312-manylinux_2_34_x86_64.whl (7.9 MB view details)

Uploaded CPython 3.12manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.9-cp312-cp312-manylinux_2_34_aarch64.whl (7.8 MB view details)

Uploaded CPython 3.12manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.9-cp311-cp311-manylinux_2_34_x86_64.whl (7.9 MB view details)

Uploaded CPython 3.11manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.9-cp311-cp311-manylinux_2_34_aarch64.whl (7.8 MB view details)

Uploaded CPython 3.11manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.9-cp310-cp310-manylinux_2_34_x86_64.whl (7.9 MB view details)

Uploaded CPython 3.10manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.9-cp310-cp310-manylinux_2_34_aarch64.whl (7.8 MB view details)

Uploaded CPython 3.10manylinux: glibc 2.34+ ARM64

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp314-cp314-manylinux_2_34_x86_64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp314-cp314-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 75fa10e34c7763b3360da82f9ae6fb1c09546770c16bb130683f1e16a0790b90
MD5 f9e37878673130f6fd377b0ce50584c7
BLAKE2b-256 407bd083d6d950287bd81c3f8a9cb967ee7eb0219c5363ba9958923ec54f4653

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp314-cp314-manylinux_2_34_aarch64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp314-cp314-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 340ca557443ea9e770af8d2888015de8c0f3da22d1447165ed271081f1020534
MD5 71f9eb025192a705aaf4dd7f49a92fe0
BLAKE2b-256 b13378260d94e2ba308cab8afa3160ce34628c55c76ac806407f7881a9d7c963

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp313-cp313-manylinux_2_34_x86_64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp313-cp313-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 0be75eba72a069cc5088cf43d8204c4d6d776f27960a3a4e54b09343cb558eef
MD5 2ecafa9c3a884190678eb10960e07f11
BLAKE2b-256 6a82041848d5e3ec07ccf1c6bbcb56c34e9f6c78b23fb54bd535a432309a84a5

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp313-cp313-manylinux_2_34_aarch64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp313-cp313-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 a636f1c71cea195112054b71e7478f91847c3bfb15844a25e9e8b3d01ee38112
MD5 36ae7916a1e334d28c46f388fedbdc40
BLAKE2b-256 95031fdd75ba53f1e8797bf558772c1aa5e48765690ff0c3e934396ee9fbb093

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp312-cp312-manylinux_2_34_x86_64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp312-cp312-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 166eeff4faac1d82c5e33cddf450d2c758ec596bb24f871c57f84ecbbe511489
MD5 b5a520991f0795f093998d4870c958ac
BLAKE2b-256 22b5a454624f3270e9484d5fc2c53de8c6f810845f166231bb821d9ea0d0da88

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp312-cp312-manylinux_2_34_aarch64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp312-cp312-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 24a391bba70ef23fed9684ccd087e9217c233298b0e8da1993a28053f972116f
MD5 47e6a5907fa118ca4cef48178672b1d4
BLAKE2b-256 72f2272080b9e5f4222fb87e8f911767e565195477ba42f1aaec92c83c800e9a

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp311-cp311-manylinux_2_34_x86_64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp311-cp311-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 05083aa1bf21016f935d99f13f1c7e0ba6a54f6044f0d3d430bccee1d1f87204
MD5 6f1ed0cf66d25a0b40979a62b25eb3a3
BLAKE2b-256 fea95459d2aff6ca0759ccf205c1f6d5d5c212f42e9d4bbb19f447316117f73d

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp311-cp311-manylinux_2_34_aarch64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp311-cp311-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 35747cb085026e08e2c82a0a1de6323761041eb4a7656e714c3f9faa3a45e412
MD5 db87cd8dca40ed5222d8ac06ca9ba8b1
BLAKE2b-256 c009894f4be5f1f27b1419094d757cb348becd674104d5fa7a7991535ea2df67

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp310-cp310-manylinux_2_34_x86_64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp310-cp310-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 405fde54c42ad8aac296b75e9a29c7770fb8277f4ba1078e815de40010eb7280
MD5 3b5710a12c38852ef53a041edc49424c
BLAKE2b-256 e32f56c0bed4e6af10c63e24fa9d5f46e678276401ff241c2937b45bf6baa312

See more details on using hashes here.

File details

Details for the file pegaflow_llm_cu13-0.23.9-cp310-cp310-manylinux_2_34_aarch64.whl.

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.9-cp310-cp310-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 73991381c7e1e1c04f0c78d2a0c73ee8c61b06bc1c038d079328c9e30b28c2c3
MD5 4ba0b38a6d192966c9f30f0cdeb18f9b
BLAKE2b-256 09a512f981d3ad677617055fc41a68b2fbde714ead8c597e8819ae40c81e55aa

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page