Write Custom Routing Strategies

Build Rust filters, scorers, and pickers for the Dynamo frontend or EPP
以 Markdown 格式查看

Experimental. A custom worker-selection policy controls how Dynamo filters and scores eligible workers, then selects one. Dynamo still owns discovery, eligibility, queueing, reservations, accounting, and metrics.

How It Works

This feature replaces the worker-ranking part of Dynamo’s routing pipeline. A WorkerFilter can exclude a host-eligible worker, a WorkerScorer assigns a finite cost to each remaining worker, Dynamo adds costs from all configured scorers, and a WorkerPicker chooses one row from the scored candidates. Filters run in declaration order before scoring, and rejecting every worker returns an error. Lower costs rank first by convention, but the picker can implement deterministic selection, sampling, tie-breaking, or policy-local state. Dynamo continues to own discovery, eligibility, score validation, accounting, and reservation.

For multiple compile-checked policies, see the custom policy examples. Repository contributors who use a coding agent must also provide the worker-selection API rules.

Use a Built-In Policy First

The Dynamo frontend ships built-in worker-selection policies, in lib/router-plugins. Selecting one needs router-policy YAML only — no policy crate, no rebuild, no private image. See Worker-Selection Policies for the available types and how to select one.

Write your own policy when no shipped policy expresses the rule you need. The rest of this page covers that case; a custom policy is added alongside the shipped ones, so you keep both.

Choose the Policy Stage

StageInputOutputUse this stage for
WorkerFilterRequest context and one host-eligible workerKeep or rejectThe rule is a hard requirement, not a ranking preference
WorkerScorerRequest context and one eligible workerOne finite costThe rule ranks workers or adds a penalty
WorkerPickerAll eligible workers and their total costsOne row indexThe rule samples, breaks ties, tracks policy state, or ignores total cost
Policy factoryRouter configuration, worker type, and routing partitionOne WorkerSelectionPolicyPrefill, decode, models, or routing groups need different components

An external policy owns its filters, scorers, and picker. Dynamo’s default scorer and picker are internal and can change with the built-in routing algorithm.

Public Plugin API and Compatibility

Use dynamo_kv_router::plugins for plugin registration and these modules for plugin contracts:

ModulePublic API
plugins::worker_selectionWorkerFilter, WorkerScorer, WorkerPicker, WorkerSelectionPolicy, request context, worker signal accessors, configuration, and factories
plugins::request_classifierRequestClassifier, ClassifyRequest, ClassifyEvent, configuration, and factories
pluginsRouterPluginRegistry and the resolved RouterPlugins bundle

The module organization changes no signal accessor, callback signature, signal meaning, or configuration format. Cache overlap, worker load, session metadata, and classification inputs retain their existing APIs. Dynamo owns scheduler queues, eligibility checks, reservations, lifecycle synchronization, and the default selector implementation. Low-level host interfaces such as WorkerSelector are used to connect Dynamo components; external plugins use WorkerSelectionPolicy.

Existing plugin imports and the WorkerSelectionPolicyRegistry name remain available through compatibility re-exports until version 1.7. They refer to the same types and traits, so existing policies can register into RouterPluginRegistry without adapters. Migrate imports to plugins::worker_selection or plugins::request_classifier, rename the registry to RouterPluginRegistry, and use register_worker_selection for worker policies. The compatibility exports, registry alias, worker-only register method, and frontend worker_selection_policy_factory setter are scheduled for removal in 1.7. The signal methods themselves are unchanged.

Build the Policy

Create the Policy and Catalog Crates

Set the Dynamo checkout and policy project paths:

# Set the source and destination paths used by the remaining commands.
export DYNAMO_DIR=/work/dynamo
export POLICY_DIR=/work/acme-routing
# Create one crate for the policy and one crate for policy registration.
mkdir -p "$POLICY_DIR"
cargo init --lib --name acme-routing-policy "$POLICY_DIR/policy"
cargo init --lib --name acme-routing-catalog "$POLICY_DIR/catalog"

Add dynamo-kv-router from the same checkout that builds the frontend or EPP:

# Add the worker-selection API to the policy crate.
cargo add \
--manifest-path "$POLICY_DIR/policy/Cargo.toml" \
--path "$DYNAMO_DIR/lib/kv-router" \
dynamo-kv-router
# Add parameter deserialization to the policy crate.
cargo add \
--manifest-path "$POLICY_DIR/policy/Cargo.toml" \
--features derive serde
# Add the policy crate to the catalog.
cargo add \
--manifest-path "$POLICY_DIR/catalog/Cargo.toml" \
--path "$POLICY_DIR/policy" \
acme-routing-policy
# Add the registry API to the catalog.
cargo add \
--manifest-path "$POLICY_DIR/catalog/Cargo.toml" \
--path "$DYNAMO_DIR/lib/kv-router" \
dynamo-kv-router

Filter Workers

A filter answers a hard yes-or-no question about one worker that already passed Dynamo’s eligibility checks. Return true to keep the worker or false to remove it. Use a filter only when the worker must not receive the request; use a scorer for preferences.

This filter keeps workers with at least the configured number of device-resident overlap blocks:

use dynamo_kv_router::plugins::worker_selection::{
WorkerCandidate, WorkerFilter, WorkerInputs, WorkerSelectionContext,
WorkerSelectionPolicyError,
};
struct MinimumDeviceOverlapFilter {
minimum_blocks: f64,
}
impl WorkerFilter for MinimumDeviceOverlapFilter {
fn required_worker_inputs(&self) -> WorkerInputs {
WorkerInputs::CACHE
}
fn keep(
&mut self,
_context: &WorkerSelectionContext<'_>,
candidate: &WorkerCandidate,
) -> Result<bool, WorkerSelectionPolicyError> {
let cache = candidate
.cache()
.ok_or_else(|| WorkerSelectionPolicyError::failed("cache input unavailable"))?;
Ok(cache.device_overlap_blocks() >= self.minimum_blocks)
}
}

Dynamo runs filters in declaration order before scoring. A worker must pass every configured filter. If no workers remain, selection returns an error.

Pass filters to WorkerSelectionPolicy::new_with_filters. If the policy has no hard requirement, omit filters and use WorkerSelectionPolicy::new.

Score Workers

A scorer expresses a preference without excluding a worker. It returns one finite cost for one worker. Lower total cost is better by convention.

This scorer uses the current number of active requests as its cost:

use dynamo_kv_router::plugins::worker_selection::{
WorkerCandidate, WorkerInputs, WorkerScorer, WorkerSelectionContext,
WorkerSelectionPolicyError,
};
struct ActiveRequestsScorer;
impl WorkerScorer for ActiveRequestsScorer {
fn required_worker_inputs(&self) -> WorkerInputs {
WorkerInputs::LOAD
}
fn score(
&mut self,
_context: &WorkerSelectionContext<'_>,
candidate: &WorkerCandidate,
) -> Result<f64, WorkerSelectionPolicyError> {
let load = candidate
.load()
.ok_or_else(|| WorkerSelectionPolicyError::failed("load input unavailable"))?;
Ok(load.active_requests() as f64)
}
}

A policy can stack multiple scorers. Dynamo calls them in declaration order and adds their costs. Dynamo rejects a non-finite contribution or total.

Pick a Worker

A picker makes the final choice after filtering and scoring. It sees every remaining worker and its total cost, then returns one row index. Most policies pick the lowest cost, but a picker can instead sample, break ties, or use policy-local state.

This picker selects the lowest-cost row:

use dynamo_kv_router::plugins::worker_selection::{
WorkerInputView, WorkerPicker, WorkerSelectionContext, WorkerSelectionPolicyError,
};
struct LowestCostPicker;
impl WorkerPicker for LowestCostPicker {
fn pick(
&mut self,
_context: &WorkerSelectionContext<'_>,
input: WorkerInputView<'_>,
) -> Result<usize, WorkerSelectionPolicyError> {
input
.candidates()
.iter()
.enumerate()
.min_by(|(_, left), (_, right)| left.cost().total_cmp(&right.cost()))
.map(|(row, _)| row)
.ok_or_else(|| WorkerSelectionPolicyError::failed("no eligible worker"))
}
}

Candidate order is unspecified, so inspect explicit values instead of relying on row order. Dynamo rejects an out-of-range index before accounting or reservation.

Parse Parameters and Create the Factory

The provider runs once at startup. Parse and validate all parameters there, then capture the validated values in the factory:

use std::sync::Arc;
use dynamo_kv_router::plugins::worker_selection::{
WorkerSelectionPolicyFactory, WorkerSelectionPolicyParameters,
WorkerSelectionPolicyProviderError,
};
use dynamo_kv_router::plugins::worker_selection::{WorkerFilter, WorkerScorer, WorkerSelectionPolicy};
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct Parameters {
min_device_overlap_blocks: f64,
}
fn provider(
parameters: &WorkerSelectionPolicyParameters,
) -> Result<WorkerSelectionPolicyFactory, WorkerSelectionPolicyProviderError> {
let parameters: Parameters = parameters.deserialize()?;
if !parameters.min_device_overlap_blocks.is_finite()
|| parameters.min_device_overlap_blocks < 0.0
{
return Err(WorkerSelectionPolicyProviderError::new(
"min_device_overlap_blocks must be a finite non-negative number",
));
}
let minimum_blocks = parameters.min_device_overlap_blocks;
Ok(Arc::new(move |config, worker_type, _partition| {
let filters: Vec<Box<dyn WorkerFilter>> = vec![
Box::new(MinimumDeviceOverlapFilter { minimum_blocks }),
];
let scorers: Vec<Box<dyn WorkerScorer>> = vec![Box::new(ActiveRequestsScorer)];
WorkerSelectionPolicy::new_with_filters(
config.clone(),
worker_type.as_str(),
filters,
scorers,
Box::new(LowestCostPicker),
)
}))
}

Dynamo calls the returned factory once per routing partition. Branch on the typed WorkerType::Aggregated, WorkerType::Prefill, WorkerType::Decode, or WorkerType::Encode role when worker pools need different components. Use the partition identity for distinct model or routing-group state.

Register the Policy

Expose a registration function from the policy crate:

use dynamo_kv_router::plugins::{RouterPluginRegistry, WorkerSelectionPolicyRegistryError};
pub fn register(
registry: &mut RouterPluginRegistry,
) -> Result<(), WorkerSelectionPolicyRegistryError> {
registry.register_worker_selection("least-busy", Arc::new(provider))
}

Call this function from the catalog:

pub fn register(
registry: &mut RouterPluginRegistry,
) -> Result<(), WorkerSelectionPolicyRegistryError> {
acme_routing_policy::register(registry)
}

Choose a stable, unique type name. Unknown types, duplicate registrations, and invalid parameters stop startup.

Configure an Instance

Create $POLICY_DIR/worker-selection.yaml:

worker_selection:
aggregated: least-busy
prefill: prefill-least-busy
decode: least-busy
encode: least-busy
instances:
- name: least-busy
type: least-busy
parameters:
min_device_overlap_blocks: 0
- name: prefill-least-busy
type: least-busy
parameters:
min_device_overlap_blocks: 1

The type selects a registered provider. The name identifies one configured instance. worker_selection.aggregated, worker_selection.prefill, worker_selection.decode, and worker_selection.encode select the matching worker pools. An omitted role uses Dynamo’s built-in policy; one role does not fall back to another role’s selection.

worker_selection.encode applies only to a surface-carrying encode worker set that constructs the standard KV chooser. It does not configure the surface-less multimodal encoder hop: EncoderRouter selects those workers independently with round-robin routing.

For prefill and decode, the selection order is the role-specific CLI or environment override, DYN_ROUTER_WORKER_SELECTION_POLICY, the matching YAML selection, and then Dynamo’s built-in policy. For aggregated and encode, the order is DYN_ROUTER_WORKER_SELECTION_POLICY, the matching YAML selection, and then the built-in policy. Set any selection to default to use the built-in policy for that scope.

Role selections are startup settings. Configure them in the worker-selection YAML, with the CLI flags or environment variables above, or through the Python KvRouterConfig constructor. They are not part of serialized KvRouterConfig JSON, and KvRouterConfig.from_json rejects those keys.

Check the Policy and Catalog

cargo check --manifest-path "$POLICY_DIR/catalog/Cargo.toml"

Add one focused test for each policy decision and one registration test for every type name.

Register a Request Classifier

Experimental. Request classifiers use the same statically linked catalog as worker filters, scorers, and pickers. Keep the catalog’s single register function and call register_request_classifier for each classifier type. RouterPluginRegistry also accepts the existing worker-selection providers; WorkerSelectionPolicyRegistry remains an alias, so existing catalogs compile unchanged.

Implement RequestClassifier in your policy crate. Its provider validates RequestClassifierParameters at startup and returns a RequestClassifierFactory: a closure that receives RequestClassifierContext and creates a fresh Box<dyn RequestClassifier> for each routed model. The shared router builder supplies the cached capacity context for both HTTP discovery and Python KvRouter.generate. Register that provider alongside your worker-selection providers:

use std::sync::Arc;
use dynamo_kv_router::plugins::RouterPluginRegistry;
pub fn register(registry: &mut RouterPluginRegistry) -> Result<(), Box<dyn std::error::Error>> {
acme_worker_policy::register(registry)?;
registry.register_request_classifier("acme-admission", Arc::new(acme_admission::provider))?;
Ok(())
}

The provider signature is:

use dynamo_kv_router::plugins::request_classifier::{
RequestClassifierFactory, RequestClassifierParameters, RequestClassifierProviderError,
};
pub fn provider(
parameters: &RequestClassifierParameters,
) -> Result<RequestClassifierFactory, RequestClassifierProviderError>;

Select the classifier in the same --router-policy-config YAML:

request_classifier:
type: acme-admission
parameters:
max_in_flight: 32

The parameter mapping belongs to the plugin; max_in_flight above is an example plugin setting. Omit request_classifier to retain pass-through admission. Unknown types and invalid parameters fail startup. Build and link the catalog with the existing custom-policy workflow below; no second catalog export or build flag is required.

The discovery-backed HTTP frontend and Python KvRouter.generate path install classifiers on their KV routers. Classification applies to decode/aggregated admission before queue ordering and worker selection. Selection is fixed at frontend startup and is process-wide; per-model queue profiles do not select different classifier types. A classifier can change the policy class, due time, and scheduling cost through the existing ClassifyRequest API. Dynamo retains queueing, worker selection, reservations, and lifecycle delivery.

Shared Plugin Construction

The catalog registers plugin providers once. RouterPluginRegistry::resolve_plugins validates the selected configuration and returns a RouterPlugins bundle containing the configured factories. Both the HTTP frontend and Python router pass this bundle to the shared router construction path. Each router receives fresh plugin instances; factories run during construction, not for each request.

Rust embedders can use the same setup:

use dynamo_kv_router::plugins::RouterPluginRegistry;
use dynamo_llm::entrypoint::input::http::HttpFrontend;
let mut registry = RouterPluginRegistry::default();
acme_routing_catalog::register(&mut registry)?;
let plugins = registry.resolve_plugins(&kv_router_config)?;
HttpFrontend::default()
.plugins(plugins)
.run(distributed_runtime, engine_config)
.await?;

The host’s RouterPluginBuilder passes worker selection through the shared selection core and installs the remaining plugins before the router is exposed to requests. The worker policy is constructed once per partition, including the probe for its required inputs. Adding another plugin trait requires a typed registration slot, a factory in the bundle, and its construction hook. Entry points continue to carry one bundle. Execution and lifecycle callbacks stay on each trait; the bundle does not add per-request dispatch.

The compatibility paths described above remain available until 1.7. An empty bundle retains the default selector and pass-through admission.

The disaggregated prefill hop and query-only probes do not run classifiers. Stateful low-level best_worker calls also lack classifier lifecycle enrollment. Offline replay and standalone selection reject request_classifier configuration because they cannot run its lifecycle.

Available Signals

Request context and worker identity are always available. Dynamo calculates optional per-worker signals only for groups that a filter, scorer, or picker requests.

Request Context

AccessorMeaning
request_blocks()Incoming prompt size in KV blocks
block_size()Tokens in one KV block
tracks_prefill_tokens()Whether the request contributes to prefill-load tracking
session_context()Optional session metadata described below
affinity_target()Optional advisory worker and data-parallel rank from soft session affinity
expected_output_tokens()Optional expected output length
priority_jump()Scheduler priority boost. Queue policies treat negative values as zero
strict_priority()Strict integer priority. The queue orders larger values first
router_temperature_override()Optional per-request router temperature override

Session Context

session_context() returns None when the request has no session metadata. This policy-facing view contains selected session metadata; it is not Dynamo’s internal request envelope. When present, it provides:

AccessorMeaning
session_id()Stable reasoning or tool-session identifier
parent_session_id()Optional parent session for subagents
session_final()Optional terminal marker for lifecycle-aware policies
input_trigger()Optional UserMessage, ToolResult, or Other request trigger

The custom policy examples use input_trigger() to give tool-result turns a cache-local picker path.

Soft Session Affinity

Set --router-session-affinity-mode soft with a session-affinity TTL to let a custom policy influence an existing session binding. The current binding enters the normal selection pipeline as affinity_target(). The policy still receives the full host-eligible candidate set, subject to its own filters, and may select another worker. Dynamo rebinds the session after dispatch returns a response stream.

Match both fields when a policy wants to recognize the target. A target with dp_rank: None matches every rank on its worker; a populated rank matches only that worker-rank pair. The target can be absent from the candidate table when the worker is unavailable or a custom filter rejects it.

Hard affinity remains the default. A hard binding and every explicit request target take the exact-target path instead of advisory custom selection. A selection, setup, or dispatch failure before a response stream leaves the previous soft binding intact. An error or cancellation after the stream is returned does not roll back a completed rebind.

The custom policy examples include an overload-aware soft pinning policy that retains the advisory target until its active-request count exceeds a configured threshold, then selects the least-loaded alternative. Its two-Mocker walkthrough holds the first request open and verifies an A -> B -> B worker sequence: overload moves the second request to B, and the third request retains B after load drains, proving that the plugin-selected dispatch updated the binding.

Worker Identity and Cost

AccessorMeaning
WorkerCandidate::worker()Candidate worker ID and data-parallel rank
ScoredWorkerCandidate::worker()Picker row worker ID and data-parallel rank
ScoredWorkerCandidate::cost()Sum of all scorer contributions for the picker row

Optional Worker Inputs

If a component needs no optional worker data, return WorkerInputs::NONE. Combine exact groups with |, such as WorkerInputs::CACHE | WorkerInputs::LOAD.

GroupAccessorMeaning
CACHEdevice_overlap_blocks()Device-resident prefix overlap in blocks
CACHEhost_overlap_blocks()Host-pinned prefix overlap in blocks
CACHEdisk_overlap_blocks()Disk prefix overlap in blocks
CACHEshared_beyond_device_blocks()Shared-cache hits beyond the device-resident prefix
LOADactive_prefill_tokens()Tokens currently active in the worker’s prefill stage
LOADdecode_cost_blocks()Projected active decode footprint, including this request’s additional active blocks
LOADactive_requests()Requests currently active on the worker
PREFERRED_TAINTpreferred_taint_multiplier()Optional cost multiplier from matching preferred routing taints

The cache accessors return raw tier counts. A missing tier or worker entry is zero; Dynamo does not substitute its weighted effective-overlap estimate. Each custom scorer chooses how to combine the raw counts.

A filter or scorer reads requested groups through WorkerCandidate::cache() or load(). A picker reads index-aligned arrays through WorkerInputView. Each component must request every optional group that it reads.

WorkerCandidate::preferred_taint_multiplier() and ScoredWorkerCandidate::preferred_taint_multiplier() return the optional cost multiplier from preferred routing constraints. A filter, scorer, or picker must request WorkerInputs::PREFERRED_TAINT before reading the multiplier. Without that declaration, Dynamo does not materialize the multiplier. Exact hard-pinned requests also do not materialize it. Required routing constraints remain host-enforced eligibility rules.

Both paths use the same policy crate, catalog, and YAML file. Choose the process that owns worker selection.

Add the catalog to the Python binding manifest. Keep the dependency alias dynamo-worker-selection-policy-catalog:

Your catalog is registered alongside the policies Dynamo ships, not instead of them, so both remain selectable. Your catalog cannot reuse a shipped policy’s type name; the registry rejects the duplicate at startup rather than overriding it.

# Link the policy catalog into the Python extension.
cargo add \
--manifest-path "$DYNAMO_DIR/lib/bindings/python/Cargo.toml" \
--optional \
--rename dynamo-worker-selection-policy-catalog \
--path "$POLICY_DIR/catalog" \
acme-routing-catalog

Build the extension with the linked catalog:

# Build the extension with the custom-policy feature.
cd "$DYNAMO_DIR/lib/bindings/python"
CARGO_TARGET_DIR="$DYNAMO_DIR/target" maturin develop --uv --features custom-policy
# Install the Python package from this checkout.
cd "$DYNAMO_DIR"
uv pip install -e .
# Start the frontend with the policy configuration.
python3 -m dynamo.frontend \
--router-mode kv \
--router-policy-config "$POLICY_DIR/worker-selection.yaml" \
--router-prefill-policy prefill-least-busy \
--router-decode-policy least-busy

The role-specific CLI flags override the matching YAML fields for this process. In this command, prefill-least-busy replaces worker_selection.prefill, and least-busy replaces worker_selection.decode.

The linked extension also applies custom policies in python3 -m dynamo.router. That process waits for a worker model card so it can choose built-in or custom policy behavior from the card’s typed role. DYN_ROUTER_MODEL_CARD_WAIT_SECS bounds the startup wait and defaults to 600 seconds.

Request Classifier Signals and Placement

Rust request classifiers can defer admission while retaining independently pollable futures. KvRouter::request_classifier_context() supplies the router’s cached worker view when constructing a classifier:

  • block_size() is the number of tokens per KV block. Each entry in workers() identifies a registered worker/rank and its advertised total_kv_blocks(), or None if capacity was not advertised. Their product is total capacity, not currently free space. The view follows existing discovery updates and does not probe worker health.
  • ClassifyRequest::progress() returns a read-only context high-water mark, initialized from input_tokens() and raised by host observations of prompt plus generated tokens. Clones share a monotonic counter across migration retries. A new request lifecycle gets a fresh counter, even if it reuses a request ID. This measures logical tokens, not physical KV occupancy.
  • set_worker_selection_target(worker) replaces the request’s soft affinity preference; clear_worker_selection_target() removes that preference. Hard pins and caller eligibility constraints remain authoritative. A custom picker can fall back when the preferred worker is ineligible. The classifier’s Sent event reports the worker actually used.

Holding a classification future pauses admission before dispatch; it does not interrupt an already running generation. Worker filters, scorers, and pickers continue to own placement decisions within Dynamo’s eligibility rules.

Policy Contract

  • Return true from a filter to keep a worker and false to reject it.
  • Expect filters to run in declaration order before scoring. Rejecting every worker returns an error.
  • Return finite scorer costs.
  • Return a valid picker row.
  • Treat candidate order as unspecified.
  • Request only the signal groups that the component reads.
  • Keep blocking I/O and panics out of keep, score, and pick.
  • Keep policy state local to the factory-created policy unless cross-partition sharing is a deliberate requirement.
  • Build the policy against the same Dynamo revision as the frontend or EPP.
  • If a signal adds work, storage, allocation, or another scan, run the worker-selection benchmark.

The example README contains the in-tree package names and build-check commands. For the built-in cost model, see Routing Concepts. For the standalone selection lifecycle, see Standalone Selection Service.