dynamo-llm 1.5.0

Dynamo LLM Library
// SPDX-FileCopyrightText: Copyright (c) 2024-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
// SPDX-License-Identifier: Apache-2.0

use std::collections::HashMap;
use std::collections::hash_map::Entry;

use dynamo_kv_router::protocols::{
    ExternalSequenceBlockHash, KvCacheRemoveData, KvCacheStoreData, ResidencyDomain, StorageTier,
};

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum EventDedupPolicy {
    RefCounted,
    SetLike,
}

/// Policy-driven deduplication filter for publisher KV cache events.
///
/// vLLM can emit multiple store/remove events for the same block hash.
/// Refcounts are tracked **per DP rank** because identical block hashes
/// on different ranks represent independent blocks.
///
/// `RefCounted` Stores increment a refcount and Removes pass only when it
/// reaches zero. `SetLike` events bypass bookkeeping when their producer
/// guarantees at most one logical residency per owner, tier, and block hash.
/// Clears reset refcounts for the emitting rank across storage tiers.
pub(super) struct EventDedupFilter {
    /// Per-(dp_rank, storage_tier, residency_domain) refcounts.
    per_rank_tier:
        HashMap<(u32, StorageTier, ResidencyDomain), HashMap<ExternalSequenceBlockHash, usize>>,
}

impl EventDedupFilter {
    pub(super) fn new() -> Self {
        Self {
            per_rank_tier: HashMap::new(),
        }
    }

    /// Track a store event. Increments refcount for each block hash on the
    /// given (DP rank, storage tier). Stores always pass through — this only
    /// updates bookkeeping.
    pub(super) fn track_store_in_domain(
        &mut self,
        dp_rank: u32,
        storage_tier: StorageTier,
        residency_domain: ResidencyDomain,
        policy: EventDedupPolicy,
        data: &KvCacheStoreData,
    ) {
        if policy == EventDedupPolicy::SetLike {
            return;
        }
        let refcounts = self
            .per_rank_tier
            .entry((dp_rank, storage_tier, residency_domain))
            .or_default();
        for block in &data.blocks {
            *refcounts.entry(block.block_hash).or_insert(0) += 1;
        }
    }

    /// Filter a remove event. Retains only block hashes whose refcount on the
    /// given (DP rank, storage tier) decrements to 0 (removing them from the
    /// map). Returns `None` if no hashes survive filtering.
    pub(super) fn filter_remove_in_domain(
        &mut self,
        dp_rank: u32,
        storage_tier: StorageTier,
        residency_domain: ResidencyDomain,
        policy: EventDedupPolicy,
        mut data: KvCacheRemoveData,
    ) -> Option<KvCacheRemoveData> {
        if policy == EventDedupPolicy::SetLike {
            return (!data.block_hashes.is_empty()).then_some(data);
        }
        let refcounts = self
            .per_rank_tier
            .entry((dp_rank, storage_tier, residency_domain))
            .or_default();
        data.block_hashes.retain(|hash| {
            match refcounts.entry(*hash) {
                Entry::Occupied(mut entry) => {
                    *entry.get_mut() -= 1;
                    if *entry.get() == 0 {
                        entry.remove();
                        true // refcount hit 0 -> pass through
                    } else {
                        false // still has references -> filter out
                    }
                }
                Entry::Vacant(_) => {
                    true // not tracked -> pass through defensively
                }
            }
        });
        if data.block_hashes.is_empty() {
            None
        } else {
            Some(data)
        }
    }

    /// Clear refcounts for one DP rank and residency domain across storage tiers.
    pub(super) fn clear_rank_domain(
        &mut self,
        dp_rank: u32,
        domain: ResidencyDomain,
        policy: EventDedupPolicy,
    ) {
        if policy == EventDedupPolicy::SetLike {
            return;
        }
        self.per_rank_tier
            .retain(|(tracked_dp_rank, _, tracked_domain), _| {
                *tracked_dp_rank != dp_rank || *tracked_domain != domain
            });
    }

    #[cfg(test)]
    pub(super) fn track_store(
        &mut self,
        dp_rank: u32,
        storage_tier: StorageTier,
        data: &KvCacheStoreData,
    ) {
        self.track_store_in_domain(
            dp_rank,
            storage_tier,
            ResidencyDomain::Worker,
            EventDedupPolicy::RefCounted,
            data,
        );
    }

    #[cfg(test)]
    pub(super) fn filter_remove(
        &mut self,
        dp_rank: u32,
        storage_tier: StorageTier,
        data: KvCacheRemoveData,
    ) -> Option<KvCacheRemoveData> {
        self.filter_remove_in_domain(
            dp_rank,
            storage_tier,
            ResidencyDomain::Worker,
            EventDedupPolicy::RefCounted,
            data,
        )
    }

    #[cfg(test)]
    pub(super) fn clear_rank(&mut self, dp_rank: u32) {
        self.per_rank_tier
            .retain(|(tracked_dp_rank, _, _), _| *tracked_dp_rank != dp_rank);
    }
}