Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

33 changes: 33 additions & 0 deletions container/build-indexer-image.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
#!/usr/bin/env bash
# Build the standalone KV-cache indexer image from COMMITTED dynamo source.
#
# ./container/build-indexer-image.sh [tag]
#
# Builds a multi-stage image (compile wheel -> slim runtime) using a git-archive
# of HEAD as the build context, so the multi-GB local target/ is never shipped
# to the docker daemon. Default tag: <REGISTRY>/dynamo-indexer:kvtest-<sha>.
# Push separately: docker push <tag>
set -euo pipefail

REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
REGISTRY="${REGISTRY:-localhost:30500}"
SHA="$(git -C "$REPO_ROOT" rev-parse --short=10 HEAD)"
TAG="${1:-${REGISTRY}/dynamo-indexer:kvtest-${SHA}}"

BUILD_DIR="$(mktemp -d)"
trap 'rm -rf "$BUILD_DIR"' EXIT

echo "==> git archive ${SHA} -> build context (committed files only)"
git -C "$REPO_ROOT" archive --format=tar HEAD -o "$BUILD_DIR/dynamo-src.tar"
cp "$REPO_ROOT/container/indexer.Dockerfile" "$BUILD_DIR/Dockerfile"

echo "==> docker build $TAG"
# --ulimit nofile: BuildKit RUN steps otherwise inherit a low fd limit, which a
# big parallel cargo build exhausts ("too many open files"). The Dockerfile also
# caps CARGO_BUILD_JOBS, but raise the ceiling here too.
DOCKER_BUILDKIT=1 docker build \
--ulimit nofile=1048576:1048576 \
-f "$BUILD_DIR/Dockerfile" -t "$TAG" "$BUILD_DIR"

echo "==> done: $TAG"
echo " push with: docker push $TAG"
93 changes: 93 additions & 0 deletions container/indexer.Dockerfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
# syntax=docker/dockerfile:1.7
# =============================================================================
# Standalone KV-cache indexer image -> `python -m dynamo.indexer`
#
# Serves DeepInfra's per-model kv-indexer services; one image, the flavor is a flag:
# kv-indexer:reality (no flags) real evictions, per-shard trees
# kv-indexer:routing --keep-evictions evictions parked, per-shard trees
# kv-indexer:h24 --h24 flat "stored by anyone" counterfactual
# Each discovers the model's engine pods, subscribes to their ZMQ KV events and
# serves /query /query_by_hash /workers /dump /health (+ /metrics).
#
# Self-contained: compiles the Rust python bindings (dynamo._core) from source,
# then installs ONLY that wheel into a slim runtime. Both stages are ubuntu:24.04,
# so the runtime libstdc++ is identical to what the extension was linked against
# -> no LD_PRELOAD shim needed (that hack was only for the mismatched conda env).
#
# The build context is a `git archive` of the dynamo source (committed files
# only), so the multi-GB local target/ never enters the build. The indexer is
# fully standalone: no etcd / NATS, just HTTP + ZMQ. See build-indexer-image.sh.
# =============================================================================

# ---------- builder: compile the wheel -------------------------------------
FROM ubuntu:24.04 AS builder
ENV DEBIAN_FRONTEND=noninteractive
RUN apt-get update && apt-get install -y --no-install-recommends \
ca-certificates curl git \
build-essential pkg-config cmake \
clang libclang-dev \
protobuf-compiler libprotobuf-dev \
libzmq3-dev \
python3 python3-dev python3-venv \
&& rm -rf /var/lib/apt/lists/*

# Rust: rustup with NO default toolchain. dynamo's rust-toolchain.toml pins
# the toolchain, which rustup auto-fetches on the first cargo run in /src.
RUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs \
| sh -s -- -y --no-modify-path --default-toolchain none --profile minimal
ENV PATH=/root/.cargo/bin:$PATH

# maturin in an isolated venv (Ubuntu 24.04 system python is externally-managed).
RUN python3 -m venv /opt/maturin && /opt/maturin/bin/pip install --no-cache-dir maturin
ENV PATH=/opt/maturin/bin:$PATH

# Source tarball (git archive), auto-extracted by ADD. rust-toolchain.toml lands
# at /src so rustup resolves the pinned toolchain for the whole workspace.
ADD dynamo-src.tar /src
# git archive stamps every file with the commit time, and cargo judges a path
# crate fresh by mtime against the shared /cargo-target cache. A sibling build
# (another branch) can therefore leave a newer artifact that cargo reuses with
# the wrong contents; stamping the sources with the build time prevents that.
RUN find /src -type f -exec touch {} +

WORKDIR /src/lib/bindings/python
ENV CARGO_TARGET_DIR=/cargo-target
# Cap parallelism: this compiles on a shared ~288-core services node. The default
# (one rustc per core) both starves prod neighbours and exhausts the build's
# file-descriptor limit (hundreds of parallel rustc -> "too many open files").
# 16 jobs is plenty and well-behaved; build-indexer-image.sh also raises nofile.
ENV CARGO_BUILD_JOBS=16
# Release build: optimized for real indexer throughput (slower to compile than a
# debug build, but the radix-tree query/apply hot paths need the optimizations).
# kv-indexer-metrics adds the `dynamo.indexer` binary + Prometheus /metrics.
# nixl-sys' build script runs bindgen, so point it at libclang (resolved
# dynamically to survive llvm bumps). Cache mounts keep the crate registry +
# compiled artifacts warm across rebuilds.
RUN --mount=type=cache,target=/root/.cargo/registry \
--mount=type=cache,target=/cargo-target \
export LIBCLANG_PATH="$(dirname "$(find /usr/lib/llvm-* -name 'libclang.so*' 2>/dev/null | head -1)")" \
&& maturin build --release --features kv-indexer-metrics --out /wheels

# ---------- runtime: install just the wheel --------------------------------
FROM ubuntu:24.04 AS runtime
ENV DEBIAN_FRONTEND=noninteractive
RUN apt-get update && apt-get install -y --no-install-recommends \
ca-certificates \
libzmq5 \
python3 python3-venv \
&& rm -rf /var/lib/apt/lists/*

# venv with only the indexer wheel and its deps (pydantic, uvloop -> manylinux
# wheels from PyPI, no compiler needed here).
COPY --from=builder /wheels/ /wheels/
RUN python3 -m venv /opt/venv \
&& /opt/venv/bin/pip install --no-cache-dir /wheels/ai_dynamo_runtime*.whl \
&& rm -rf /wheels
ENV PATH=/opt/venv/bin:$PATH

# Smoke-test the import + CLI at build time so a broken wheel fails the build.
RUN python -m dynamo.indexer --help >/dev/null

EXPOSE 8090
ENTRYPOINT ["python", "-m", "dynamo.indexer"]
CMD ["--port", "8090"]
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,10 @@ Returns:
Register a ZMQ endpoint for an instance. Each call creates or reuses the indexer for the given `(model_name, routing_group)` pair.
Registration is non-blocking: if the worker is not up yet, the listener is accepted in `pending` state and transitions to `active` once the initial ZMQ connection succeeds.

Repeating `/register` for an existing `instance_id` and `dp_rank` returns an error;
it does not replace the listener or reset its sequence watermark. For worker restarts,
see [Worker Restarts](#worker-restarts).

```bash
# Single model, default routing group
curl -X POST http://localhost:8090/register \
Expand Down Expand Up @@ -490,6 +494,29 @@ recovers history still retained by the engine. Without a replay endpoint, or whe
range is no longer retained, restart with a healthy indexer in `--peers` to recover its startup
snapshot, or perform another full resynchronization before relying on the rebuilt worker state.

## Worker Restarts

A ZMQ reconnect does not identify a new worker or cache-event publisher lifetime.
If the publisher restarts with its sequence reset, an existing listener retains its
watermark and discards batches at or below it. Previously indexed ownership can also
remain stale.

On a worker or publisher restart, call `/unregister` and wait for the response before
calling `/register` with the current worker metadata. For a whole-worker restart,
omit `dp_rank` from `/unregister` to remove all ranks. Apply this lifecycle change to
every indexer replica; peer registration does not propagate worker lifecycle changes.
Serialize lifecycle operations for the same worker.

This sequence requests removal of old cache ownership and starts a fresh sequence
watermark. It does not certify a complete cache view or provide an atomic boundary
against in-flight old events. Registration success and listener status `active` do
not establish that startup events were received or replayed. Recover missing history
before relying on cache overlap, subject to the limits in
[Gap Detection and Replay](#gap-detection-and-replay).

The [standalone selector](standalone-selection.md#worker-restarts) has a different
registration API: updating a schedulable worker reconciles its indexer registration.

## Limitations

- **Standalone mode is ZMQ only**: Workers must publish KV events via ZMQ PUB sockets.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,26 @@ both default to `"default"` when omitted.
`GET /health` is process liveness. `GET /ready` returns `200` only after at
least one worker is schedulable, otherwise `503` with lifecycle details.

### Worker Restarts

When a worker or its cache-event publisher restarts, update its registration on every
selector replica. A reconnect at the same address does not reset the existing event
sequence watermark or invalidate old cache ownership.

With KV events enabled, `POST /workers` for an existing schedulable worker drains it,
removes its indexer registration, and creates new listeners with fresh sequence
watermarks. Supply the complete current worker record because this endpoint replaces
the catalog record. `PATCH /workers/{worker_id}` also reconciles a schedulable worker;
even an unchanged update can rebuild its cache view, so avoid using registration
updates as periodic heartbeats. Serialize lifecycle operations for the same worker.

This differs from the standalone indexer's `/register`, which rejects duplicate
registrations and requires `/unregister` first. See
[Indexer Worker Restarts](standalone-indexer.md#worker-restarts) for cleanup and
replay limits. Neither a successful catalog update nor `/ready` certifies a complete
cache view: missed startup events must still be recovered, and expired replay
history requires another source of complete state.

## Selection API

### `POST /select`
Expand Down
3 changes: 3 additions & 0 deletions lib/bindings/python/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion lib/bindings/python/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ custom-policy = [
"dynamo-kv-router/standalone-selection",
]
media-ffmpeg = ["dynamo-llm/media-ffmpeg"]
kv-indexer = ["dep:clap", "dep:tracing-subscriber", "dynamo-kv-router/standalone-indexer"]
kv-indexer = ["dep:clap", "dep:tracing-subscriber", "dynamo-kv-router/standalone-indexer", "dynamo-kv-router/kube-discovery"]
kv-indexer-metrics = ["kv-indexer", "dynamo-kv-router/metrics"]
# Exposes the optional Relay stats/snapshot methods to Python diagnostic builds.
ckf-diagnostics = ["dynamo-llm/ckf-diagnostics"]
Expand Down
131 changes: 131 additions & 0 deletions lib/bindings/python/rust/llm/kv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,122 @@ struct KvIndexerCli {
/// Use local timezone for access log timestamps (default: UTC)
#[arg(long)]
access_log_local_time: bool,

/// Total timeout in seconds for one engine `GET /kv_recover` download
/// (connect + body). A full TreeDump of a large engine is tens of MB.
#[arg(long, default_value_t = indexer::kv_recover::DEFAULT_RECOVER_TIMEOUT_S)]
recover_timeout_secs: u64,

/// Maximum concurrent `/kv_recover` downloads across all listeners of this
/// indexer. Bounds the load a fleet-wide (re)subscription puts on engines.
#[arg(long, default_value_t = indexer::kv_recover::DEFAULT_RECOVER_CONCURRENCY)]
recover_concurrency: usize,

/// Park Removed events in a buffer (and drop Cleared) instead of applying
/// them, approximating infinite memory (the kv-indexer:routing flavor).
/// Buffered evictions older than --evict-retention-secs are replayed into
/// the tree once memory use crosses --evict-memory-threshold of the limit.
#[arg(long, default_value_t = false)]
keep_evictions: bool,

/// Minimum age (seconds) a parked eviction must reach before the memory
/// sweep may apply it. Only meaningful with --keep-evictions.
#[arg(long, default_value_t = 1800)]
evict_retention_secs: u64,

/// Fraction of the memory limit (0..1] above which the sweep replays aged
/// evictions. Only meaningful with --keep-evictions.
#[arg(long, default_value_t = 0.75)]
evict_memory_threshold: f64,

/// Emit verbose audit logs on the `kv_audit` tracing target: one line per
/// query (block hashes + full response) and one per store/evict/clear
/// event ingested from the engines. Filter with RUST_LOG=kv_audit=info.
#[arg(long, default_value_t = false)]
enable_logging: bool,

/// Run as the flat h24 counterfactual (kv-indexer:h24): one map per model
/// of every prefix stored by any worker within --h24-horizon-secs,
/// ignoring evictions. Exclusive with --keep-evictions.
#[arg(long, default_value_t = false)]
h24: bool,

/// Retention horizon of the h24 expiry sweep, in seconds.
#[arg(long, default_value_t = 172800)]
h24_horizon_secs: u64,

/// Kubernetes namespace to watch for engine pods. Together with
/// --watch-model-name (or --watch-label) this enables pod auto-discovery:
/// subscribe on Ready, unsubscribe on delete.
#[arg(long)]
watch_namespace: Option<String>,

/// Raw label selector for this model's engine pods, e.g.
/// "engine_hash=d4b7a85131172ca6". Overrides the selector derived from
/// --watch-model-name.
#[arg(long)]
watch_label: Option<String>,

/// ZMQ KV-event port the discovered engines publish on.
#[arg(long, default_value_t = 5557)]
watch_zmq_port: u16,

/// HTTP port serving `GET /kv_recover` on the engines, used for gap
/// recovery at `http://<pod-ip>:<port>/kv_recover`.
#[arg(long)]
watch_recover_port: Option<u16>,

/// Data-parallel ranks per engine pod (vLLM --data-parallel-size). Rank r
/// publishes on --watch-zmq-port + r and recovers on
/// --watch-recover-port + r. Leaving this at 1 on a DP engine indexes
/// only rank 0's cache.
#[arg(long, default_value_t = 1)]
watch_dp_size: u32,

/// Model whose engine pods to discover, e.g. "openai/gpt-oss-120b". The
/// watch uses `di/model_name=<sanitized name>` unless --watch-label is
/// given; discovered pods register under this name (default:
/// --model-name).
#[arg(long)]
watch_model_name: Option<String>,
}

#[cfg(feature = "kv-indexer")]
impl KvIndexerCli {
/// Pod-discovery config, `None` unless --watch-namespace is set.
fn kube_discovery(&self) -> anyhow::Result<Option<indexer::discovery::KubeDiscoveryConfig>> {
let Some(namespace) = self.watch_namespace.clone() else {
anyhow::ensure!(
self.watch_label.is_none() && self.watch_model_name.is_none(),
"--watch-label/--watch-model-name require --watch-namespace"
);
return Ok(None);
};
let label_selector = match (&self.watch_label, &self.watch_model_name) {
(Some(raw), _) => raw.clone(),
(None, Some(model_name)) => indexer::discovery::model_name_label_selector(model_name),
(None, None) => {
anyhow::bail!("--watch-namespace requires --watch-model-name or --watch-label")
}
};
let block_size = self
.block_size
.ok_or_else(|| anyhow::anyhow!("--block-size is required with --watch-namespace"))?;
anyhow::ensure!(self.watch_dp_size >= 1, "--watch-dp-size must be at least 1");
Ok(Some(indexer::discovery::KubeDiscoveryConfig {
namespace,
label_selector,
zmq_port: self.watch_zmq_port,
recover_port: self.watch_recover_port,
dp_size: self.watch_dp_size,
model_name: self
.watch_model_name
.clone()
.unwrap_or_else(|| self.model_name.clone()),
routing_group: self.routing_group.clone(),
block_size,
}))
}
}

pub fn run_kv_indexer_cli<I, T>(args: I) -> anyhow::Result<()>
Expand All @@ -248,6 +364,8 @@ where
anyhow::anyhow!("invalid --trace-id-header '{}': {e}", cli.trace_id_header)
})?;

let kube_discovery = cli.kube_discovery()?;

init_standalone_logging();

let rt = tokio::runtime::Runtime::new()?;
Expand All @@ -262,6 +380,19 @@ where
access_log: cli.access_log,
trace_id_header,
access_log_local_time: cli.access_log_local_time,
kv_recover: indexer::kv_recover::KvRecoverSettings {
timeout_s: cli.recover_timeout_secs,
concurrency: cli.recover_concurrency,
},
kube_discovery,
audit_log: cli.enable_logging,
h24_horizon_s: cli.h24.then_some(cli.h24_horizon_secs),
keep_evictions: cli.keep_evictions.then_some(
indexer::evictions::KeepEvictionsConfig {
retention_s: cli.evict_retention_secs,
memory_threshold: cli.evict_memory_threshold,
},
),
}))
}

Expand Down
7 changes: 7 additions & 0 deletions lib/kv-router/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ profile = []
standalone-indexer = ["dep:axum", "dep:serde_json", "dep:reqwest", "dep:zmq", "dep:chrono", "dep:tracing-appender"]
standalone-slot-tracker = ["dep:axum", "dep:serde_json", "dep:zmq"]
standalone-selection = ["standalone-indexer", "dep:once_cell"]
# Kubernetes pod auto-discovery for the standalone indexer (pulls in the kube client).
kube-discovery = ["standalone-indexer", "dep:kube", "dep:k8s-openapi", "dep:futures"]

[dependencies]
# repo
Expand Down Expand Up @@ -76,6 +78,11 @@ chrono = { workspace = true, optional = true }
reqwest = { version = "0.12.24", default-features = false, features = ["json"], optional = true }
tracing-appender = { version = "0.2", optional = true }

# kube-discovery (optional)
kube = { workspace = true, optional = true }
k8s-openapi = { workspace = true, optional = true }
futures = { workspace = true, optional = true }

[dev-dependencies]
criterion = "0.5"
rstest = "0.18.2"
Expand Down
Loading
Loading