diff --git a/common/arg.cpp b/common/arg.cpp index bdc2e9eb4fc..3bf8a301ba3 100644 --- a/common/arg.cpp +++ b/common/arg.cpp @@ -2301,6 +2301,20 @@ common_params_context common_params_parser_init(common_params & params, llama_ex params.devices = parse_device_list(value); } ).set_env("LLAMA_ARG_DEVICE")); + add_opt(common_arg( + {"--moe-n-slots"}, "N", + "number of expert slots in GPU memory for MoE disk offloading (0 = disabled)", + [](common_params & params, const std::string & value) { + params.moe.n_slots = std::stoi(value); + } + ).set_env("LLAMA_ARG_MOE_N_SLOTS")); + add_opt(common_arg( + {"--moe-n-layers"}, "N", + "number of MoE layers to disk-offload (0 = all layers)", + [](common_params & params, const std::string & value) { + params.moe.n_layers = std::stoi(value); + } + ).set_env("LLAMA_ARG_MOE_N_LAYERS")); add_opt(common_arg( {"--list-devices"}, "print list of available devices and exit", diff --git a/common/common.cpp b/common/common.cpp index 97daf281783..ed5d7a53787 100644 --- a/common/common.cpp +++ b/common/common.cpp @@ -1553,6 +1553,7 @@ struct llama_model_params common_model_params_to_llama(common_params & params) { mparams.progress_callback = params.load_progress_callback; mparams.progress_callback_user_data = params.load_progress_callback_user_data; mparams.no_alloc = params.no_alloc; + mparams.moe = params.moe; return mparams; } @@ -1568,6 +1569,7 @@ struct llama_context_params common_context_params_to_llama(const common_params & cparams.n_threads = params.cpuparams.n_threads; cparams.n_threads_batch = params.cpuparams_batch.n_threads == -1 ? params.cpuparams.n_threads : params.cpuparams_batch.n_threads; + cparams.moe = params.moe; cparams.embeddings = params.embedding; cparams.rope_scaling_type = params.rope_scaling_type; cparams.rope_freq_base = params.rope_freq_base; diff --git a/common/common.h b/common/common.h index 8a0e5eed5ee..7047d7e7ebb 100644 --- a/common/common.h +++ b/common/common.h @@ -444,6 +444,7 @@ struct common_params { // offload params std::vector devices; // devices to use for offloading + llama_moe_params moe = { 0, 0 }; int32_t n_gpu_layers = -1; // number of layers to store in VRAM, -1 is auto, <= -2 is all int32_t main_gpu = 0; // the GPU that is used for scratch and small tensors diff --git a/ggml/include/ggml-metal.h b/ggml/include/ggml-metal.h index 433838f0d6d..5e789c48325 100644 --- a/ggml/include/ggml-metal.h +++ b/ggml/include/ggml-metal.h @@ -56,6 +56,41 @@ GGML_BACKEND_API void ggml_backend_metal_capture_next_compute(ggml_backend_t bac GGML_BACKEND_API ggml_backend_reg_t ggml_backend_metal_reg(void); +typedef struct ggml_backend_metal_event * ggml_backend_metal_event_t; + +GGML_BACKEND_API ggml_backend_metal_event_t ggml_backend_metal_event_new(ggml_backend_t backend); +GGML_BACKEND_API void ggml_backend_metal_event_free(ggml_backend_metal_event_t event); +GGML_BACKEND_API void ggml_backend_metal_event_signal(ggml_backend_metal_event_t event, uint64_t value); +GGML_BACKEND_API void * ggml_backend_metal_event_raw(ggml_backend_metal_event_t event); + +// Shared memory layout (MTLStorageModeShared). +const int MOE_MAX_IDS = 4096; +const size_t MOE_OFF_REQ = 0; // atomic_uint: request seq +const size_t MOE_OFF_N = 8; // int32: id count +const size_t MOE_OFF_SELECTED = 16; // int32[n]: expert ids (GPU writes) +const size_t MOE_OFF_REMAPPED = MOE_OFF_SELECTED + MOE_MAX_IDS * 4; // int32[n]: slot ids (CPU writes) +const size_t MOE_MSG_NBYTES = MOE_OFF_REMAPPED + MOE_MAX_IDS * 4; + +struct ggml_metal_moe_intercept { + int n; + uint32_t seq; + const struct ggml_tensor * msg_tensor; + ggml_backend_metal_event_t event; + bool reuse; +}; + +typedef bool (*ggml_metal_moe_query_fn)(void * user_data, + const struct ggml_tensor * src0, + const struct ggml_tensor * ids, + struct ggml_metal_moe_intercept * out); + +struct ggml_metal_moe_handler { + ggml_metal_moe_query_fn fn; + void * user_data; +}; + +GGML_BACKEND_API void ggml_backend_metal_set_moe_handler(ggml_backend_t backend, struct ggml_metal_moe_handler handler); + #ifdef __cplusplus } #endif diff --git a/ggml/src/ggml-metal/ggml-metal-context.h b/ggml/src/ggml-metal/ggml-metal-context.h index abf4b06ed2a..69984fadfe6 100644 --- a/ggml/src/ggml-metal/ggml-metal-context.h +++ b/ggml/src/ggml-metal/ggml-metal-context.h @@ -1,6 +1,7 @@ #pragma once #include "ggml-metal-device.h" +#include "ggml-metal.h" #ifdef __cplusplus extern "C" { @@ -32,6 +33,7 @@ void ggml_metal_event_wait (ggml_metal_t ctx, ggml_metal_event_t ev); ggml_metal_event_t ggml_metal_get_ev_cpy(ggml_metal_t ctx); void ggml_metal_set_n_cb (ggml_metal_t ctx, int n_cb); +void ggml_metal_set_moe_handler (ggml_metal_t ctx, struct ggml_metal_moe_handler moe_handler); void ggml_metal_set_abort_callback (ggml_metal_t ctx, ggml_abort_callback abort_callback, void * user_data); bool ggml_metal_supports_family (ggml_metal_t ctx, int family); void ggml_metal_capture_next_compute(ggml_metal_t ctx); diff --git a/ggml/src/ggml-metal/ggml-metal-context.m b/ggml/src/ggml-metal/ggml-metal-context.m index 32d97cd5d0a..4c5158f301e 100644 --- a/ggml/src/ggml-metal/ggml-metal-context.m +++ b/ggml/src/ggml-metal/ggml-metal-context.m @@ -79,6 +79,8 @@ // error state - set when a command buffer fails during synchronize // once set, graph_compute will return GGML_STATUS_FAILED until the backend is recreated bool has_error; + + struct ggml_metal_moe_handler moe_handler; }; ggml_metal_t ggml_metal_init(ggml_metal_device_t dev) { @@ -660,6 +662,11 @@ ggml_metal_event_t ggml_metal_get_ev_cpy(ggml_metal_t ctx) { return ctx->ev_cpy; } +void ggml_metal_set_moe_handler(ggml_metal_t ctx, struct ggml_metal_moe_handler moe_handler) { + GGML_ASSERT(moe_handler.fn == NULL || ctx->moe_handler.fn == NULL); + ctx->moe_handler = moe_handler; +} + void ggml_metal_set_n_cb(ggml_metal_t ctx, int n_cb) { if (ctx->n_cb != n_cb) { ctx->n_cb = MIN(n_cb, GGML_METAL_MAX_COMMAND_BUFFERS); @@ -702,7 +709,8 @@ void ggml_metal_set_n_cb(ggml_metal_t ctx, int n_cb) { ctx->use_concurrency, ctx->capture_compute, ctx->debug_graph, - ctx->debug_fusion); + ctx->debug_fusion, + ctx->moe_handler); for (int idx = 0; idx < ggml_metal_op_n_nodes(ctx_op); ++idx) { const int res = ggml_metal_op_encode(ctx_op, idx); diff --git a/ggml/src/ggml-metal/ggml-metal-device.cpp b/ggml/src/ggml-metal/ggml-metal-device.cpp index ba006d9b31a..95f725b18b1 100644 --- a/ggml/src/ggml-metal/ggml-metal-device.cpp +++ b/ggml/src/ggml-metal/ggml-metal-device.cpp @@ -176,6 +176,24 @@ ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_set_rows(ggml_me return res; } +ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_moe_interceptor(ggml_metal_library_t lib) { + char base[256]; + char name[256]; + + snprintf(base, 256, "kernel_moe_interceptor"); + snprintf(name, 256, "%s", base); + + ggml_metal_pipeline_with_params res = ggml_metal_library_get_pipeline(lib, name); + if (!res.pipeline) { + res = ggml_metal_library_compile_pipeline(lib, base, name, nullptr); + } + + res.nsg = 1; + res.smem = 0; + + return res; +} + ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_diag(ggml_metal_library_t lib, const ggml_tensor * op) { char base[256]; char name[256]; diff --git a/ggml/src/ggml-metal/ggml-metal-device.h b/ggml/src/ggml-metal/ggml-metal-device.h index 4a3ebb5569d..3491b0ee922 100644 --- a/ggml/src/ggml-metal/ggml-metal-device.h +++ b/ggml/src/ggml-metal/ggml-metal-device.h @@ -91,6 +91,8 @@ void ggml_metal_encoder_memory_barrier(ggml_metal_encoder_t encoder); void ggml_metal_encoder_end_encoding(ggml_metal_encoder_t encoder); +void ggml_metal_encoder_wait_for_event(ggml_metal_encoder_t enc, void * event, uint64_t value); + // // MTLLibrary wrapper // @@ -161,6 +163,7 @@ struct ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_opt_step_ struct ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_opt_step_sgd (ggml_metal_library_t lib, const struct ggml_tensor * op); struct ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_memset (ggml_metal_library_t lib, const struct ggml_tensor * op); struct ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_count_equal (ggml_metal_library_t lib, const struct ggml_tensor * op); +struct ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_moe_interceptor (ggml_metal_library_t lib); struct ggml_metal_pipeline_with_params ggml_metal_library_get_pipeline_flash_attn_ext_pad( ggml_metal_library_t lib, @@ -267,6 +270,8 @@ typedef struct ggml_metal_event * ggml_metal_event_t; void ggml_metal_event_encode_signal(ggml_metal_event_t ev, ggml_metal_cmd_buf_t cmd_buf); void ggml_metal_event_encode_wait (ggml_metal_event_t ev, ggml_metal_cmd_buf_t cmd_buf); +void ggml_metal_event_cpu_signal (ggml_metal_event_t ev, uint64_t value); +void * ggml_metal_event_get_obj (ggml_metal_event_t ev); ggml_metal_device_t ggml_metal_device_init(int device); void ggml_metal_device_free(ggml_metal_device_t dev); diff --git a/ggml/src/ggml-metal/ggml-metal-device.m b/ggml/src/ggml-metal/ggml-metal-device.m index 885344ec670..0a2a4c68f43 100644 --- a/ggml/src/ggml-metal/ggml-metal-device.m +++ b/ggml/src/ggml-metal/ggml-metal-device.m @@ -458,6 +458,8 @@ struct ggml_metal_pipeline_with_params ggml_metal_library_compile_pipeline(ggml_ struct ggml_metal_encoder { id obj; + id cmd_buf; + bool concurrent; }; ggml_metal_encoder_t ggml_metal_encoder_init(ggml_metal_cmd_buf_t cmd_buf_raw, bool concurrent) { @@ -465,6 +467,10 @@ ggml_metal_encoder_t ggml_metal_encoder_init(ggml_metal_cmd_buf_t cmd_buf_raw, b id cmd_buf = (id) cmd_buf_raw; + res->cmd_buf = cmd_buf; + res->concurrent = concurrent; + [res->cmd_buf retain]; + if (concurrent) { res->obj = [cmd_buf computeCommandEncoderWithDispatchType: MTLDispatchTypeConcurrent]; } else { @@ -476,8 +482,23 @@ ggml_metal_encoder_t ggml_metal_encoder_init(ggml_metal_cmd_buf_t cmd_buf_raw, b return res; } +void ggml_metal_encoder_wait_for_event(ggml_metal_encoder_t enc, void *event, uint64_t value) { + id ev = (__bridge id)event; + [enc->obj endEncoding]; + [enc->obj release]; + [enc->cmd_buf encodeWaitForEvent:ev value:value]; + if (enc->concurrent) { + enc->obj = [enc->cmd_buf + computeCommandEncoderWithDispatchType:MTLDispatchTypeConcurrent]; + } else { + enc->obj = [enc->cmd_buf computeCommandEncoder]; + } + [enc->obj retain]; +} + void ggml_metal_encoder_free(ggml_metal_encoder_t encoder) { [encoder->obj release]; + [encoder->cmd_buf release]; free(encoder); } @@ -1004,6 +1025,14 @@ void ggml_metal_event_encode_wait(ggml_metal_event_t ev, ggml_metal_cmd_buf_t cm [cmd_buf encodeWaitForEvent:event value:atomic_load_explicit(&ev->value, memory_order_relaxed)]; } +void ggml_metal_event_cpu_signal(ggml_metal_event_t ev, uint64_t value) { + ((id)ev->obj).signaledValue = value; +} + +void *ggml_metal_event_get_obj(ggml_metal_event_t ev) { + return ev->obj; +} + ggml_metal_event_t ggml_metal_device_event_init(ggml_metal_device_t dev) { id event = [dev->mtl_device newSharedEvent]; diff --git a/ggml/src/ggml-metal/ggml-metal-ops.cpp b/ggml/src/ggml-metal/ggml-metal-ops.cpp index 206af227a2c..e2e758975da 100644 --- a/ggml/src/ggml-metal/ggml-metal-ops.cpp +++ b/ggml/src/ggml-metal/ggml-metal-ops.cpp @@ -1,5 +1,6 @@ #include "ggml-metal-ops.h" +#include "ggml-metal.h" #include "ggml.h" #include "ggml-impl.h" #include "ggml-backend-impl.h" @@ -36,7 +37,8 @@ struct ggml_metal_op { bool use_concurrency, bool use_capture, int debug_graph, - int debug_fusion) { + int debug_fusion, + ggml_metal_moe_handler moe_handler) { this->dev = dev; this->lib = ggml_metal_device_get_library(dev); this->enc = ggml_metal_encoder_init(cmd_buf, use_concurrency); @@ -48,6 +50,7 @@ struct ggml_metal_op { this->use_capture = use_capture; this->debug_graph = debug_graph; this->debug_fusion = debug_fusion; + this->moe_handler = moe_handler; this->gf = gf; idxs.reserve(gf->n_nodes); @@ -100,6 +103,8 @@ struct ggml_metal_op { int debug_graph; int debug_fusion; + ggml_metal_moe_handler moe_handler; + private: ggml_cgraph * gf; @@ -120,7 +125,8 @@ ggml_metal_op_t ggml_metal_op_init( bool use_concurrency, bool use_capture, int debug_graph, - int debug_fusion) { + int debug_fusion, + ggml_metal_moe_handler moe_handler) { ggml_metal_op_t res = new ggml_metal_op( dev, cmd_buf, @@ -131,7 +137,8 @@ ggml_metal_op_t ggml_metal_op_init( use_concurrency, use_capture, debug_graph, - debug_fusion); + debug_fusion, + moe_handler); return res; } @@ -2310,6 +2317,35 @@ int ggml_metal_op_mul_mat_id(ggml_metal_op_t ctx, int idx) { ggml_metal_buffer_id bid_src2 = ggml_metal_get_buffer_id(op->src[2]); ggml_metal_buffer_id bid_dst = ggml_metal_get_buffer_id(op); + // moe interceptor + if (ctx->moe_handler.fn) { + ggml_metal_moe_intercept mi; + if (ctx->moe_handler.fn(ctx->moe_handler.user_data, op->src[0], op->src[2], &mi)) { + ggml_metal_buffer_id moe_base = ggml_metal_get_buffer_id(mi.msg_tensor); + + if (!mi.reuse) { + auto pipeline = ggml_metal_library_get_pipeline_moe_interceptor(lib); + + ggml_metal_buffer_id b_req = { moe_base.metal, moe_base.offs + MOE_OFF_REQ }; + ggml_metal_buffer_id b_msel = { moe_base.metal, moe_base.offs + MOE_OFF_SELECTED }; + uint32_t n_u = (uint32_t) mi.n; + + ggml_metal_encoder_set_pipeline(enc, pipeline); + ggml_metal_encoder_set_buffer(enc, bid_src2, 0); + ggml_metal_encoder_set_buffer(enc, b_req, 1); + ggml_metal_encoder_set_buffer(enc, b_msel, 2); + ggml_metal_encoder_set_bytes(enc, &n_u, sizeof(n_u), 3); + ggml_metal_encoder_set_bytes(enc, &mi.seq, sizeof(mi.seq), 4); + ggml_metal_encoder_dispatch_threadgroups(enc, 1, 1, 1, 32, 1, 1); + + ggml_metal_encoder_wait_for_event(enc, ggml_backend_metal_event_raw(mi.event), mi.seq); + } + + bid_src2.metal = moe_base.metal; + bid_src2.offs = moe_base.offs + MOE_OFF_REMAPPED; + } + } + const uint32_t r2 = 1; const uint32_t r3 = 1; diff --git a/ggml/src/ggml-metal/ggml-metal-ops.h b/ggml/src/ggml-metal/ggml-metal-ops.h index 36c61071b4f..bb10d01f04e 100644 --- a/ggml/src/ggml-metal/ggml-metal-ops.h +++ b/ggml/src/ggml-metal/ggml-metal-ops.h @@ -1,6 +1,7 @@ #pragma once #include "ggml-metal-device.h" +#include "ggml-metal.h" #ifdef __cplusplus extern "C" { @@ -18,7 +19,8 @@ ggml_metal_op_t ggml_metal_op_init( bool use_concurrency, bool use_capture, int debug_graph, - int debug_fusion); + int debug_fusion, + struct ggml_metal_moe_handler); void ggml_metal_op_free(ggml_metal_op_t ctx); diff --git a/ggml/src/ggml-metal/ggml-metal.cpp b/ggml/src/ggml-metal/ggml-metal.cpp index a1003b3acff..96b66ff1135 100644 --- a/ggml/src/ggml-metal/ggml-metal.cpp +++ b/ggml/src/ggml-metal/ggml-metal.cpp @@ -565,6 +565,42 @@ static void ggml_backend_metal_set_n_cb(ggml_backend_t backend, int n_cb) { ggml_metal_set_n_cb(ctx, n_cb); } +void ggml_backend_metal_set_moe_handler(ggml_backend_t backend, struct ggml_metal_moe_handler moe_handler) { + GGML_ASSERT(ggml_backend_is_metal(backend)); + + ggml_metal_t ctx = (ggml_metal_t) backend->context; + + ggml_metal_set_moe_handler(ctx, moe_handler); +} + +struct ggml_backend_metal_event { + ggml_metal_device_t dev; + ggml_metal_event_t ev; +}; + +ggml_backend_metal_event_t ggml_backend_metal_event_new(ggml_backend_t backend) { + GGML_ASSERT(ggml_backend_is_metal(backend)); + ggml_metal_device_t dev = (ggml_metal_device_t) backend->device->context; + + auto * e = new ggml_backend_metal_event; + e->dev = dev; + e->ev = ggml_metal_device_event_init(dev); + return e; +} + +void ggml_backend_metal_event_free(ggml_backend_metal_event_t event) { + ggml_metal_device_event_free(event->dev, event->ev); + delete event; +} + +void ggml_backend_metal_event_signal(ggml_backend_metal_event_t event, uint64_t value) { + ggml_metal_event_cpu_signal(event->ev, value); +} + +void * ggml_backend_metal_event_raw(ggml_backend_metal_event_t event) { + return ggml_metal_event_get_obj(event->ev); +} + static ggml_backend_i ggml_backend_metal_i = { /* .get_name = */ ggml_backend_metal_name, /* .free = */ ggml_backend_metal_free, diff --git a/ggml/src/ggml-metal/ggml-metal.metal b/ggml/src/ggml-metal/ggml-metal.metal index e772664ba91..f751d3b17c2 100644 --- a/ggml/src/ggml-metal/ggml-metal.metal +++ b/ggml/src/ggml-metal/ggml-metal.metal @@ -1010,6 +1010,22 @@ template<> inline float4 elu_approx(float4 x) { return res; } +kernel void kernel_moe_interceptor( + device const int * selected [[buffer(0)]], + device atomic_uint * req_seq [[buffer(1)]], + device int * msg_sel [[buffer(2)]], + constant uint & n [[buffer(3)]], + constant uint & seq [[buffer(4)]], + uint tid [[thread_position_in_threadgroup]]) { + for (uint i = tid; i < n; i += 32) { + msg_sel[i] = selected[i]; + } + threadgroup_barrier(mem_flags::mem_device); + if (tid == 0) { + atomic_store_explicit(req_seq, seq, memory_order_relaxed); + } +} + constant short FC_unary_op [[function_constant(FC_UNARY + 0)]]; constant bool FC_unary_cnt[[function_constant(FC_UNARY + 1)]]; diff --git a/include/llama.h b/include/llama.h index e8374c53b70..e7ee359c599 100644 --- a/include/llama.h +++ b/include/llama.h @@ -288,6 +288,11 @@ extern "C" { ggml_backend_buffer_type_t buft; }; + typedef struct llama_moe_params { + int32_t n_slots; + int32_t n_layers; + } llama_moe_params; + struct llama_model_params { // NULL-terminated list of devices to use for offloading (if NULL, all available devices are used) ggml_backend_dev_t * devices; @@ -324,6 +329,8 @@ extern "C" { bool use_extra_bufts; // use extra buffer types (used for weight repacking) bool no_host; // bypass host buffer allowing extra buffers to be used bool no_alloc; // only load metadata and simulate memory allocations + + llama_moe_params moe; }; struct llama_sampler_seq_config { @@ -342,6 +349,8 @@ extern "C" { int32_t n_threads; // number of threads to use for generation int32_t n_threads_batch; // number of threads to use for batch processing + llama_moe_params moe; + enum llama_context_type ctx_type; // set the context type (e.g. MTP) enum llama_rope_scaling_type rope_scaling_type; // RoPE scaling type, from `enum llama_rope_scaling_type` enum llama_pooling_type pooling_type; // whether to pool (sum) embedding results by sequence id diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 7b1fcfca0ad..9e13f99601e 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -32,6 +32,7 @@ add_library(llama llama-model-loader.cpp llama-model-saver.cpp llama-model.cpp + llama-moe-offloader.cpp llama-quant.cpp llama-sampler.cpp llama-vocab.cpp diff --git a/src/llama-context.cpp b/src/llama-context.cpp index ad36c06667d..58501eeed31 100644 --- a/src/llama-context.cpp +++ b/src/llama-context.cpp @@ -60,6 +60,10 @@ llama_context::llama_context( cparams.n_threads = params.n_threads; cparams.n_threads_batch = params.n_threads_batch; + cparams.moe = params.moe; + if (cparams.moe.n_layers == 0) { + cparams.moe.n_layers = INT32_MAX; + } cparams.yarn_ext_factor = params.yarn_ext_factor >= 0.0f ? params.yarn_ext_factor : hparams.yarn_ext_factor; cparams.yarn_attn_factor = params.yarn_attn_factor >= 0.0f ? params.yarn_attn_factor : hparams.yarn_attn_factor; cparams.yarn_beta_fast = params.yarn_beta_fast >= 0.0f ? params.yarn_beta_fast : hparams.yarn_beta_fast; @@ -248,6 +252,24 @@ llama_context::llama_context( backends.emplace_back(backend); } + if (cparams.moe.n_slots > 0) { +#if defined(__APPLE__) && defined(GGML_USE_METAL) + ggml_backend_t metal_be = nullptr; + for (auto & backend : backends) { + if (ggml_backend_is_metal(backend.get())) { + metal_be = backend.get(); + break; + } + } + if (metal_be) { + moe_offloader = std::make_unique(metal_be); + moe_offloader->start(); + } +#else + LLAMA_LOG_WARN("%s: moe offloader not supported on this platform, ignoring --moe-n-slots\n", __func__); +#endif + } + // add ACCEL backends (such as BLAS) for (size_t i = 0; i < ggml_backend_dev_count(); ++i) { ggml_backend_dev_t dev = ggml_backend_dev_get(i); @@ -389,6 +411,9 @@ llama_context::llama_context( } llama_context::~llama_context() { + if (moe_offloader) { + moe_offloader->stop(); + } if (!model.hparams.no_alloc) { for (size_t i = 0; i < backend_ptrs.size(); ++i) { ggml_backend_t backend = backend_ptrs[i]; @@ -2286,21 +2311,25 @@ llm_graph_params llama_context::graph_params( const llama_memory_context_i * mctx, llm_graph_type gtype) const { return { - /*.arch =*/ model.arch, - /*.hparams =*/ model.hparams, - /*.cparams =*/ cparams, - /*.ubatch =*/ ubatch, - /*.gtype =*/ gtype, - /*.sched =*/ sched.get(), - /*.backend_cpu =*/ backend_cpu, - /*.cvec =*/ cvec.get(), - /*.loras =*/ loras.get(), - /*.mctx =*/ mctx, - /*.cross =*/ &cross, - /*.samplers =*/ sampling.samplers, - /*.n_outputs =*/ n_outputs, - /*.cb =*/ graph_get_cb(), - /*.res =*/ res, + /*.arch =*/ model.arch, + /*.hparams =*/ model.hparams, + /*.cparams =*/ cparams, + /*.ubatch =*/ ubatch, + /*.gtype =*/ gtype, + /*.sched =*/ sched.get(), + /*.backend_cpu =*/ backend_cpu, + /*.cvec =*/ cvec.get(), + /*.loras =*/ loras.get(), + /*.mctx =*/ mctx, + /*.cross =*/ &cross, + /*.samplers =*/ sampling.samplers, + /*.moe_n_slots =*/ cparams.moe.n_slots, + /*.moe_n_layers =*/ cparams.moe.n_layers, + /*.moe_file_idx =*/ &model.moe_file_idx, + /*.moe_offloader =*/ moe_offloader.get(), + /*.n_outputs =*/ n_outputs, + /*.cb =*/ graph_get_cb(), + /*.res =*/ res, }; } @@ -3339,6 +3368,7 @@ llama_context_params llama_context_default_params() { /*.n_rs_seq =*/ 0, /*.n_threads =*/ GGML_DEFAULT_N_THREADS, // TODO: better default /*.n_threads_batch =*/ GGML_DEFAULT_N_THREADS, + /*.moe =*/ {0, 0}, /*.ctx_type =*/ LLAMA_CONTEXT_TYPE_DEFAULT, /*.rope_scaling_type =*/ LLAMA_ROPE_SCALING_TYPE_UNSPECIFIED, /*.pooling_type =*/ LLAMA_POOLING_TYPE_UNSPECIFIED, diff --git a/src/llama-context.h b/src/llama-context.h index d03f681d4a1..0853769154a 100644 --- a/src/llama-context.h +++ b/src/llama-context.h @@ -1,5 +1,6 @@ #pragma once +#include "llama-moe-offloader.h" #include "llama.h" #include "llama-ext.h" #include "llama-cparams.h" @@ -361,6 +362,8 @@ struct llama_context { // env: LLAMA_GRAPH_REUSE_DISABLE bool graph_reuse_disable = false; + std::unique_ptr moe_offloader; + // perf mutable int64_t t_start_us = 0; mutable int64_t t_load_us = 0; diff --git a/src/llama-cparams.h b/src/llama-cparams.h index 20ec59fe335..65c310d107a 100644 --- a/src/llama-cparams.h +++ b/src/llama-cparams.h @@ -16,6 +16,8 @@ struct llama_cparams { int32_t n_threads; // number of threads to use for generation int32_t n_threads_batch; // number of threads to use for batch processing + llama_moe_params moe; + float rope_freq_base; float rope_freq_scale; diff --git a/src/llama-graph.cpp b/src/llama-graph.cpp index fc027de8b39..9cd8d91898a 100644 --- a/src/llama-graph.cpp +++ b/src/llama-graph.cpp @@ -5,6 +5,8 @@ #include "llama-batch.h" #include "llama-cparams.h" +#include "llama-moe-offloader.h" + #include "llama-kv-cache.h" #include "llama-kv-cache-iswa.h" #include "llama-memory-hybrid.h" @@ -968,6 +970,10 @@ llm_graph_context::llm_graph_context(const llm_graph_params & params) : mctx (params.mctx), cross (params.cross), samplers (params.samplers), + moe_n_slots (params.moe_n_slots), + moe_n_layers (params.moe_n_layers), + moe_file_idx (params.moe_file_idx), + moe_offloader (params.moe_offloader), cb_func (params.cb), res (params.res), ctx0 (res->get_ctx()), @@ -1530,9 +1536,48 @@ ggml_tensor * llm_graph_context::build_moe_ffn( ggml_tensor * up = nullptr; ggml_tensor * experts = nullptr; + const bool moe_layer_enabled = moe_offloader != nullptr && (il < moe_n_layers) && (n_expert_used < n_expert); + int64_t n_slots = moe_n_slots > 0 ? (int64_t) moe_n_slots : n_expert; + + ggml_tensor * routing_ids = selected_experts; + + ggml_tensor * gate_up_exps_src = gate_up_exps; + ggml_tensor * up_exps_src = up_exps; + ggml_tensor * gate_exps_src = gate_exps; + ggml_tensor * down_exps_src = down_exps; + + if (moe_layer_enabled) { + // find weight by name and create a slot for it + auto bind_pool = [&](ggml_tensor * orig) -> ggml_tensor * { + if (!orig) { + return nullptr; + } + auto it = moe_file_idx->find(orig->name); + if (it == moe_file_idx->end()) { + GGML_ABORT("moe: tensor not found in file index: %s", orig->name); + } + const auto & fi = it->second; + return moe_offloader->bind_pool(il, orig, n_slots, fi.offset, fi.fd); + }; + + // must share the same expert->slot mapping + // see resolve() which picks victims in lockstep across all pools. + // nullptr tensors (merged gate_up path) are skipped by bind_pool. + gate_up_exps_src = bind_pool(gate_up_exps); + up_exps_src = bind_pool(up_exps); + gate_exps_src = bind_pool(gate_exps); + down_exps_src = bind_pool(down_exps); + + // selected experts is a non-contiguous view, + // but the interceptor expects flat i32 array + routing_ids = ggml_cont(ctx0, selected_experts); + ggml_set_name(routing_ids, "moe_slot_ids"); + cb(routing_ids, "ffn_moe_slot_ids", il); + } + if (gate_up_exps) { // merged gate_up path: one mul_mat_id, then split into gate and up views - ggml_tensor * gate_up = build_lora_mm_id(gate_up_exps, cur, selected_experts); // [n_ff*2, n_expert_used, n_tokens] + ggml_tensor * gate_up = build_lora_mm_id(gate_up_exps_src, cur, routing_ids); // [n_ff*2, n_expert_used, n_tokens] cb(gate_up, "ffn_moe_gate_up", il); if (gate_up_exps_b) { @@ -1556,7 +1601,7 @@ ggml_tensor * llm_graph_context::build_moe_ffn( cb(up, "ffn_moe_up", il); } else { // separate gate and up path - up = build_lora_mm_id(up_exps, cur, selected_experts); // [n_ff, n_expert_used, n_tokens] + up = build_lora_mm_id(up_exps_src, cur, routing_ids); // [n_ff, n_expert_used, n_tokens] cb(up, "ffn_moe_up", il); if (up_exps_b) { @@ -1574,7 +1619,7 @@ ggml_tensor * llm_graph_context::build_moe_ffn( } if (gate_exps) { - cur = build_lora_mm_id(gate_exps, cur, selected_experts); // [n_ff, n_expert_used, n_tokens] + cur = build_lora_mm_id(gate_exps_src, cur, routing_ids); // [n_ff, n_expert_used, n_tokens] cb(cur, "ffn_moe_gate", il); } else { cur = up; @@ -1664,7 +1709,7 @@ ggml_tensor * llm_graph_context::build_moe_ffn( GGML_ABORT("fatal error"); } - experts = build_lora_mm_id(down_exps, cur, selected_experts); // [n_embd, n_expert_used, n_tokens] + experts = build_lora_mm_id(down_exps_src, cur, routing_ids); // [n_embd, n_expert_used, n_tokens] cb(experts, "ffn_moe_down", il); if (down_exps_b) { diff --git a/src/llama-graph.h b/src/llama-graph.h index bf6778237e6..806e8f19f18 100644 --- a/src/llama-graph.h +++ b/src/llama-graph.h @@ -4,6 +4,8 @@ #include "llama-batch.h" #include "llama-hparams.h" #include "llama-adapter.h" +#include "llama-moe-offloader.h" +#include "llama-model-loader.h" #include #include @@ -548,6 +550,12 @@ struct llm_graph_params { std::map samplers; + int32_t moe_n_slots = 0; + int32_t moe_n_layers = INT32_MAX; + + const std::unordered_map *moe_file_idx; + llama_moe_offloader *moe_offloader; + static bool samplers_equal( const std::map & lhs, const std::map & rhs) { @@ -765,6 +773,12 @@ struct llm_graph_context { std::map samplers; + int32_t moe_n_slots = 0; + int32_t moe_n_layers = INT32_MAX; + + const std::unordered_map *moe_file_idx; + llama_moe_offloader * moe_offloader; + const llm_graph_cb & cb_func; llm_graph_result * res; diff --git a/src/llama-model-loader.cpp b/src/llama-model-loader.cpp index c645d0785ab..66a149da381 100644 --- a/src/llama-model-loader.cpp +++ b/src/llama-model-loader.cpp @@ -1249,6 +1249,25 @@ struct ggml_tensor * llama_model_loader::create_tensor( } ggml_tensor * t_meta = get_tensor_meta(tn.str().c_str()); + + if (t_meta && skip_tensor_set.count(t_meta->name)) { + GGML_ASSERT(ctx_skip != nullptr); + + ggml_tensor * t = ggml_get_tensor(ctx_skip, t_meta->name); + if (t) { + return t; + } + + t = ggml_dup_tensor(ctx_skip, t_meta); + ggml_set_name(t, t_meta->name); + + if (!(flags & TENSOR_DUPLICATED)) { + n_created++; + } + + return t; + } + ggml_backend_buffer_type_t buft = buft_for_tensor(t_meta); if (buft == nullptr) { return nullptr; // return type is ggml_tensor * diff --git a/src/llama-model-loader.h b/src/llama-model-loader.h index c476026d3e5..505f9feb0d4 100644 --- a/src/llama-model-loader.h +++ b/src/llama-model-loader.h @@ -14,6 +14,7 @@ #include #include #include +#include using llama_buf_map = std::unordered_map; @@ -28,6 +29,11 @@ enum llama_fver { const char * llama_file_version_name(llama_fver version); +struct llm_tensor_file_info { + int fd; + uint64_t offset; +}; + struct llama_model_loader { // Holds information on a model weight struct llama_tensor_weight { @@ -86,6 +92,9 @@ struct llama_model_loader { llama_mmaps mappings; + std::unordered_set skip_tensor_set; + ggml_context * ctx_skip = nullptr; + std::map weights_map; std::unordered_map kv_overrides; const llama_model_tensor_buft_override * tensor_buft_overrides; diff --git a/src/llama-model.cpp b/src/llama-model.cpp index 0c3e03a61dc..44c0b72da89 100644 --- a/src/llama-model.cpp +++ b/src/llama-model.cpp @@ -34,6 +34,12 @@ #include #include +#ifdef __has_include +# if __has_include() +# include +# endif +#endif + static llama_model * llama_model_mapping(llm_arch arch, const llama_model_params & params) { switch (arch) { case LLM_ARCH_LLAMA: @@ -967,6 +973,16 @@ llama_model::~llama_model() { for (auto * lora : loras) { delete lora; } + +#if defined(__APPLE__) + for (int fd : moe_duped_fds) { + close(fd); + } + + if (moe_ctx_skip) { + ggml_free(moe_ctx_skip); + } +#endif } void llama_model_base::load_stats(llama_model_loader & ml) { @@ -1264,6 +1280,38 @@ bool llama_model_base::load_tensors(llama_model_loader & ml) { layers.resize(n_layer); + for (const auto & [name, w] : ml.weights_map) { + GGML_UNUSED(w); + + // TODO: come up with a better way, same for layer_idx parsing + if (name.find("_exps.") == std::string::npos) { + continue; + } + + if (params.moe.n_slots <= 0) { + continue; + } + + int layer_idx = -1; + sscanf(name.c_str(), "blk.%d.", &layer_idx); + if (layer_idx < 0 || layer_idx >= params.moe.n_layers) { + continue; + } + + ml.skip_tensor_set.insert(name); + } + + if (!ml.skip_tensor_set.empty()) { + const size_t n_skip = ml.skip_tensor_set.size(); + + ggml_init_params p{}; + p.mem_size = ggml_tensor_overhead() * (n_skip + 32); + p.mem_buffer = nullptr; + p.no_alloc = true; + + ml.ctx_skip = ggml_init(p); + } + // call the per-model loading function load_arch_tensors(ml); @@ -1422,6 +1470,27 @@ bool llama_model_base::load_tensors(llama_model_loader & ml) { } ml.done_getting_tensors(); +#if defined(__APPLE__) + for (const auto & file : ml.files) { + const int fd = dup(file->file_id()); + if (fd < 0) { + throw std::runtime_error(format("failed to dup fd for model file: %s", strerror(errno))); + } + moe_duped_fds.push_back(fd); + } + + for (const auto & [name, w] : ml.weights_map) { + if (name.find("_exps.") == std::string::npos) { + continue; + } + moe_file_idx[name] = { moe_duped_fds.at(w.idx), w.offs }; + } +#endif + + // pass ownership to model + moe_ctx_skip = ml.ctx_skip; + ml.ctx_skip = nullptr; + GGML_ASSERT(!(output && tok_embd && strcmp(output->name, tok_embd->name) == 0 && output->type == GGML_TYPE_NVFP4)); @@ -1432,6 +1501,12 @@ bool llama_model_base::load_tensors(llama_model_loader & ml) { } } + if (moe_ctx_skip) { + for (auto * cur = ggml_get_first_tensor(moe_ctx_skip); cur != NULL; cur = ggml_get_next_tensor(moe_ctx_skip, cur)) { + tensors_by_name.emplace_back(ggml_get_name(cur), cur); + } + } + ml.init_mappings(true, use_mlock ? &pimpl->mlock_mmaps : nullptr); pimpl->mappings.reserve(ml.mappings.size()); @@ -2147,6 +2222,7 @@ llama_model_params llama_model_default_params() { /*.use_extra_bufts =*/ true, /*.no_host =*/ false, /*.no_alloc =*/ false, + /*.moe =*/ {0, 0}, }; return result; diff --git a/src/llama-model.h b/src/llama-model.h index b797b8966ac..6b5ce196276 100644 --- a/src/llama-model.h +++ b/src/llama-model.h @@ -6,6 +6,7 @@ #include "llama-hparams.h" #include "llama-memory.h" #include "llama-vocab.h" +#include "llama-model-loader.h" #include #include @@ -586,6 +587,11 @@ struct llama_model { int64_t t_load_us = 0; int64_t t_start_us = 0; + // for moe disk offloader + std::vector moe_duped_fds; + std::unordered_map moe_file_idx; + ggml_context * moe_ctx_skip = nullptr; + explicit llama_model(const llama_model_params & params); virtual ~llama_model(); diff --git a/src/llama-moe-offloader.cpp b/src/llama-moe-offloader.cpp new file mode 100644 index 00000000000..7896a43620b --- /dev/null +++ b/src/llama-moe-offloader.cpp @@ -0,0 +1,353 @@ +#if defined(__APPLE__) && defined(GGML_USE_METAL) + +#include "llama-moe-offloader.h" + +#include "ggml-backend.h" +#include "ggml-metal.h" + +#include +#include +#include + +#include +#include +#include +#include +#include +#include + +moe_layer::~moe_layer() { + for (auto & p : pools) { + ggml_backend_buffer_free(p.own_buf); + ggml_free(p.own_ctx); + } + ggml_backend_buffer_free(msg_buf); + ggml_free(msg_ctx); + if (shared_event) { + ggml_backend_metal_event_free(shared_event); + } +} + +llama_moe_offloader::llama_moe_offloader(ggml_backend_t backend) : backend(backend) {} + +llama_moe_offloader::~llama_moe_offloader() { + clear(); +} + +void llama_moe_offloader::start() { + start_sidecar(); + ggml_backend_metal_set_moe_handler(backend, { hook, this }); +} + +void llama_moe_offloader::stop() { + stop_sidecar(); + ggml_backend_synchronize(backend); + ggml_backend_metal_set_moe_handler(backend, { nullptr, nullptr }); + backend = nullptr; +} + +ggml_tensor * llama_moe_offloader::bind_pool(int layer_idx, + ggml_tensor * orig, + int64_t n_slots, + uint64_t file_offset, + int fd) { + GGML_ASSERT(file_offset != 0); + + if (!orig) { + return nullptr; + } + + std::lock_guard g(build_mtx); + + if ((int) layers.size() <= layer_idx) { + layers.resize((size_t) layer_idx + 1); + } + + auto & layer_ptr = layers[(size_t) layer_idx]; + + if (!layer_ptr) { + auto layer = std::make_unique(); + layer->layer_idx = layer_idx; + layer->n_slots = n_slots; + layer->n_expert = orig->ne[2]; + + // lru + layer->expert_to_slot.assign((size_t) layer->n_expert, -1); + layer->slot_to_expert.assign((size_t) n_slots, -1); + layer->lru_clock.assign((size_t) n_slots, 0); + layer->in_use.assign((size_t) n_slots, 0); + + // shmem message setup + { + ggml_init_params p{}; + p.mem_size = ggml_tensor_overhead() + 64; + p.no_alloc = true; + layer->msg_ctx = ggml_init(p); + layer->msg_tensor = ggml_new_tensor_1d(layer->msg_ctx, GGML_TYPE_I8, MOE_MSG_NBYTES); + char nm[64]; + snprintf(nm, sizeof(nm), "moe_msg_L%d", layer_idx); + ggml_set_name(layer->msg_tensor, nm); + layer->msg_buf = ggml_backend_alloc_ctx_tensors(layer->msg_ctx, backend); + GGML_ASSERT(layer->msg_buf); + layer->mapped = (uint8_t *) layer->msg_tensor->data; + memset(layer->mapped, 0, MOE_MSG_NBYTES); + } + + layer->shared_event = ggml_backend_metal_event_new(backend); + + sidecar_layers.push_back(layer.get()); + layer_ptr = std::move(layer); + } + + moe_layer & layer = *layer_ptr; + + auto it = layer.name_to_pool_idx.find(orig->name); + if (it != layer.name_to_pool_idx.end()) { + return layer.pools[it->second].tensor; + } + + moe_pool pool; + pool.file_offset = file_offset; + pool.fd = fd; + pool.stride = orig->nb[2]; + { + ggml_init_params p{}; + p.mem_size = ggml_tensor_overhead() * 2; + p.no_alloc = true; + pool.own_ctx = ggml_init(p); + } + pool.tensor = ggml_new_tensor_3d(pool.own_ctx, orig->type, orig->ne[0], orig->ne[1], n_slots); + GGML_ASSERT(pool.tensor->nb[2] == pool.stride && "pool/orig per-expert stride mismatch"); + { + char nm[256]; + snprintf(nm, sizeof(nm), "%s_pool", orig->name); + ggml_set_name(pool.tensor, nm); + } + pool.own_buf = ggml_backend_alloc_ctx_tensors(pool.own_ctx, backend); + GGML_ASSERT(pool.own_buf); + + size_t pool_idx = layer.pools.size(); + pool_to_layer[pool.tensor] = &layer; + layer.name_to_pool_idx[orig->name] = pool_idx; + + layer.pools.push_back(pool); + + return layer.pools.back().tensor; +} + +bool llama_moe_offloader::hook(void * user_data, + const ggml_tensor * src0, + const ggml_tensor * src2, + ggml_metal_moe_intercept * out) { + auto * me = (llama_moe_offloader *) user_data; + + auto it = me->pool_to_layer.find(src0); + if (it == me->pool_to_layer.end()) { + return false; + } + + moe_layer & layer = *it->second; + + out->msg_tensor = layer.msg_tensor; + out->event = layer.shared_event; + + int n = (int) (src2->ne[0] * src2->ne[1]); + GGML_ASSERT(n <= MOE_MAX_IDS); + out->n = n; + + const int expected_uses = (int) layer.pools.size(); + + if (layer.last_src2 == src2 && layer.last_src2_uses < expected_uses) { + out->reuse = true; + out->seq = layer.last_src2_seq; + layer.last_src2_uses++; + } else { + out->reuse = false; + uint32_t seq = layer.next_seq.fetch_add(1, std::memory_order_relaxed) + 1; + layer.last_src2 = src2; + layer.last_src2_seq = seq; + layer.last_src2_uses = 1; + out->seq = seq; + *(int32_t *) (layer.mapped + MOE_OFF_N) = n; + } + + return true; +} + +int64_t llama_moe_offloader::lru_evict(moe_layer & layer) { + int64_t victim = -1; + uint32_t oldest = UINT32_MAX; + for (int64_t i = 0; i < layer.n_slots; ++i) { + if (!layer.in_use[i] && layer.lru_clock[i] < oldest) { + oldest = layer.lru_clock[i]; + victim = i; + } + } + return victim; +} + +bool llama_moe_offloader::pread_pool(moe_layer & layer, size_t pool_idx, int32_t expert_id, int64_t dst_slot) { + moe_pool & p = layer.pools[pool_idx]; + uint64_t off = p.file_offset + (uint64_t) expert_id * p.stride; + void * dst = (uint8_t *) p.tensor->data + (size_t) dst_slot * p.stride; + + uint8_t * d = (uint8_t *) dst; + size_t got = 0; + while (got < p.stride) { + ssize_t r = pread(p.fd, d + got, p.stride - got, (off_t) (off + got)); + if (r < 0 && errno == EINTR) { + continue; + } + if (r <= 0) { + return false; + } + got += (size_t) r; + } + return true; +} + +struct moe_pread_task { + moe_layer * layer; + size_t pool_idx; + int32_t expert_id; + int64_t dst_slot; +}; + +void llama_moe_offloader::resolve(moe_layer & layer, const int32_t * ids, int32_t * out, int n) { + std::vector tasks; + tasks.reserve((size_t) n * layer.pools.size()); + + std::fill(layer.in_use.begin(), layer.in_use.end(), 0); + + for (int i = 0; i < n; ++i) { + const int32_t e = ids[i]; + GGML_ASSERT(e >= 0 && e < layer.n_expert); + + int32_t slot = layer.expert_to_slot[(size_t) e]; + + if (slot >= 0) { + layer.lru_clock[slot] = ++layer.lru_time; + layer.total_hits++; + } else { + const int64_t victim = lru_evict(layer); + if (victim < 0) { + fprintf(stderr, "moe L%d: n_slots=%lld too small - ubatch has more unique experts than slots\n", + layer.layer_idx, (long long) layer.n_slots); + GGML_ABORT("moe_offloader: n_slots too small for ubatch"); + } + slot = (int32_t) victim; + + const int32_t old_e = layer.slot_to_expert[(size_t) victim]; + if (old_e >= 0) { + layer.expert_to_slot[(size_t) old_e] = -1; + } + + layer.expert_to_slot[(size_t) e] = slot; + layer.slot_to_expert[(size_t) victim] = e; + layer.lru_clock[victim] = ++layer.lru_time; + layer.total_misses++; + + for (size_t pi = 0; pi < layer.pools.size(); ++pi) { + tasks.push_back({ &layer, pi, e, victim }); + } + } + + layer.in_use[(size_t) slot] = 1; + out[i] = slot; + } + + if (tasks.empty()) { + return; + } + + if (tasks.size() == 1) { + if (!pread_pool(*tasks[0].layer, tasks[0].pool_idx, tasks[0].expert_id, tasks[0].dst_slot)) { + GGML_ABORT("moe_offloader: pread failed"); + } + } else { + __block bool failed = false; + dispatch_apply(tasks.size(), DISPATCH_APPLY_AUTO, ^(size_t i) { + const moe_pread_task & t = tasks[i]; + if (!pread_pool(*t.layer, t.pool_idx, t.expert_id, t.dst_slot)) { + failed = true; + } + }); + if (failed) { + GGML_ABORT("moe_offloader: parallel pread failed"); + } + } + + std::atomic_thread_fence(std::memory_order_release); +} + +void llama_moe_offloader::sidecar_loop() { + std::vector remap_buf(MOE_MAX_IDS); + + std::vector layers_snapshot; + + while (sidecar_run.load(std::memory_order_relaxed)) { + bool any = false; + + { + // will have zero contention after graph finalization + // and 48 pointers copy is cheap (no realloc) + std::lock_guard g(build_mtx); + layers_snapshot = sidecar_layers; + } + + for (auto * lp : layers_snapshot) { + moe_layer & layer = *lp; + + auto * req_a = (std::atomic *) (layer.mapped + MOE_OFF_REQ); + uint32_t req = req_a->load(std::memory_order_acquire); + uint32_t done = layer.last_processed.load(std::memory_order_relaxed); + + if (req <= done) { + continue; + } + + int n_ids = *(int32_t *) (layer.mapped + MOE_OFF_N); + GGML_ASSERT(n_ids > 0 && n_ids <= MOE_MAX_IDS); + const int32_t * ids = (const int32_t *) (layer.mapped + MOE_OFF_SELECTED); + + resolve(layer, ids, remap_buf.data(), n_ids); + + memcpy(layer.mapped + MOE_OFF_REMAPPED, remap_buf.data(), (size_t) n_ids * sizeof(int32_t)); + + layer.last_processed.store(req, std::memory_order_relaxed); + ggml_backend_metal_event_signal(layer.shared_event, (uint64_t) req); + + any = true; + } + + if (!any) { + usleep(10); + } + } +} + +void llama_moe_offloader::start_sidecar() { + if (sidecar_run.exchange(true)) { + return; + } + + sidecar_thread = std::thread([this] { sidecar_loop(); }); +} + +void llama_moe_offloader::stop_sidecar() { + if (!sidecar_run.exchange(false)) { + return; + } + if (sidecar_thread.joinable()) { + sidecar_thread.join(); + } +} + +void llama_moe_offloader::clear() { + std::lock_guard g(build_mtx); + layers.clear(); + sidecar_layers.clear(); + pool_to_layer.clear(); +} + +#endif diff --git a/src/llama-moe-offloader.h b/src/llama-moe-offloader.h new file mode 100644 index 00000000000..d042e67bb95 --- /dev/null +++ b/src/llama-moe-offloader.h @@ -0,0 +1,125 @@ +#pragma once + +#include "ggml-backend.h" +#include "ggml.h" + +#include +#include +#include +#include +#include +#include +#include +#include + +#if defined(__APPLE__) && defined(GGML_USE_METAL) + +#include "ggml-metal.h" + + + +struct moe_pool { + ggml_tensor * tensor = nullptr; + ggml_context * own_ctx = nullptr; + ggml_backend_buffer_t own_buf = nullptr; + uint64_t file_offset = 0; + int fd = -1; + size_t stride = 0; +}; + +struct moe_layer { + int layer_idx = -1; + int64_t n_slots = 0; + int64_t n_expert = 0; + + std::vector pools; + + std::unordered_map name_to_pool_idx; + + // lru state + std::vector expert_to_slot; + std::vector slot_to_expert; + std::vector lru_clock; + uint32_t lru_time = 0; + std::vector in_use; + + // shmem message + ggml_tensor * msg_tensor = nullptr; + ggml_context * msg_ctx = nullptr; + ggml_backend_buffer_t msg_buf = nullptr; + uint8_t * mapped = nullptr; // == msg_tensor->data + + // gpu-cpu sync + ggml_backend_metal_event_t shared_event = nullptr; + std::atomic next_seq{ 0 }; + std::atomic last_processed{ 0 }; + + // reuse detection for gate_up/down etc. + const ggml_tensor * last_src2 = nullptr; + uint32_t last_src2_seq = 0; + int last_src2_uses = 0; + + uint64_t total_hits = 0; + uint64_t total_misses = 0; + + ~moe_layer(); +}; + +class llama_moe_offloader { + public: + llama_moe_offloader(ggml_backend_t backend); + ~llama_moe_offloader(); + + void start(); + void stop(); + + ggml_tensor * bind_pool(int layer_idx, ggml_tensor * orig, int64_t n_slots, uint64_t file_offset, int fd); + + static bool hook(void * user_data, + const ggml_tensor * src0, + const ggml_tensor * src2, + ggml_metal_moe_intercept * out); + + void clear(); + + private: + static bool pread_pool(moe_layer & layer, size_t pool_idx, int32_t expert_id, int64_t dst_slot); + + static int64_t lru_evict(moe_layer & layer); + + static void resolve(moe_layer & layer, const int32_t * ids, int32_t * out, int n); + + void start_sidecar(); + void stop_sidecar(); + void sidecar_loop(); + + ggml_backend_t backend; + + std::unordered_map pool_to_layer; + + std::mutex build_mtx; + std::vector> layers; + std::vector sidecar_layers; + + std::thread sidecar_thread; + std::atomic sidecar_run{ false }; +}; + +#else // !(__APPLE__ && GGML_USE_METAL) + +class llama_moe_offloader { + public: + llama_moe_offloader(ggml_backend_t) {} + + ~llama_moe_offloader() {} + + void start() {} + + void stop() {} + + ggml_tensor * bind_pool(int, ggml_tensor *, int64_t, uint64_t, int) { return nullptr; } + + void clear() {} +}; + +#endif // __APPLE__ && GGML_USE_METAL