Skip to content
Merged
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
23 changes: 19 additions & 4 deletions src/dynamic_batch_scheduler.cc
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
// Copyright 2018-2024, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
// Copyright 2018-2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
//
// Redistribution and use in source and binary forms, with or without
// modification, are permitted provided that the following conditions
Expand Down Expand Up @@ -673,14 +673,31 @@ void
DynamicBatchScheduler::DelegateResponse(
std::unique_ptr<InferenceRequest>& request)
{
// Reserve this request's slot in the completion queue now, while we are
// still in scheduler-receive order. FinalizeResponses() stalls on an empty
// front slot, which is what holds later responses back until earlier ones
// are ready. Filling the slot only at completion time would order the queue
// by completion instead, defeating preserve_ordering.
//
// Only reserve when ordering is required: when just the response cache is
// enabled the delegator sends directly and would never fill the slot,
// leaking it for the lifetime of the model.
std::vector<std::pair<std::unique_ptr<InferenceResponse>, uint32_t>>*
queue_slot = nullptr;
if (preserve_ordering_) {
std::lock_guard<std::mutex> lock(completion_queue_mtx_);
completion_queue_.emplace_back();
queue_slot = &completion_queue_.back();
}
Comment on lines +687 to +691

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Slot reservation breaks arrival order

When an earlier cache miss remains in the dynamic-batching queue and a later request hits the cache, the hit reserves its completion slot immediately while the miss reserves only after BatcherThread dequeues it, causing the later response to occupy an earlier slot and violating scheduler-wide preserve_ordering.


// Cache plumbing
const std::string& key = request->CacheKey();
const bool is_key_set = request->CacheKeyIsSet();
const uint64_t lookup_end_ns = request->CacheLookupEndNs();
const uint64_t lookup_start_ns = request->CacheLookupStartNs();

request->SetResponseDelegator(
[this, key, is_key_set, lookup_end_ns, lookup_start_ns](
[this, queue_slot, key, is_key_set, lookup_end_ns, lookup_start_ns](
std::unique_ptr<InferenceResponse>&& response, const uint32_t flags) {
if (response_cache_enabled_) {
// Logical error, the key should be set if caching is enabled
Expand Down Expand Up @@ -731,8 +748,6 @@ DynamicBatchScheduler::DelegateResponse(
if (preserve_ordering_) {
{
std::lock_guard<std::mutex> lock(completion_queue_mtx_);
completion_queue_.emplace_back();
auto queue_slot = &completion_queue_.back();
queue_slot->emplace_back(std::move(response), flags);
}
FinalizeResponses();
Expand Down
Loading