From 160bd031b25fb93eda4b1ab1a86d860751b9e444 Mon Sep 17 00:00:00 2001 From: Xuan-Son Nguyen Date: Mon, 7 Sep 2026 15:50:46 +0200 Subject: [PATCH] server: fix LRU hang on multiple requests same model (#28539) * server: fix LRU hang on multiple requests same model * server: keep a queued model out of the victim pool until its waiters leave A waiter that gave up while its model was still loading left the model idle with no request behind it, and nothing recounted the free slots, so a second request queued behind it stayed queued forever. tick() was only driven by requests: join, claim and the end of a proxied request. Keep the queue entry alive after a successful claim so the model coming up is never picked as a victim before its waiters use it, and recount the slots on every status change and whenever a waiter abandons the queue. The model is then evicted as soon as it comes up with nobody left to serve. --------- Co-authored-by: Pascal --- tools/server/server-models.cpp | 165 ++++++++++--------------- tools/server/server-models.h | 4 + tools/server/tests/unit/test_router.py | 20 +++ 3 files changed, 92 insertions(+), 97 deletions(-) diff --git a/tools/server/server-models.cpp b/tools/server/server-models.cpp index db0fac995..4d2592b25 100644 --- a/tools/server/server-models.cpp +++ b/tools/server/server-models.cpp @@ -80,18 +80,19 @@ struct server_lru_sched { } // returns "" if no model can be given up - std::string pick_victim(std::unique_lock & lk, const std::string & exclude) { + std::string pick_victim(std::unique_lock & lk) { check_lock(lk); std::string victim; int64_t victim_last_used = 0; for (const auto & m : models.mapping) { - if (m.first == exclude) { - continue; - } // a busy model is mid-request, one still coming up has no request to finish if (m.second.req_count != 0 || !m.second.meta.is_ready_or_sleep()) { continue; } + // already on its way out, or a queued request wants it + if (models.stopping_models.count(m.first) || find(m.first)) { + continue; + } if (victim.empty() || m.second.meta.last_used < victim_last_used) { victim = m.first; victim_last_used = m.second.meta.last_used; @@ -109,7 +110,7 @@ struct server_lru_sched { SRV_INF("request for name=%s joined the queue, %d waiting\n", model_id.c_str(), e->n_waiters); return; } - queue.push_back({ model_id, 1, false, false }); + queue.push_back({ model_id, 1, false }); SRV_INF("models_max reached, request for name=%s queued at position %zu\n", model_id.c_str(), queue.size()); } @@ -144,85 +145,67 @@ struct server_lru_sched { return true; } - // ok means the model is up: drop the entry, the other waiters just watch its status now + // on failure the entry is back in line; on success it stays until its waiters leave, + // so the model coming up is never picked as a victim before they use it void claim_done(std::unique_lock & lk, const std::string & model_id, bool ok) { check_lock(lk); + if (ok) { + return; + } for (auto it = queue.begin(); it != queue.end(); ++it) { if (it->model_id == model_id) { - if (ok) { - queue.erase(it); - } else { - it->loading = false; - } + it->loading = false; return; } } } - // a model is on its way out for this entry, so other requests do not also give up one - void mark_slot_pending(std::unique_lock & lk, const std::string & model_id) { + // evict idle models while queued requests outnumber the slots that are free or being freed + // caller must hold models.mutex; never blocks, so it is safe from any thread + void tick(std::unique_lock & lk) { check_lock(lk); - if (entry_t * e = find(model_id)) { - e->slot_pending = true; + if (models.base_params.models_max <= 0 || queue.empty()) { + return; } - } - - // model_id went idle: give up its slot if a queued request needs one - // thread-safe, caller must NOT hold models.mutex - void on_model_idle(const std::string & model_id) { - if (models.base_params.models_max <= 0) { - return; // no limit, nothing is ever queued - } - { - std::unique_lock lk(models.mutex); - if (queue.empty()) { - return; - } - size_t promised = 0; - bool has_unserved = false; - for (const auto & e : queue) { - if (e.needs_slot()) { - has_unserved = true; - } else { - promised++; - } - } - if (!has_unserved) { - return; - } - if ((int) count_running() - (int) promised < models.base_params.models_max) { - return; // a slot is already on its way - } - // never give up a model that a queued request wants - for (const auto & e : queue) { - if (e.model_id == model_id) { - return; - } - } - auto it = models.mapping.find(model_id); - if (it == models.mapping.end() || it->second.req_count != 0 || !it->second.meta.is_ready_or_sleep()) { - return; - } - for (auto & e : queue) { - if (!e.slot_pending) { - e.slot_pending = true; - break; + int n_running = 0; + int n_stopping = 0; + for (const auto & m : models.mapping) { + if (m.second.meta.is_running()) { + n_running++; + if (models.stopping_models.count(m.first)) { + n_stopping++; } } } - SRV_INF("model name=%s went idle, giving up its slot to a queued request\n", model_id.c_str()); - models.unload(model_id); + int n_needed = 0; + int n_claimed = 0; // claimed the slot, but load() has not spawned yet + for (const auto & e : queue) { + if (!e.loading) { + n_needed++; + continue; + } + auto it = models.mapping.find(e.model_id); + if (it != models.mapping.end() && !it->second.meta.is_running()) { + n_claimed++; + } + } + int n_free = models.base_params.models_max - n_running + n_stopping - n_claimed; + while (n_free < n_needed) { + std::string victim = pick_victim(lk); + if (victim.empty()) { + return; // all remaining models are busy, wait for a request to end + } + SRV_INF("evicting idle LRU name=%s for a queued request\n", victim.c_str()); + models.request_stop(victim); + n_free++; + } } private: struct entry_t { std::string model_id; - int n_waiters; // requests waiting for this model - bool slot_pending; // a model is already being evicted for this entry - bool loading; // one of the waiters is doing the load right now - - // a slot is already coming, or already taken by the load in flight - bool needs_slot() const { return !slot_pending && !loading; } + int n_waiters; // requests waiting for this model + bool loading; // one of the waiters is doing the load right now }; entry_t * find(const std::string & model_id) { @@ -946,7 +929,7 @@ void server_models::unload_lru() { if (sched->has_capacity(lk)) { return; } - lru_model_name = sched->pick_victim(lk, ""); + lru_model_name = sched->pick_victim(lk); } if (!lru_model_name.empty()) { SRV_INF("models_max limit reached, removing LRU name=%s\n", lru_model_name.c_str()); @@ -1169,6 +1152,11 @@ void server_models::load(const std::string & name, const load_options & opts) { cv.notify_all(); } +void server_models::request_stop(const std::string & name) { + stopping_models.insert(name); + cv_stop.notify_all(); +} + void server_models::unload(const std::string & name) { std::unique_lock lk(mutex); auto it = mapping.find(name); @@ -1182,13 +1170,12 @@ void server_models::unload(const std::string & name) { }); } else if (it->second.meta.is_running()) { SRV_INF("stopping model instance name=%s\n", name.c_str()); - stopping_models.insert(name); if (it->second.meta.status == SERVER_MODEL_STATUS_LOADING) { // special case: if model is in loading state, unloading means force-killing it SRV_WRN("model name=%s is still loading, force-killing\n", name.c_str()); it->second.subproc->terminate(); } - cv_stop.notify_all(); + request_stop(name); // status change will be handled by the managing thread } else { SRV_WRN("model instance name=%s is not running\n", name.c_str()); @@ -1206,8 +1193,7 @@ void server_models::unload_all() { inst.subproc->stopped.store(true, std::memory_order_relaxed); } else if (inst.meta.is_running()) { SRV_INF("stopping model instance name=%s\n", name.c_str()); - stopping_models.insert(name); - cv_stop.notify_all(); + request_stop(name); // status change will be handled by the managing thread } // moving the thread to join list to avoid deadlock @@ -1234,6 +1220,8 @@ void server_models::update_status(const std::string & name, const update_status_ if (!args.progress.is_null()) { meta.progress = args.progress; } + // a model that comes up idle or goes down changes the slot count for queued requests + sched->tick(lk); } // broadcast status change to SSE { @@ -1380,13 +1368,11 @@ bool server_models::ensure_model_ready(const std::string & name, const std::func bool queued = false; bool did_load = false; - std::string victim; { std::unique_lock lk(mutex); auto it = mapping.find(name); if (it != mapping.end() && it->second.meta.status == SERVER_MODEL_STATUS_UNLOADED) { - bool has_capacity = sched->has_capacity(lk); - if (has_capacity && sched->queue_empty(lk)) { + if (sched->has_capacity(lk) && sched->queue_empty(lk)) { lk.unlock(); SRV_INF("model name=%s is not loaded, loading...\n", name.c_str()); load(name); @@ -1394,21 +1380,11 @@ bool server_models::ensure_model_ready(const std::string & name, const std::func } else { // also queue when a slot looks free but others wait already, else they starve sched->join(lk, name); + sched->tick(lk); queued = true; - if (!has_capacity) { - // an idle model may sit here right now, do not wait for a request to end - victim = sched->pick_victim(lk, name); - if (!victim.empty()) { - sched->mark_slot_pending(lk, name); - } - } } } } - if (!victim.empty()) { - SRV_INF("evicting idle LRU name=%s to make room for name=%s\n", victim.c_str(), name.c_str()); - unload(victim); - } // while queued, this is also where the load happens: the head of the queue does it SRV_INF("waiting until model name=%s is fully loaded...\n", name.c_str()); @@ -1470,9 +1446,7 @@ bool server_models::ensure_model_ready(const std::string & name, const std::func } lk.lock(); sched->claim_done(lk, name, ok); - if (ok) { - queued = false; // entry is gone, the other waiters watch the status now - } + sched->tick(lk); continue; } @@ -1480,6 +1454,7 @@ bool server_models::ensure_model_ready(const std::string & name, const std::func } } catch (...) { leave_queue(); + sched->tick(lk); // a slot freed for this waiter goes to the next one throw; } leave_queue(); @@ -1529,18 +1504,14 @@ server_http_res_ptr server_models::proxy_request(const server_http_req & req, co ); proxy->cleanup = [this, name]() { - bool went_idle = false; - { - std::unique_lock lk(mutex); - auto it = mapping.find(name); - if (it != mapping.end() && it->second.req_count > 0) { - it->second.req_count--; - went_idle = it->second.req_count == 0; + std::unique_lock lk(mutex); + auto it = mapping.find(name); + if (it != mapping.end() && it->second.req_count > 0) { + it->second.req_count--; + if (it->second.req_count == 0) { + sched->tick(lk); } } - if (went_idle) { - sched->on_model_idle(name); - } }; return proxy; diff --git a/tools/server/server-models.h b/tools/server/server-models.h index 5cbb6a801..7f6c26b35 100644 --- a/tools/server/server-models.h +++ b/tools/server/server-models.h @@ -216,6 +216,10 @@ private: // not thread-safe, caller must hold mutex void add_model(server_model_meta && meta); + // ask the monitoring thread to stop a running instance + // not thread-safe, caller must hold mutex + void request_stop(const std::string & name); + // notify SSE clients void notify_sse(const std::string & event, const std::string & model_id, const json & data = nullptr); diff --git a/tools/server/tests/unit/test_router.py b/tools/server/tests/unit/test_router.py index 96eb87978..e4b7f9fe4 100644 --- a/tools/server/tests/unit/test_router.py +++ b/tools/server/tests/unit/test_router.py @@ -297,6 +297,26 @@ def test_router_queue_is_fifo(): assert first.done_at < second.done_at, "queue was not served in arrival order" +def test_router_queue_two_waiters_share_one_eviction(): + """two requests that both find the same idle model must both be served in the end""" + global server + server.models_max = 1 + server.start() + + _load_model_and_wait(MODEL_A, timeout=120) + + # both arrive while MODEL_A is idle, so both want its slot; only one eviction can happen + first = _Bg(lambda: _tokenize(MODEL_B)).start() + second = _Bg(lambda: _tokenize(MODEL_C)).start() + + first.join(90) + second.join(90) + + first.assert_ok("first queued request") + second.assert_ok("second queued request") + assert _get_model_status(MODEL_A) == "unloaded" + + def test_router_no_models_autoload(): global server server.no_models_autoload = True