Team Ai
Datasetpublic

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.

sourceHugging Faceupdated 2mo agoView on Hugging Face
0likes3kdownloads
server-models.cpp2508 linesDownload Raw Back to server
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;

Showing the first 1,200 of 2508 lines. Download the file for the rest.

Brunobkr/llama.cpp_AlgMor24_github · Team Ai