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.13-cp314-cp314-manylinux_2_34_x86_64.whl (8.0 MB view details)

Uploaded CPython 3.14manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.13-cp314-cp314-manylinux_2_34_aarch64.whl (7.9 MB view details)

Uploaded CPython 3.14manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.13-cp313-cp313-manylinux_2_34_x86_64.whl (8.0 MB view details)

Uploaded CPython 3.13manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.13-cp313-cp313-manylinux_2_34_aarch64.whl (7.9 MB view details)

Uploaded CPython 3.13manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.13-cp312-cp312-manylinux_2_34_x86_64.whl (8.0 MB view details)

Uploaded CPython 3.12manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.13-cp312-cp312-manylinux_2_34_aarch64.whl (7.9 MB view details)

Uploaded CPython 3.12manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.13-cp311-cp311-manylinux_2_34_x86_64.whl (8.0 MB view details)

Uploaded CPython 3.11manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.13-cp311-cp311-manylinux_2_34_aarch64.whl (7.9 MB view details)

Uploaded CPython 3.11manylinux: glibc 2.34+ ARM64

pegaflow_llm_cu13-0.23.13-cp310-cp310-manylinux_2_34_x86_64.whl (8.0 MB view details)

Uploaded CPython 3.10manylinux: glibc 2.34+ x86-64

pegaflow_llm_cu13-0.23.13-cp310-cp310-manylinux_2_34_aarch64.whl (7.9 MB view details)

Uploaded CPython 3.10manylinux: glibc 2.34+ ARM64

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp314-cp314-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 c0ec7141facb876b80bf93bcef54a818c252c6e6e4a11ad611be46380e2c00f7
MD5 4cc22b4e92817b9fd102527e1cb40b71
BLAKE2b-256 d09ef1668a7a48d7f2171a4b76a51105c110106f22adc094182310f76fc93ece

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp314-cp314-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 1caf6ee37151369be312c6d3c6723dcec8a1ef0876e77a18b1db057afc25048e
MD5 39effbcec413c5d94312c6075bed8b4a
BLAKE2b-256 845d8ba75cce020c2f3539ffc409c1cfab952b31efc13176f725b607c2fd6719

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp313-cp313-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 d996efb1cbc81c9ab238d911bed40f52c5da105ac9c3bfc39e5ec171da0a81bd
MD5 05caf1d7f5985179959806625825ffaa
BLAKE2b-256 a7b30b45201b68c121f026a30a3a9757ca56dedac8d02999d93d08ebd5467476

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp313-cp313-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 fbbbf3b0d5a77d034374c842bdb30ef81b23581235d7661a377772ec8dbe78c0
MD5 529b880d532b33c4fb4afec02ced5290
BLAKE2b-256 99bbab9d45a6582fab1fea895f3cbdd989754f98191b2e65988e3cdd0343a532

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp312-cp312-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 574b8a476673e092e075d32e80490231200ed5a39e30dd36b48378e169495378
MD5 0ad34b4d0611add4c8fdc1e2b47de001
BLAKE2b-256 efac22b9de41f415e8eeaf4ff7b04497f166379f8f8d7a60ffd7f02c05c2a86d

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp312-cp312-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 674865dffcc201551a6f133fee5ee36f721c1ba9744728e416bf6cb961c56f65
MD5 49011f49b0f4223bde83db6080463a83
BLAKE2b-256 87bcd7493b39789c51457643bccf6e1226740b7f25d45536bc93245d858d859f

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp311-cp311-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 2c4c5bd39d01fb187e80472c75057948ab99dfecf939e968ce6bad5500946ac3
MD5 c8328c2fd63e21ff9ab2216883acaca3
BLAKE2b-256 db9688af580eef47e6fff240d3f070f83d2ad53a8f83c4a4542b65d2b8e04c9d

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp311-cp311-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 11e411c1026aae23f30c9db4b14212545f7eaf636ea5f79d8aba041cca70ccb0
MD5 2e42be5b8373b66254ae02147f18c7fe
BLAKE2b-256 a0fdf2c69661c785ebe93207c780929b127a5b93ad14a7e8327e3820d36c5415

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp310-cp310-manylinux_2_34_x86_64.whl
Algorithm Hash digest
SHA256 724903f1b930bc8731adac8f78b9836d20c0c17550f06782627258f46e5214c6
MD5 7657eaaf78342cde9158d7b85ea3d3b1
BLAKE2b-256 ea7f54b13154b6c3b2451f93d8a28d85ddc2de32b3afaaf0150e8f2a583c7996

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for pegaflow_llm_cu13-0.23.13-cp310-cp310-manylinux_2_34_aarch64.whl
Algorithm Hash digest
SHA256 c3cc603bba1bfc2b2ce843af54da7caf5bc6553961f157239d0059f9c4f4e3b9
MD5 95156611c1ecc12b8837ea9d6c6452a1
BLAKE2b-256 aadbfabe90f6868b5c591947fcf0465da54741f43a3892b4070ddd5541424d2f

See more details on using hashes here.

Release history Release notifications | RSS feed

0.24.1

10 files

0.24.0

10 files

This release

0.23.13 This release

10 files

0.23.12

10 files

0.23.11

10 files

0.23.10

10 files

0.23.9

10 files

0.23.8

10 files

0.23.7

5 files

0.23.6

5 files

0.23.5

5 files

0.23.4

5 files

0.23.3

5 files

0.23.2

5 files

0.23.1

5 files

0.23.0

5 files

0.22.10

5 files

0.22.9

5 files

0.22.8

5 files

0.22.7

5 files

0.22.6

5 files

0.22.5

5 files

0.22.4

5 files

0.22.3

5 files

0.22.2

5 files

0.22.1

5 files

0.22.0

5 files

0.21.2

5 files

0.21.1

4 files

0.21.0

4 files

0.20.0

4 files

0.19.1

4 files

0.19.0

4 files

0.18.0

4 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page