Brunobkr/llama.cpp_AlgMor24_github
ΩFFFΣLLIa • llama.cpp • AlgMor24 ██████╗ ███████╗███████╗███████╗██╗ ██╗ ██╗ █████╗ ██╔═══██╗██╔════╝██╔════╝██╔════╝██║ ██║ ██║██╔══██╗ ██║ ██║█████╗ █████╗ █████╗ ██║ ██║ ██║███████║ ██║ ██║██╔══╝ ██╔══╝ ██╔══╝ ██║ ██║ ██║██╔══██║ ╚██████╔╝██║ ██║ ███████╗███████╗███████╗██║██║ ██║ ╚═════╝ ╚═╝ ╚═╝ ╚══════╝╚══════╝╚══════╝╚═╝╚═╝ ╚═╝ High-Performance LLM / VLM Inference & Autonomous Agentic Ecosystem… See the full description on the dataset page: https://huggingface.co/datasets/Brunobkr/llama.cpp_AlgMor24_github.
03k
1#include "server-common.h"2#include "http.h"3#include "server-models.h"4#include "server-context.h"5#include "server-stream.h"6 7#include "build-info.h"8#include "preset.h"9#include "download.h"10#include "http.h"11#include "subproc.h"12 13#include <cpp-httplib/httplib.h> // TODO: remove this once we use HTTP client from download.h14#include <optional>15 16#include <functional>17#include <optional>18#include <algorithm>19#include <thread>20#include <mutex>21#include <condition_variable>22#include <cstring>23#include <cstdlib>24#include <atomic>25#include <chrono>26#include <queue>27#include <filesystem>28#include <random>29#include <sstream>30#include <cstring>31 32#ifndef _WIN3233extern char **environ;34#endif35 36#if defined(__APPLE__) && defined(__MACH__)37// macOS: use _NSGetExecutablePath to get the executable path38#include <mach-o/dyld.h>39#include <limits.h>40#endif41 42#define DEFAULT_STOP_TIMEOUT 10 // seconds43 44#define CMD_ROUTER_TO_CHILD_EXIT "cmd_router_to_child:exit"45#define CMD_CHILD_TO_ROUTER_STATE "cmd_child_to_router:state:" // followed by json string46 47// address for child process, this is needed because router may run on 0.0.0.048// ref: https://github.com/ggml-org/llama.cpp/issues/1786249#define CHILD_ADDR "127.0.0.1"50 51struct server_subproc {52 common_subproc sproc; // not yet spawned while in DOWNLOADING state53 std::atomic<bool> stopped{false}; // set to cancel a download or signal child process exit54 55 bool is_alive() {56 return sproc.alive();57 }58 59 void request_exit() {60 FILE * stdin_file = sproc.stdin_file();61 if (stdin_file) {62 fprintf(stdin_file, "%s\n", CMD_ROUTER_TO_CHILD_EXIT);63 fflush(stdin_file);64 }65 stopped.store(true, std::memory_order_relaxed);66 }67 68 void terminate() {69 sproc.terminate();70 }71};72 73struct server_lru_sched {74 server_lru_sched(server_models & models) : models(models) {}75 76 bool has_capacity(std::unique_lock<std::mutex> & lk) {77 check_lock(lk);78 return models.base_params.models_max <= 079 || count_running() < (size_t) models.base_params.models_max;80 }81 82 // returns "" if no model can be given up83 std::string pick_victim(std::unique_lock<std::mutex> & lk, const std::string & exclude) {84 check_lock(lk);85 std::string victim;86 int64_t victim_last_used = 0;87 for (const auto & m : models.mapping) {88 if (m.first == exclude) {89 continue;90 }91 // a busy model is mid-request, one still coming up has no request to finish92 if (m.second.req_count != 0 || !m.second.meta.is_ready_or_sleep()) {93 continue;94 }95 if (victim.empty() || m.second.meta.last_used < victim_last_used) {96 victim = m.first;97 victim_last_used = m.second.meta.last_used;98 }99 }100 return victim;101 }102 103 // requests wanting the same model share one entry, so they all need only one slot104 // and all get unblocked by the single load that entry performs105 void join(std::unique_lock<std::mutex> & lk, const std::string & model_id) {106 check_lock(lk);107 if (entry_t * e = find(model_id)) {108 e->n_waiters++;109 SRV_INF("request for name=%s joined the queue, %d waiting\n", model_id.c_str(), e->n_waiters);110 return;111 }112 queue.push_back({ model_id, 1, false, false });113 SRV_INF("models_max reached, request for name=%s queued at position %zu\n",114 model_id.c_str(), queue.size());115 }116 117 void leave(std::unique_lock<std::mutex> & lk, const std::string & model_id) {118 check_lock(lk);119 for (auto it = queue.begin(); it != queue.end(); ++it) {120 if (it->model_id == model_id) {121 if (--it->n_waiters <= 0) {122 queue.erase(it); // last one waiting for this model went away123 }124 return;125 }126 }127 }128 129 bool queue_empty(std::unique_lock<std::mutex> & lk) {130 check_lock(lk);131 return queue.empty();132 }133 134 // true if it is this model's turn to load, and nobody is loading it yet135 bool try_claim(std::unique_lock<std::mutex> & lk, const std::string & model_id) {136 check_lock(lk);137 if (queue.empty() || queue.front().model_id != model_id || queue.front().loading) {138 return false;139 }140 if (!has_capacity(lk)) {141 return false;142 }143 queue.front().loading = true;144 return true;145 }146 147 // ok means the model is up: drop the entry, the other waiters just watch its status now148 void claim_done(std::unique_lock<std::mutex> & lk, const std::string & model_id, bool ok) {149 check_lock(lk);150 for (auto it = queue.begin(); it != queue.end(); ++it) {151 if (it->model_id == model_id) {152 if (ok) {153 queue.erase(it);154 } else {155 it->loading = false;156 }157 return;158 }159 }160 }161 162 // a model is on its way out for this entry, so other requests do not also give up one163 void mark_slot_pending(std::unique_lock<std::mutex> & lk, const std::string & model_id) {164 check_lock(lk);165 if (entry_t * e = find(model_id)) {166 e->slot_pending = true;167 }168 }169 170 // model_id went idle: give up its slot if a queued request needs one171 // thread-safe, caller must NOT hold models.mutex172 void on_model_idle(const std::string & model_id) {173 if (models.base_params.models_max <= 0) {174 return; // no limit, nothing is ever queued175 }176 {177 std::unique_lock<std::mutex> lk(models.mutex);178 if (queue.empty()) {179 return;180 }181 size_t promised = 0;182 bool has_unserved = false;183 for (const auto & e : queue) {184 if (e.needs_slot()) {185 has_unserved = true;186 } else {187 promised++;188 }189 }190 if (!has_unserved) {191 return;192 }193 if ((int) count_running() - (int) promised < models.base_params.models_max) {194 return; // a slot is already on its way195 }196 // never give up a model that a queued request wants197 for (const auto & e : queue) {198 if (e.model_id == model_id) {199 return;200 }201 }202 auto it = models.mapping.find(model_id);203 if (it == models.mapping.end() || it->second.req_count != 0 || !it->second.meta.is_ready_or_sleep()) {204 return;205 }206 for (auto & e : queue) {207 if (!e.slot_pending) {208 e.slot_pending = true;209 break;210 }211 }212 }213 SRV_INF("model name=%s went idle, giving up its slot to a queued request\n", model_id.c_str());214 models.unload(model_id);215 }216 217 private:218 struct entry_t {219 std::string model_id;220 int n_waiters; // requests waiting for this model221 bool slot_pending; // a model is already being evicted for this entry222 bool loading; // one of the waiters is doing the load right now223 224 // a slot is already coming, or already taken by the load in flight225 bool needs_slot() const { return !slot_pending && !loading; }226 };227 228 entry_t * find(const std::string & model_id) {229 for (auto & e : queue) {230 if (e.model_id == model_id) {231 return &e;232 }233 }234 return nullptr;235 }236 237 void check_lock(std::unique_lock<std::mutex> & lk) {238 GGML_ASSERT(lk.owns_lock() && lk.mutex() == &models.mutex);239 }240 241 size_t count_running() {242 size_t count = 0;243 for (const auto & m : models.mapping) {244 if (m.second.meta.is_running()) {245 count++;246 }247 }248 return count;249 }250 251 server_models & models;252 std::deque<entry_t> queue;253};254 255// short loopback budget for the resumable stream router to child JSON calls (probe, lookup,256// delete). distinct from params.timeout_read/write which only applies to the generation proxy257static constexpr int STREAM_LOOKUP_TIMEOUT_MS = 250;258 259static std::filesystem::path get_server_exec_path() {260#if defined(_WIN32)261 wchar_t buf[32768] = { 0 }; // Large buffer to handle long paths262 DWORD len = GetModuleFileNameW(nullptr, buf, _countof(buf));263 if (len == 0 || len >= _countof(buf)) {264 throw std::runtime_error("GetModuleFileNameW failed or path too long");265 }266 return std::filesystem::path(buf);267#elif defined(__APPLE__) && defined(__MACH__)268 char small_path[PATH_MAX];269 uint32_t size = sizeof(small_path);270 271 if (_NSGetExecutablePath(small_path, &size) == 0) {272 // resolve any symlinks to get absolute path273 try {274 return std::filesystem::canonical(std::filesystem::path(small_path));275 } catch (...) {276 return std::filesystem::path(small_path);277 }278 } else {279 // buffer was too small, allocate required size and call again280 std::vector<char> buf(size);281 if (_NSGetExecutablePath(buf.data(), &size) == 0) {282 try {283 return std::filesystem::canonical(std::filesystem::path(buf.data()));284 } catch (...) {285 return std::filesystem::path(buf.data());286 }287 }288 throw std::runtime_error("_NSGetExecutablePath failed after buffer resize");289 }290#else291 char path[FILENAME_MAX];292 ssize_t count = readlink("/proc/self/exe", path, FILENAME_MAX);293 if (count <= 0) {294 throw std::runtime_error("failed to resolve /proc/self/exe");295 }296 return std::filesystem::path(std::string(path, count));297#endif298}299 300static void unset_reserved_args(common_preset & preset, bool unset_model_args) {301 preset.unset_option("LLAMA_ARG_SSL_KEY_FILE");302 preset.unset_option("LLAMA_ARG_SSL_CERT_FILE");303 preset.unset_option("LLAMA_API_KEY");304 preset.unset_option("LLAMA_ARG_MODELS_DIR");305 preset.unset_option("LLAMA_ARG_MODELS_MAX");306 preset.unset_option("LLAMA_ARG_MODELS_PRESET");307 preset.unset_option("LLAMA_ARG_MODELS_AUTOLOAD");308 if (unset_model_args) {309 preset.unset_option("LLAMA_ARG_MODEL");310 preset.unset_option("LLAMA_ARG_MMPROJ");311 preset.unset_option("LLAMA_ARG_ALIAS");312 preset.unset_option("LLAMA_ARG_HF_REPO");313 }314}315 316#ifdef _WIN32317static std::string wide_to_utf8(const wchar_t * ws) {318 if (!ws || !*ws) {319 return {};320 }321 322 const int len = static_cast<int>(std::wcslen(ws));323 const int bytes = WideCharToMultiByte(CP_UTF8, 0, ws, len, nullptr, 0, nullptr, nullptr);324 if (bytes == 0) {325 return {};326 }327 328 std::string utf8(bytes, '\0');329 WideCharToMultiByte(CP_UTF8, 0, ws, len, utf8.data(), bytes, nullptr, nullptr);330 331 return utf8;332}333#endif334 335static std::vector<std::string> get_environment() {336 std::vector<std::string> env;337 338#ifdef _WIN32339 LPWCH env_block = GetEnvironmentStringsW();340 if (!env_block) {341 return env;342 }343 for (LPWCH e = env_block; *e; e += wcslen(e) + 1) {344 env.emplace_back(wide_to_utf8(e));345 }346 FreeEnvironmentStringsW(env_block);347#else348 if (environ == nullptr) {349 return env;350 }351 for (char ** e = environ; *e != nullptr; e++) {352 env.emplace_back(*e);353 }354#endif355 356 return env;357}358 359void server_model_meta::update_args(common_preset_context & ctx_preset, std::string bin_path) {360 // update params361 unset_reserved_args(preset, false);362 preset.set_option(ctx_preset, "LLAMA_ARG_HOST", CHILD_ADDR);363 preset.set_option(ctx_preset, "LLAMA_ARG_PORT", std::to_string(port));364 preset.set_option(ctx_preset, "LLAMA_ARG_ALIAS", name);365 // TODO: maybe validate preset before rendering ?366 // render args367 args = preset.to_args(bin_path);368 369 // unified binary dispatches by subcommand, re-inject it right after the370 // binary path so the child starts as 'llama serve ...' not 'llama ...'371 const char * app_cmd = std::getenv("LLAMA_APP_CMD");372 if (app_cmd != nullptr && app_cmd[0] != '\0' && !bin_path.empty()) {373 args.insert(args.begin() + 1, app_cmd);374 }375}376 377void server_model_meta::update_caps() {378 try {379 common_params params;380 preset.apply_to_params(params, {381 "LLAMA_ARG_MODEL",382 "LLAMA_ARG_MODEL_URL",383 "LLAMA_ARG_MMPROJ",384 "LLAMA_ARG_MMPROJ_URL",385 "LLAMA_ARG_MMPROJ_AUTO",386 "LLAMA_ARG_HF_REPO",387 "LLAMA_ARG_HF_REPO_FILE",388 });389 params.offline = true;390 common_models_handler handler = common_models_handler_init(params, LLAMA_EXAMPLE_SERVER);391 common_models_handler_apply(handler, params); // note: this won't download the model because offline=true392 if (params.no_mmproj || params.mmproj.path.empty()) {393 multimodal = { false, false };394 } else {395 multimodal = mtmd_get_cap_from_file(params.mmproj.path.c_str());396 }397 } catch (const std::exception & e) {398 LOG_WRN("failed to initialize common_params for multimodal capability detection: %s\n", e.what());399 multimodal = { false, false };400 }401}402 403//404// server_models405//406 407server_models::server_models(408 const common_params & params,409 int argc,410 char ** argv)411 : ctx_preset(LLAMA_EXAMPLE_SERVER),412 base_params(params),413 base_env(get_environment()),414 base_preset(ctx_preset.load_from_args(argc, argv)),415 sched(std::make_unique<server_lru_sched>(*this)) {416 // clean up base preset417 unset_reserved_args(base_preset, true);418 // set binary path419 try {420 bin_path = get_server_exec_path().string();421 } catch (const std::exception & e) {422 bin_path = argv[0];423 LOG_WRN("failed to get server executable path: %s\n", e.what());424 LOG_WRN("using original argv[0] as fallback: %s\n", argv[0]);425 }426 load_models();427 debug_fake_timing = !common_get_env("LLAMA_SERVER_DEBUG_FAKE_TIMING").empty();428}429 430server_models::~server_models() = default;431 432void server_models::add_model(server_model_meta && meta) {433 if (mapping.find(meta.name) != mapping.end()) {434 throw std::runtime_error(string_format("model '%s' appears multiple times", meta.name.c_str()));435 }436 437 // check model name does not conflict with existing aliases438 for (const auto & [key, inst] : mapping) {439 if (inst.meta.aliases.count(meta.name)) {440 throw std::runtime_error(string_format("model name '%s' conflicts with alias of model '%s'",441 meta.name.c_str(), key.c_str()));442 }443 }444 445 // parse aliases from preset's --alias option (comma-separated)446 std::string alias_str;447 if (meta.preset.get_option("LLAMA_ARG_ALIAS", alias_str) && !alias_str.empty()) {448 for (auto & alias : string_split<std::string>(alias_str, ',')) {449 alias = string_strip(alias);450 if (!alias.empty()) {451 meta.aliases.insert(alias);452 }453 }454 }455 456 // parse tags from preset's --tags option (comma-separated)457 std::string tags_str;458 if (meta.preset.get_option("LLAMA_ARG_TAGS", tags_str) && !tags_str.empty()) {459 for (auto & tag : string_split<std::string>(tags_str, ',')) {460 tag = string_strip(tag);461 if (!tag.empty()) {462 meta.tags.insert(tag);463 }464 }465 }466 467 // validate aliases do not conflict with existing names or aliases468 for (const auto & alias : meta.aliases) {469 if (mapping.find(alias) != mapping.end()) {470 throw std::runtime_error(string_format("alias '%s' for model '%s' conflicts with existing model name",471 alias.c_str(), meta.name.c_str()));472 }473 for (const auto & [key, inst] : mapping) {474 if (inst.meta.aliases.count(alias)) {475 throw std::runtime_error(string_format("alias '%s' for model '%s' conflicts with alias of model '%s'",476 alias.c_str(), meta.name.c_str(), key.c_str()));477 }478 }479 }480 481 meta.update_args(ctx_preset, bin_path); // render args482 meta.update_caps();483 std::string name = meta.name;484 mapping[name] = instance_t{485 /* subproc */ std::make_shared<server_subproc>(),486 /* th */ std::thread(),487 /* meta */ std::move(meta)488 };489}490 491void server_models::notify_sse(const std::string & event, const std::string & model_id, const json & data) {492 std::unique_ptr<server_task_result_router> result = std::make_unique<server_task_result_router>();493 result->data = {494 {"model", model_id},495 {"event", event},496 };497 if (!data.is_null()) {498 result->data["data"] = data;499 }500 SRV_DBG("notifying SSE clients about event '%s' for model '%s': %s\n", event.c_str(), model_id.c_str(), safe_json_to_str(result->data).c_str());501 sse.broadcast(std::move(result));502}503 504void server_models::load_models() {505 // Phase 1: load presets from all sources - pure I/O, no lock needed506 // 1. cached models507 common_presets cached_models = ctx_preset.load_from_cache();508 SRV_INF("Loaded %zu cached model presets\n", cached_models.size());509 // 2. local models from --models-dir510 common_presets local_models;511 if (!base_params.models_dir.empty()) {512 local_models = ctx_preset.load_from_models_dir(base_params.models_dir);513 SRV_INF("Loaded %zu local model presets from %s\n", local_models.size(), base_params.models_dir.c_str());514 }515 // 3. custom-path models from presets516 common_preset global = {};517 common_presets custom_presets = {};518 if (!base_params.models_preset.empty()) {519 custom_presets = ctx_preset.load_from_ini(base_params.models_preset, global);520 SRV_INF("Loaded %zu custom model presets from %s\n", custom_presets.size(), base_params.models_preset.c_str());521 }522 523 // cascade, apply global preset first524 cached_models = ctx_preset.cascade(global, cached_models);525 local_models = ctx_preset.cascade(global, local_models);526 custom_presets = ctx_preset.cascade(global, custom_presets);527 528 // note: if a model exists in both cached and local, local takes precedence529 common_presets final_presets;530 std::unordered_map<std::string, server_model_source> source_map;531 for (const auto & [name, preset] : cached_models) {532 final_presets[name] = preset;533 source_map[name] = SERVER_MODEL_SOURCE_CACHE;534 }535 for (const auto & [name, preset] : local_models) {536 final_presets[name] = preset;537 source_map[name] = SERVER_MODEL_SOURCE_MODELS_DIR;538 }539 for (const auto & [name, custom] : custom_presets) {540 if (final_presets.find(name) != final_presets.end()) {541 final_presets[name].merge(custom);542 } else {543 final_presets[name] = custom;544 }545 source_map[name] = SERVER_MODEL_SOURCE_PRESET;546 }547 548 // overlay router's own CLI args on top of every model preset so that549 // e.g. `llama-server --temp 0` is honoured by all child processes550 for (auto & [name, preset] : final_presets) {551 preset.merge(base_preset);552 }553 554 auto get_source = [&](const std::string & name) {555 return source_map.count(name) ? source_map.at(name) : SERVER_MODEL_SOURCE_PRESET;556 };557 558 // Helpers that read `mapping` - must be called while holding the lock.559 std::unordered_set<std::string> custom_names;560 for (const auto & [name, preset] : custom_presets) custom_names.insert(name);561 auto join_set = [](const std::set<std::string> & s) {562 std::string result;563 for (const auto & v : s) {564 if (!result.empty()) result += ", ";565 result += v;566 }567 return result;568 };569 auto log_available_models = [&]() {570 SRV_INF("Available models (%zu) (*: custom preset)\n", mapping.size());571 for (const auto & [name, inst] : mapping) {572 bool has_custom = custom_names.find(name) != custom_names.end();573 std::string info;574 if (!inst.meta.aliases.empty()) info += " (aliases: " + join_set(inst.meta.aliases) + ")";575 if (!inst.meta.tags.empty()) info += " [tags: " + join_set(inst.meta.tags) + "]";576 SRV_INF(" %c %s%s\n", has_custom ? '*' : ' ', name.c_str(), info.c_str());577 }578 };579 auto apply_stop_timeout = [&]() {580 for (auto & [name, inst] : mapping) {581 std::string val;582 if (inst.meta.preset.get_option(COMMON_ARG_PRESET_STOP_TIMEOUT, val)) {583 try {584 inst.meta.stop_timeout = std::stoi(val);585 } catch (...) {586 SRV_WRN("invalid stop-timeout value '%s' for model '%s', using default %d seconds\n",587 val.c_str(), name.c_str(), DEFAULT_STOP_TIMEOUT);588 inst.meta.stop_timeout = DEFAULT_STOP_TIMEOUT;589 }590 }591 }592 };593 // update_args() injects HOST/PORT/ALIAS, so strip them before comparing presets594 auto preset_options_for_compare = [](common_preset p) {595 p.unset_option("LLAMA_ARG_HOST");596 p.unset_option("LLAMA_ARG_PORT");597 p.unset_option("LLAMA_ARG_ALIAS");598 return p.options;599 };600 601 // Phase 2: acquire the lock once for all mapping mutations.602 // We temporarily release it only when calling functions that acquire it internally603 // (unload, load) or when joining threads (the monitoring thread calls update_status604 // which locks the mutex, so joining while holding it would deadlock).605 std::unique_lock<std::mutex> lk(mutex);606 607 need_reload = false;608 bool is_first_load = mapping.empty();609 610 if (is_first_load) {611 // FIRST LOAD: add all models, then unlock for autoloading612 for (const auto & [name, preset] : final_presets) {613 server_model_meta meta{614 /* source */ get_source(name),615 /* preset */ preset,616 /* name */ name,617 /* aliases */ {},618 /* tags */ {},619 /* port */ 0,620 /* status */ SERVER_MODEL_STATUS_UNLOADED,621 /* last_used */ 0,622 /* args */ std::vector<std::string>(),623 /* loaded_info */ {},624 /* progress */ {},625 /* exit_code */ 0,626 /* stop_timeout */ DEFAULT_STOP_TIMEOUT,627 /* multimodal */ mtmd_caps{false, false},628 // /* need_download */ false,629 };630 add_model(std::move(meta));631 }632 apply_stop_timeout();633 log_available_models();634 635 std::vector<std::string> models_to_load;636 for (const auto & [name, inst] : mapping) {637 std::string val;638 if (inst.meta.preset.get_option(COMMON_ARG_PRESET_LOAD_ON_STARTUP, val) && common_arg_utils::is_truthy(val)) {639 models_to_load.push_back(name);640 }641 }642 if ((int)models_to_load.size() > base_params.models_max) {643 throw std::runtime_error(string_format(644 "number of models to load on startup (%zu) exceeds models_max (%d)",645 models_to_load.size(), base_params.models_max));646 }647 648 lk.unlock();649 for (const auto & name : models_to_load) {650 SRV_INF("(startup) loading model %s\n", name.c_str());651 load(name);652 }653 } else {654 // RELOAD: diff the new preset list against the current mapping and reconcile655 is_reloading = true;656 657 // find running models whose source was removed or whose preset changed658 std::vector<std::string> to_unload;659 for (const auto & [name, inst] : mapping) {660 if (!inst.meta.is_running()) continue;661 auto it = final_presets.find(name);662 if (it == final_presets.end()) {663 to_unload.push_back(name); // removed from source664 } else if (preset_options_for_compare(inst.meta.preset) != preset_options_for_compare(it->second)) {665 to_unload.push_back(name); // preset changed666 }667 }668 669 // unload() acquires the lock internally, so release before each call670 for (const auto & name : to_unload) {671 SRV_INF("(reload) unloading model name=%s (source updated or removed)\n", name.c_str());672 lk.unlock();673 unload(name);674 lk.lock();675 }676 677 // wait for all targeted models to reach UNLOADED; cv.wait handles unlock/relock678 cv.wait(lk, [&]() {679 for (const auto & name : to_unload) {680 auto it = mapping.find(name);681 if (it != mapping.end() && it->second.meta.is_running()) return false;682 }683 return true;684 });685 686 // collect all threads to join in one pass while the lock is held:687 // - monitoring threads from just-unloaded models (to_unload)688 // - threads of finished downloads (DOWNLOADED), they acquire the mutex on exit689 // - threads of already-UNLOADED models that are being removed from source690 std::vector<std::thread> threads_to_join;691 for (const auto & name : to_unload) {692 auto it = mapping.find(name);693 if (it != mapping.end() && it->second.th.joinable()) {694 threads_to_join.push_back(std::move(it->second.th));695 }696 }697 for (auto & [name, inst] : mapping) {698 if (inst.meta.status == SERVER_MODEL_STATUS_DOWNLOADING) {699 continue; // downloading models are not from config sources, leave them alone700 }701 if (inst.meta.status == SERVER_MODEL_STATUS_DOWNLOADED) {702 // joining this thread under the lock deadlocks: it locks the mutex on its way out703 if (inst.th.joinable()) {704 threads_to_join.push_back(std::move(inst.th));705 }706 continue;707 }708 if (final_presets.find(name) == final_presets.end() && !inst.meta.is_running() && inst.th.joinable()) {709 threads_to_join.push_back(std::move(inst.th));710 }711 }712 713 // join outside the lock - monitoring thread calls update_status (needs lock)714 lk.unlock();715 for (auto & th : threads_to_join) th.join();716 lk.lock();717 718 // erase models no longer in any source719 for (auto it = mapping.begin(); it != mapping.end(); ) {720 if (it->second.meta.status == SERVER_MODEL_STATUS_DOWNLOADING) {721 ++it; // download thread is still busy, skip722 } else if (it->second.meta.status == SERVER_MODEL_STATUS_DOWNLOADED) {723 // download finished, thread is joined above, safe to erase724 GGML_ASSERT(!it->second.th.joinable());725 it = mapping.erase(it);726 } else if (final_presets.find(it->first) == final_presets.end()) {727 SRV_INF("(reload) removing model name=%s (no longer in source)\n", it->first.c_str());728 GGML_ASSERT(!it->second.th.joinable()); // must have been joined above729 it = mapping.erase(it);730 } else {731 ++it;732 }733 }734 735 // update presets for non-running models still in source736 for (auto & [name, inst] : mapping) {737 if (inst.meta.is_running()) continue;738 auto it = final_presets.find(name);739 if (it == final_presets.end()) continue; // erased above740 741 inst.meta.preset = it->second;742 743 // re-parse aliases, then validate against other models744 std::set<std::string> new_aliases;745 std::string alias_str;746 if (inst.meta.preset.get_option("LLAMA_ARG_ALIAS", alias_str) && !alias_str.empty()) {747 for (auto & alias : string_split<std::string>(alias_str, ',')) {748 alias = string_strip(alias);749 if (!alias.empty()) new_aliases.insert(alias);750 }751 }752 inst.meta.aliases.clear();753 for (const auto & alias : new_aliases) {754 bool conflict = false;755 for (const auto & [other_name, other_inst] : mapping) {756 if (other_name == name) continue;757 if (other_name == alias || other_inst.meta.aliases.count(alias)) {758 SRV_WRN("(reload) alias '%s' for model '%s' conflicts with model '%s', skipping\n",759 alias.c_str(), name.c_str(), other_name.c_str());760 conflict = true;761 break;762 }763 }764 if (!conflict) inst.meta.aliases.insert(alias);765 }766 767 // re-parse tags768 inst.meta.tags.clear();769 std::string tags_str;770 if (inst.meta.preset.get_option("LLAMA_ARG_TAGS", tags_str) && !tags_str.empty()) {771 for (auto & tag : string_split<std::string>(tags_str, ',')) {772 tag = string_strip(tag);773 if (!tag.empty()) inst.meta.tags.insert(tag);774 }775 }776 777 inst.meta.exit_code = 0; // clear failed state so the model can be reloaded778 inst.meta.update_args(ctx_preset, bin_path);779 inst.meta.update_caps();780 }781 782 // add models that are new in this reload783 std::vector<std::string> newly_added;784 for (const auto & [name, preset] : final_presets) {785 if (mapping.find(name) == mapping.end()) {786 server_model_meta meta{787 /* source */ get_source(name),788 /* preset */ preset,789 /* name */ name,790 /* aliases */ {},791 /* tags */ {},792 /* port */ 0,793 /* status */ SERVER_MODEL_STATUS_UNLOADED,794 /* last_used */ 0,795 /* args */ std::vector<std::string>(),796 /* loaded_info */ {},797 /* progress */ {},798 /* exit_code */ 0,799 /* stop_timeout */ DEFAULT_STOP_TIMEOUT,800 /* multimodal */ mtmd_caps{false, false},801 // /* need_download */ false,802 };803 add_model(std::move(meta));804 newly_added.push_back(name);805 }806 }807 808 apply_stop_timeout();809 810 // clear reload flag before unlocking for autoload - load() blocks on !is_reloading,811 // so clearing it here (while still locked) prevents a deadlock in the autoload calls below812 is_reloading = false;813 cv.notify_all();814 815 log_available_models();816 817 // collect autoload candidates while still under the lock818 std::vector<std::string> to_autoload;819 for (const auto & name : newly_added) {820 auto it = mapping.find(name);821 if (it != mapping.end()) {822 std::string val;823 if (it->second.meta.preset.get_option(COMMON_ARG_PRESET_LOAD_ON_STARTUP, val) && common_arg_utils::is_truthy(val)) {824 to_autoload.push_back(name);825 }826 }827 }828 829 lk.unlock();830 for (const auto & name : to_autoload) {831 SRV_INF("(reload) loading new model %s\n", name.c_str());832 load(name);833 }834 835 notify_sse("models_reload", "*");836 }837}838 839void server_models::update_meta(const std::string & name, const server_model_meta & meta) {840 std::lock_guard<std::mutex> lk(mutex);841 auto it = mapping.find(name);842 if (it != mapping.end()) {843 it->second.meta = meta;844 }845 cv.notify_all(); // notify wait_until_loading_finished846}847 848bool server_models::has_model(const std::string & name) {849 std::lock_guard<std::mutex> lk(mutex);850 if (mapping.find(name) != mapping.end()) {851 return true;852 }853 for (const auto & [key, inst] : mapping) {854 if (inst.meta.aliases.count(name)) {855 return true;856 }857 }858 return false;859}860 861std::optional<server_model_meta> server_models::get_meta(const std::string & name) {862 std::unique_lock<std::mutex> lk(mutex);863 if (need_reload) {864 lk.unlock();865 load_models();866 lk.lock();867 }868 869 auto it = mapping.find(name);870 if (it != mapping.end()) {871 return it->second.meta;872 }873 for (const auto & [key, inst] : mapping) {874 if (inst.meta.aliases.count(name)) {875 return inst.meta;876 }877 }878 return std::nullopt;879}880 881std::vector<server_model_meta> server_models::get_all_meta() {882 std::unique_lock<std::mutex> lk(mutex);883 if (need_reload) {884 lk.unlock();885 load_models();886 lk.lock();887 }888 889 std::vector<server_model_meta> result;890 result.reserve(mapping.size());891 for (const auto & [name, inst] : mapping) {892 result.push_back(inst.meta);893 }894 return result;895}896 897void server_models::unload_lru() {898 if (base_params.models_max <= 0) {899 return; // no limit900 }901 // remove one of the servers if we passed the models_max (least recently used - LRU)902 std::string lru_model_name;903 {904 std::unique_lock<std::mutex> lk(mutex);905 if (sched->has_capacity(lk)) {906 return;907 }908 lru_model_name = sched->pick_victim(lk, "");909 }910 if (!lru_model_name.empty()) {911 SRV_INF("models_max limit reached, removing LRU name=%s\n", lru_model_name.c_str());912 unload(lru_model_name);913 // wait for unload to complete914 {915 std::unique_lock<std::mutex> lk(mutex);916 cv.wait(lk, [this, &lru_model_name]() {917 return mapping[lru_model_name].meta.status == SERVER_MODEL_STATUS_UNLOADED;918 });919 }920 }921}922 923void server_models::load(const std::string & name) {924 load(name, load_options{});925}926 927void server_models::load(const std::string & name, const load_options & opts) {928 if (debug_fake_timing) {929 // do not hold the mutex here, other requests must keep making progress930 std::this_thread::sleep_for(std::chrono::seconds(2));931 }932 933 if (!opts.custom_meta.has_value()) {934 if (!has_model(name)) {935 throw std::runtime_error("model name=" + name + " is not found");936 }937 unload_lru();938 }939 940 std::unique_lock<std::mutex> lk(mutex);941 // edge case: block until any in-progress reload has finished so we always load942 // against the freshest preset and a consistent mapping state943 cv.wait(lk, [this]() { return !is_reloading; });944 945 auto meta = opts.custom_meta.has_value() ? *opts.custom_meta : mapping[name].meta;946 if (meta.status != SERVER_MODEL_STATUS_UNLOADED) {947 SRV_INF("model %s is not ready\n", name.c_str());948 return;949 }950 951 // Re-check capacity under the lock to prevent concurrent loads from952 // exceeding models_max. Without this, the window between unload_lru()953 // releasing its lock and this lock_guard acquiring allows multiple954 // threads to each observe capacity and all proceed to load.955 if (base_params.models_max > 0) {956 size_t count_active = 0;957 for (const auto & m : mapping) {958 if (m.second.meta.is_running()) {959 count_active++;960 }961 }962 if (count_active >= (size_t)base_params.models_max) {963 throw std::runtime_error("model limit reached, try again later");964 }965 }966 967 // prepare new instance info968 instance_t inst;969 inst.meta = meta;970 inst.meta.port = common_http_get_free_port();971 inst.meta.status = SERVER_MODEL_STATUS_LOADING;972 inst.meta.loaded_info = json{};973 inst.meta.last_used = ggml_time_ms();974 975 if (inst.meta.port <= 0) {976 throw std::runtime_error("failed to get a port number");977 }978 979 inst.subproc = std::make_shared<server_subproc>();980 {981 SRV_INF("spawning server instance with name=%s on port %d\n", inst.meta.name.c_str(), inst.meta.port);982 983 inst.meta.update_args(ctx_preset, bin_path); // render args984 985 std::vector<std::string> child_args = inst.meta.args; // copy986 std::vector<std::string> child_env = base_env; // copy987 child_env.push_back("LLAMA_SERVER_ROUTER_PORT=" + std::to_string(base_params.port));988 989 if (opts.mode == SERVER_CHILD_MODE_DOWNLOAD) {990 inst.meta.status = SERVER_MODEL_STATUS_DOWNLOADING;991 child_env.push_back("LLAMA_SERVER_CHILD_MODE=download");992 child_env.push_back("LLAMA_ARG_HF_REPO=" + name);993 }994 995 SRV_INF("%s", "spawning server instance with args:\n");996 for (const auto & arg : child_args) {997 SRV_INF(" %s\n", arg.c_str());998 }999 inst.meta.args = child_args; // save for debugging1000 1001 // TODO @ngxson : maybe separate stdout and stderr in the future1002 // so that we can use stdout for commands and stderr for logging1003 int options = subprocess_option_no_window | subprocess_option_combined_stdout_stderr;1004 if (!inst.subproc->sproc.create(child_args, options, child_env)) {1005 throw std::runtime_error("failed to spawn server instance");1006 }1007 }1008 1009 // start a thread to manage the child process1010 // captured variables are guaranteed to be destroyed only after the thread is joined1011 inst.th = std::thread([1012 this, name,1013 child_proc = inst.subproc,1014 port = inst.meta.port,1015 stop_timeout = inst.meta.stop_timeout,1016 child_mode = opts.mode1017 ]() {1018 FILE * stdin_file = child_proc->sproc.stdin_file();1019 FILE * stdout_file = child_proc->sproc.stdout_file(); // combined stdout/stderr1020 1021 std::thread log_thread([&]() {1022 // read stdout/stderr and forward to main server log1023 // also handle status report from child process1024 std::vector<char> vec_buf(128 * 1024); // large buffer for storing info1025 char * buffer = vec_buf.data();1026 if (stdout_file) {1027 while (fgets(buffer, vec_buf.size(), stdout_file) != nullptr) {1028 LOG("[%5d] %s", port, buffer);1029 std::string str(buffer);1030 if (string_starts_with(buffer, CMD_CHILD_TO_ROUTER_STATE)) {1031 this->handle_child_state(name, str);1032 }1033 }1034 } else {1035 SRV_ERR("failed to get stdout/stderr of child process for name=%s\n", name.c_str());1036 }1037 });1038 1039 std::thread stopping_thread([&]() {1040 // thread to monitor explicit stop requests; child crash is signalled via child_proc->stopped1041 auto is_stopping = [this, &name]() {1042 return this->stopping_models.find(name) != this->stopping_models.end();1043 };1044 {1045 std::unique_lock<std::mutex> lk(this->mutex);1046 this->cv_stop.wait(lk, [&]() {1047 return is_stopping() || child_proc->stopped.load(std::memory_order_acquire);1048 });1049 }1050 // child crashed or finished on its own, skip graceful shutdown sequence1051 if (child_proc->stopped.load(std::memory_order_acquire)) {1052 return;1053 }1054 SRV_INF("stopping model instance name=%s\n", name.c_str());1055 fprintf(stdin_file, "%s\n", CMD_ROUTER_TO_CHILD_EXIT);1056 fflush(stdin_file);1057 int64_t start_time = ggml_time_ms();1058 while (true) {1059 std::unique_lock<std::mutex> lk(this->mutex);1060 if (!is_stopping() || child_proc->stopped.load(std::memory_order_acquire)) {1061 return;1062 }1063 int64_t elapsed = ggml_time_ms() - start_time;1064 if (elapsed >= stop_timeout * 1000) {1065 lk.unlock();1066 SRV_WRN("force-killing model instance name=%s after %d seconds timeout\n", name.c_str(), stop_timeout);1067 child_proc->terminate();1068 return;1069 }1070 this->cv_stop.wait_for(lk, std::chrono::seconds(1), [&]() {1071 return !is_stopping() || child_proc->stopped.load(std::memory_order_acquire);1072 });1073 }1074 });1075 1076 // we reach here when the child process exits (stdout EOF)1077 // note: we cannot join() prior to this point because it will close stdin_file1078 if (log_thread.joinable()) {1079 log_thread.join();1080 }1081 1082 child_proc->stopped.store(true, std::memory_order_release);1083 {1084 std::lock_guard<std::mutex> lk(this->mutex);1085 stopping_models.erase(name);1086 cv_stop.notify_all();1087 }1088 if (stopping_thread.joinable()) {1089 stopping_thread.join();1090 }1091 1092 // get the exit code1093 int exit_code = child_proc->sproc.join();1094 1095 // update status and exit code1096 if (child_mode == SERVER_CHILD_MODE_DOWNLOAD) {1097 // instance will be cleaned up on next load_models() call1098 } else {1099 this->update_status(name, {1100 SERVER_MODEL_STATUS_UNLOADED,1101 exit_code1102 });1103 }1104 SRV_INF("instance name=%s exited with status %d\n", name.c_str(), exit_code);1105 });1106 1107 // clean up old process/thread if exists1108 {1109 auto & old_instance = mapping[name];1110 // old process should have exited already, but just in case, we clean it up here1111 if (old_instance.subproc && old_instance.subproc->is_alive()) {1112 SRV_WRN("old process for model name=%s is still alive, this is unexpected\n", name.c_str());1113 old_instance.subproc->terminate(); // force kill1114 }1115 if (old_instance.th.joinable()) {1116 old_instance.th.join();1117 }1118 }1119 1120 notify_sse("model_status", name, {1121 {"status", server_model_status_to_string(inst.meta.status)},1122 });1123 1124 mapping[name] = std::move(inst);1125 cv.notify_all();1126}1127 1128void server_models::unload(const std::string & name) {1129 std::unique_lock<std::mutex> lk(mutex);1130 auto it = mapping.find(name);1131 if (it != mapping.end()) {1132 if (it->second.meta.status == SERVER_MODEL_STATUS_DOWNLOADING) {1133 SRV_INF("cancelling download for model name=%s\n", name.c_str());1134 it->second.subproc->request_exit();1135 // for convenience, we wait the status change here1136 wait(lk, name, [](const server_model_meta & new_meta) {1137 return new_meta.status != SERVER_MODEL_STATUS_DOWNLOADING;1138 });1139 } else if (it->second.meta.is_running()) {1140 SRV_INF("stopping model instance name=%s\n", name.c_str());1141 stopping_models.insert(name);1142 if (it->second.meta.status == SERVER_MODEL_STATUS_LOADING) {1143 // special case: if model is in loading state, unloading means force-killing it1144 SRV_WRN("model name=%s is still loading, force-killing\n", name.c_str());1145 it->second.subproc->terminate();1146 }1147 cv_stop.notify_all();1148 // status change will be handled by the managing thread1149 } else {1150 SRV_WRN("model instance name=%s is not running\n", name.c_str());1151 }1152 }1153}1154 1155void server_models::unload_all() {1156 std::vector<std::thread> to_join;1157 {1158 std::lock_guard<std::mutex> lk(mutex);1159 for (auto & [name, inst] : mapping) {1160 if (inst.meta.status == SERVER_MODEL_STATUS_DOWNLOADING) {1161 SRV_INF("cancelling download for model name=%s\n", name.c_str());1162 inst.subproc->stopped.store(true, std::memory_order_relaxed);1163 } else if (inst.meta.is_running()) {1164 SRV_INF("stopping model instance name=%s\n", name.c_str());1165 stopping_models.insert(name);1166 cv_stop.notify_all();1167 // status change will be handled by the managing thread1168 }1169 // moving the thread to join list to avoid deadlock1170 to_join.push_back(std::move(inst.th));1171 }1172 }1173 for (auto & th : to_join) {1174 if (th.joinable()) {1175 th.join();1176 }1177 }1178}1179 1180void server_models::update_status(const std::string & name, const update_status_args & args) {1181 std::unique_lock<std::mutex> lk(mutex);1182 auto it = mapping.find(name);1183 if (it != mapping.end()) {1184 auto & meta = it->second.meta;1185 meta.status = args.status;1186 meta.exit_code = args.exit_code;1187 if (!args.loaded_info.is_null()) {1188 meta.loaded_info = args.loaded_info;1189 }1190 if (!args.progress.is_null()) {1191 meta.progress = args.progress;1192 }1193 }1194 // broadcast status change to SSE1195 {1196 json data = {1197 {"status", server_model_status_to_string(args.status)},1198 };1199 if (args.status == SERVER_MODEL_STATUS_UNLOADED) {1200 data["exit_code"] = args.exit_code;