# Applies to llama.cpp commit a3b1eff # MTP recurrent-state rollback only; excludes the cumulative Eaman ROCm/MTP/VEC/MoE patch # Derived from experimental commits e04c6d9b, cccc8c78, be535675, and 732dd501 # Adds --spec-mtp-rs-depth, on-device replay checkpoints, fitter reservation, and host fallback # Experimental and opt-in; full recurrent depth remains the default # CPU/default build validated as llama.cpp build 1171; llama-server and focused test targets linked # test-arg-parser passed; isolated GPU/model runtime validation remains pending diff --git a/common/arg.cpp b/common/arg.cpp index 10b517a84efd514e4c340b444ee11f3c33a32a5c..f231a4473fd4e93a2d91eafae1831d12bf4d1930 100644 --- a/common/arg.cpp +++ b/common/arg.cpp @@ -1292,6 +1292,9 @@ bool common_params_parse(int argc, char ** argv, common_params & params, llama_e ctx_arg.params = params_org; return false; } + if (ctx_arg.params.speculative.draft.n_rs_seq > ctx_arg.params.speculative.draft.n_max) { + throw std::invalid_argument("--spec-mtp-rs-depth must not exceed --spec-draft-n-max"); + } if (ctx_arg.params.usage) { common_params_print_usage(ctx_arg); if (ctx_arg.print_usage) { @@ -4098,6 +4101,16 @@ common_params_context common_params_parser_init(common_params & params, llama_ex params.speculative.draft.n_max = value; } ).set_spec().set_examples({LLAMA_EXAMPLE_SPECULATIVE, LLAMA_EXAMPLE_LOOKUP, LLAMA_EXAMPLE_SERVER, LLAMA_EXAMPLE_CLI}).set_env("LLAMA_ARG_SPEC_DRAFT_N_MAX")); + add_opt(common_arg( + {"--spec-mtp-rs-depth"}, "N", + "number of target recurrent-state rollback snapshots for MTP; lower values save memory but replay accepted tokens after deep rejection (default: --spec-draft-n-max)", + [](common_params & params, int value) { + if (value < 1) { + throw std::invalid_argument("--spec-mtp-rs-depth must be at least 1"); + } + params.speculative.draft.n_rs_seq = value; + } + ).set_spec().set_examples({LLAMA_EXAMPLE_SPECULATIVE, LLAMA_EXAMPLE_SERVER, LLAMA_EXAMPLE_CLI}).set_env("LLAMA_ARG_SPEC_MTP_RS_DEPTH")); add_opt(common_arg( {"--spec-draft-n-min"}, "N", string_format("minimum number of draft tokens to use for speculative decoding (default: %d)", params.speculative.draft.n_min), diff --git a/common/common.cpp b/common/common.cpp index 6060e81b030fd47e29983007f716a9506a8d2399..68fc14e12db9957b5099a5797293c1107e378362 100644 --- a/common/common.cpp +++ b/common/common.cpp @@ -2245,6 +2245,8 @@ void common_prompt_checkpoint::clear() { data_tgt.clear(); data_dft.clear(); data_spec.clear(); + data_tgt_on_device = false; + data_dft_on_device = false; } void common_prompt_checkpoint::update_pos( @@ -2256,6 +2258,50 @@ void common_prompt_checkpoint::update_pos( this->pos_max = pos_max; } +static void common_prompt_checkpoint_save( + std::vector & data, + bool & on_device, + llama_context * ctx, + llama_seq_id seq_id, + llama_state_seq_flags flags, + const char * label) { + const bool requested_on_device = flags & LLAMA_STATE_SEQ_FLAGS_ON_DEVICE; + on_device = false; + + auto save = [&](llama_state_seq_flags save_flags) { + const size_t ckpt_size = llama_state_seq_get_size_ext(ctx, seq_id, save_flags); + if (ckpt_size == 0) { + return false; + } + + std::vector saved(ckpt_size); + const size_t n = llama_state_seq_get_data_ext(ctx, saved.data(), ckpt_size, seq_id, save_flags); + if (n != ckpt_size) { + return false; + } + + data.swap(saved); + return true; + }; + + bool saved = save(flags); + if (requested_on_device && !saved) { + COM_WRN("%s: ON_DEVICE %s checkpoint save failed; retrying with host storage\n", + __func__, label); + flags &= ~LLAMA_STATE_SEQ_FLAGS_ON_DEVICE; + saved = save(flags); + } + + if (!saved) { + if (flags & LLAMA_STATE_SEQ_FLAGS_ON_DEVICE) { + GGML_ABORT("checkpoint device save failed for %s\n", label); + } + GGML_ABORT("checkpoint size mismatch while saving %s\n", label); + } + + on_device = flags & LLAMA_STATE_SEQ_FLAGS_ON_DEVICE; +} + void common_prompt_checkpoint::update_tgt( llama_context * ctx, llama_seq_id seq_id, @@ -2264,14 +2310,7 @@ void common_prompt_checkpoint::update_tgt( return; } - const size_t ckpt_size = llama_state_seq_get_size_ext(ctx, seq_id, flags); - - data_tgt.resize(ckpt_size); - - const size_t n = llama_state_seq_get_data_ext(ctx, data_tgt.data(), ckpt_size, seq_id, flags); - if (n != ckpt_size) { - GGML_ABORT("checkpoint size mismatch: expected %zu, got %zu\n", ckpt_size, n); - } + common_prompt_checkpoint_save(data_tgt, data_tgt_on_device, ctx, seq_id, flags, "target"); } void common_prompt_checkpoint::update_dft( @@ -2282,14 +2321,7 @@ void common_prompt_checkpoint::update_dft( return; } - const size_t ckpt_size = llama_state_seq_get_size_ext(ctx, seq_id, flags); - - data_dft.resize(ckpt_size); - - const size_t n = llama_state_seq_get_data_ext(ctx, data_dft.data(), ckpt_size, seq_id, flags); - if (n != ckpt_size) { - GGML_ABORT("checkpoint size mismatch: expected %zu, got %zu\n", ckpt_size, n); - } + common_prompt_checkpoint_save(data_dft, data_dft_on_device, ctx, seq_id, flags, "draft"); } void common_prompt_checkpoint::load_tgt( @@ -2304,6 +2336,9 @@ void common_prompt_checkpoint::load_tgt( return; } + flags = (flags & ~LLAMA_STATE_SEQ_FLAGS_ON_DEVICE) | + (data_tgt_on_device ? LLAMA_STATE_SEQ_FLAGS_ON_DEVICE : 0); + const size_t n = llama_state_seq_set_data_ext(ctx, data_tgt.data(), data_tgt.size(), seq_id, flags); if (n != data_tgt.size()) { GGML_ABORT("checkpoint size mismatch: expected %zu, got %zu\n", data_tgt.size(), n); @@ -2322,6 +2357,9 @@ void common_prompt_checkpoint::load_dft( return; } + flags = (flags & ~LLAMA_STATE_SEQ_FLAGS_ON_DEVICE) | + (data_dft_on_device ? LLAMA_STATE_SEQ_FLAGS_ON_DEVICE : 0); + const size_t n = llama_state_seq_set_data_ext(ctx, data_dft.data(), data_dft.size(), seq_id, flags); if (n != data_dft.size()) { GGML_ABORT("checkpoint size mismatch: expected %zu, got %zu\n", data_dft.size(), n); @@ -2330,9 +2368,11 @@ void common_prompt_checkpoint::load_dft( void common_prompt_checkpoint::clear_tgt() { data_tgt.clear(); + data_tgt_on_device = false; } void common_prompt_checkpoint::clear_dft() { data_dft.clear(); data_spec.clear(); + data_dft_on_device = false; } diff --git a/common/common.h b/common/common.h index 1cab232ce3cab9b3264b143debcbbd5d69e2f6b3..966953dc669b1035ae78311cfbd85734013b96c2 100644 --- a/common/common.h +++ b/common/common.h @@ -324,6 +324,7 @@ struct common_params_model { struct common_params_speculative_draft { int32_t n_max = 3; // maximum number of tokens to draft during speculative decoding int32_t n_min = 0; // minimum number of draft tokens to use for speculative decoding + int32_t n_rs_seq = -1; // MTP recurrent-state rollback depth (-1 = n_max) float p_split = 0.1f; // speculative decoding split probability float p_min = 0.0f; // minimum speculative decoding probability (greedy) @@ -384,11 +385,18 @@ struct common_params_speculative { } uint32_t need_n_rs_seq() const { - bool needs_rs_seq = std::any_of(types.begin(), types.end(), [&](auto t) { - return t == COMMON_SPECULATIVE_TYPE_DRAFT_MTP || t == COMMON_SPECULATIVE_TYPE_DRAFT_EAGLE3 || t == COMMON_SPECULATIVE_TYPE_DRAFT_DFLASH || t == COMMON_SPECULATIVE_TYPE_DRAFT_DSPARK; + const bool needs_mtp = std::find(types.begin(), types.end(), COMMON_SPECULATIVE_TYPE_DRAFT_MTP) != types.end(); + const bool needs_other_rs = std::any_of(types.begin(), types.end(), [&](auto t) { + return t == COMMON_SPECULATIVE_TYPE_DRAFT_EAGLE3 || t == COMMON_SPECULATIVE_TYPE_DRAFT_DFLASH || t == COMMON_SPECULATIVE_TYPE_DRAFT_DSPARK; }); - return needs_rs_seq ? draft.n_max : 0u; + if (needs_other_rs) { + return draft.n_max; + } + if (needs_mtp) { + return draft.n_rs_seq >= 0 ? draft.n_rs_seq : draft.n_max; + } + return 0u; } }; @@ -1146,6 +1154,11 @@ struct common_prompt_checkpoint { std::vector data_tgt; std::vector data_dft; + // The actual storage mode used for each checkpoint. This can differ from + // the flags passed to update_* when an on-device save is not possible. + bool data_tgt_on_device = false; + bool data_dft_on_device = false; + // (optional) speculative-decoding implementation state stashed with the checkpoint // (e.g. eagle3's deferred-boundary g_embd row) std::vector data_spec; diff --git a/src/llama-memory-recurrent.cpp b/src/llama-memory-recurrent.cpp index ef82eb976ca76f1b999e225118320e7f1e379fcb..f7638544759a03f4c4d8cb682927e3ac0f23733b 100644 --- a/src/llama-memory-recurrent.cpp +++ b/src/llama-memory-recurrent.cpp @@ -788,7 +788,7 @@ void llama_memory_recurrent::state_write(llama_io_write_i & io, llama_seq_id seq } if ((flags & LLAMA_STATE_SEQ_FLAGS_ON_DEVICE) && cell_ranges.size() > 1) { - GGML_ABORT("cannot save/load multiple ranges of cells to/from device memory\n"); + throw std::runtime_error("cannot save/load multiple ranges of cells to/from device memory"); } // DEBUG CHECK: Sum of cell counts in ranges should equal the total cell count diff --git a/tests/test-arg-parser.cpp b/tests/test-arg-parser.cpp index ba58f852eb4f772ffad61ab6925738b331709aea..3fad40d422320dabbc859ff3713f34f840eda918 100644 --- a/tests/test-arg-parser.cpp +++ b/tests/test-arg-parser.cpp @@ -3,6 +3,7 @@ #include "download.h" #include "llama.h" #include "speculative.h" +#include "ggml-backend.h" #include #include @@ -197,6 +198,132 @@ static void test(void) { assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), params, LLAMA_EXAMPLE_SPECULATIVE)); assert(params.speculative.draft.n_max == 123); + { + common_params params_mtp; + argv = {"binary_name", "-m", "abc.gguf", "--spec-draft-n-max", "5", "--spec-mtp-rs-depth", "1"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), params_mtp, LLAMA_EXAMPLE_SPECULATIVE)); + assert(params_mtp.speculative.draft.n_max == 5); + assert(params_mtp.speculative.draft.n_rs_seq == 1); + + params_mtp.speculative.types = { COMMON_SPECULATIVE_TYPE_DRAFT_MTP }; + assert(params_mtp.speculative.need_n_rs_seq() == 1); + + params_mtp.speculative.types.push_back(COMMON_SPECULATIVE_TYPE_DRAFT_EAGLE3); + assert(params_mtp.speculative.need_n_rs_seq() == 5); + } + + { + common_params params_mtp; + argv = {"binary_name", "-m", "abc.gguf", "--spec-draft-n-max", "5"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), params_mtp, LLAMA_EXAMPLE_SPECULATIVE)); + assert(params_mtp.speculative.draft.n_rs_seq == -1); + + params_mtp.speculative.types = { COMMON_SPECULATIVE_TYPE_DRAFT_MTP }; + assert(params_mtp.speculative.need_n_rs_seq() == 5); + } + + { + common_params params_mtp; + argv = {"binary_name", "-m", "abc.gguf", "--spec-mtp-rs-depth", "1", "--spec-draft-n-max", "5"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), params_mtp, LLAMA_EXAMPLE_SPECULATIVE)); + assert(params_mtp.speculative.draft.n_rs_seq == 1); + } + + { + common_params params_mtp; + argv = {"binary_name", "-m", "abc.gguf", "--spec-draft-n-max", "5", "--spec-mtp-rs-depth", "0"}; + assert(false == common_params_parse(argv.size(), list_str_to_char(argv).data(), params_mtp, LLAMA_EXAMPLE_SPECULATIVE)); + } + + { + common_params params_mtp; + argv = {"binary_name", "-m", "abc.gguf", "--spec-draft-n-max", "5", "--spec-mtp-rs-depth=-1"}; + assert(false == common_params_parse(argv.size(), list_str_to_char(argv).data(), params_mtp, LLAMA_EXAMPLE_SPECULATIVE)); + } + + { + common_params params_mtp; + argv = {"binary_name", "-m", "abc.gguf", "--spec-mtp-rs-depth", "6", "--spec-draft-n-max", "5"}; + assert(false == common_params_parse(argv.size(), list_str_to_char(argv).data(), params_mtp, LLAMA_EXAMPLE_SPECULATIVE)); + } + + { + printf("test-arg-parser: test recurrent checkpoint reservation inputs\n\n"); + + common_params fit_params; + assert(fit_params.fit_params); + assert(fit_params.fit_params_target.size() == llama_max_devices()); + for (const size_t target : fit_params.fit_params_target) { + assert(target == 1024 * 1024 * 1024ULL); + } + + argv = {"binary_name", "-m", "abc.gguf", "--fit", "off", "--split-mode", "none", "--main-gpu", "1", + "--fit-target", "256,64", "--tensor-split", "3,1"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), fit_params, LLAMA_EXAMPLE_COMMON)); + assert(!fit_params.fit_params); + assert(fit_params.split_mode == LLAMA_SPLIT_MODE_NONE); + assert(fit_params.main_gpu == 1); + assert(fit_params.fit_params_target[0] == 256 * 1024 * 1024ULL); + assert(fit_params.fit_params_target[1] == 64 * 1024 * 1024ULL); + assert(fit_params.tensor_split[0] == 3.0f); + assert(fit_params.tensor_split[1] == 1.0f); + + auto checkpoint_reservation_needed = [](const common_params & p) { + return p.speculative.need_n_rs_seq() < static_cast(p.speculative.draft.n_max); + }; + + common_params mtp_full; + argv = {"binary_name", "-m", "abc.gguf", "--spec-draft-n-max", "5"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), mtp_full, LLAMA_EXAMPLE_SPECULATIVE)); + mtp_full.speculative.types = { COMMON_SPECULATIVE_TYPE_DRAFT_MTP }; + assert(mtp_full.speculative.need_n_rs_seq() == 5); + assert(!checkpoint_reservation_needed(mtp_full)); + + common_params mtp_reduced; + argv = {"binary_name", "-m", "abc.gguf", "--spec-draft-n-max", "5", "--spec-mtp-rs-depth", "2"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), mtp_reduced, LLAMA_EXAMPLE_SPECULATIVE)); + mtp_reduced.speculative.types = { COMMON_SPECULATIVE_TYPE_DRAFT_MTP }; + assert(mtp_reduced.speculative.need_n_rs_seq() == 2); + assert(checkpoint_reservation_needed(mtp_reduced)); + + common_params mtp_with_other_rs; + argv = {"binary_name", "-m", "abc.gguf", "--spec-draft-n-max", "5", "--spec-mtp-rs-depth", "2"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), mtp_with_other_rs, LLAMA_EXAMPLE_SPECULATIVE)); + mtp_with_other_rs.speculative.types = { + COMMON_SPECULATIVE_TYPE_DRAFT_MTP, + COMMON_SPECULATIVE_TYPE_DRAFT_EAGLE3, + }; + assert(mtp_with_other_rs.speculative.need_n_rs_seq() == 5); + assert(!checkpoint_reservation_needed(mtp_with_other_rs)); + + ggml_backend_load_all(); + std::vector gpu_devices; + for (size_t i = 0; i < ggml_backend_dev_count(); ++i) { + ggml_backend_dev_t dev = ggml_backend_dev_get(i); + if (ggml_backend_dev_type(dev) != GGML_BACKEND_DEVICE_TYPE_CPU) { + gpu_devices.push_back(dev); + } + } + if (gpu_devices.size() >= 2) { + const std::string device_arg = string_format("%s,%s", + ggml_backend_dev_name(gpu_devices[0]), ggml_backend_dev_name(gpu_devices[1])); + common_params multi_gpu; + argv = {"binary_name", "-m", "abc.gguf", "--device", device_arg, "--split-mode", "layer", + "--fit-target", "128,256", "--tensor-split", "3,1"}; + assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), multi_gpu, LLAMA_EXAMPLE_COMMON)); + assert(multi_gpu.devices.size() == 3); + assert(multi_gpu.devices[0] == gpu_devices[0]); + assert(multi_gpu.devices[1] == gpu_devices[1]); + assert(multi_gpu.devices[2] == nullptr); + assert(multi_gpu.fit_params_target[0] == 128 * 1024 * 1024ULL); + assert(multi_gpu.fit_params_target[1] == 256 * 1024 * 1024ULL); + assert(multi_gpu.tensor_split[0] == 3.0f); + assert(multi_gpu.tensor_split[1] == 1.0f); + } else { + printf("test-arg-parser: skip asymmetric multi-GPU input test (fewer than two non-CPU devices)\n\n"); + } + } + argv = {"binary_name", "-lm", "none"}; assert(true == common_params_parse(argv.size(), list_str_to_char(argv).data(), params, LLAMA_EXAMPLE_COMMON)); assert(params.load_mode == LLAMA_LOAD_MODE_NONE); diff --git a/tests/test-recurrent-state-rollback.cpp b/tests/test-recurrent-state-rollback.cpp index 5d1f0140b623cb9a44476f8e15a0279e9e3a9887..90ef3522d3af21dab92a040725633755fdfee7d9 100644 --- a/tests/test-recurrent-state-rollback.cpp +++ b/tests/test-recurrent-state-rollback.cpp @@ -8,15 +8,60 @@ #include #include -static llama_context * make_ctx(const common_params & params, llama_model * model) { +static llama_context * make_ctx( + const common_params & params, + llama_model * model, + uint32_t n_rs_seq = 8, + uint32_t n_seq_max = 1) { auto cparams = common_context_params_to_llama(params); - cparams.n_seq_max = 1; - cparams.n_rs_seq = 8; + cparams.n_seq_max = n_seq_max; + cparams.n_rs_seq = n_rs_seq; cparams.n_batch = std::max(cparams.n_batch, (uint32_t) (cparams.n_rs_seq + 1)); cparams.n_ubatch = std::max(cparams.n_ubatch, (uint32_t) (cparams.n_rs_seq + 1)); return llama_init_from_model(model, cparams); } +static bool test_fragmented_on_device_fallback( + const common_params & params, + llama_model * model, + const std::vector & tokens) { + llama_context * ctx = make_ctx(params, model, 1, 3); + if (ctx == nullptr) { + fprintf(stderr, "%s : failed to initialize fragmented context\n", __func__); + return false; + } + + llama_batch batch = llama_batch_init(3, 0, 3); + for (llama_seq_id seq_id = 0; seq_id < 3; ++seq_id) { + common_batch_add(batch, tokens[seq_id % tokens.size()], 0, { seq_id }, seq_id == 2); + } + bool ok = llama_decode(ctx, batch) == 0; + llama_batch_free(batch); + + if (ok) { + ok = llama_memory_seq_rm(llama_get_memory(ctx), 1, -1, -1); + } + + common_prompt_checkpoint ckpt; + constexpr llama_state_seq_flags flags = + LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY | LLAMA_STATE_SEQ_FLAGS_ON_DEVICE; + if (ok) { + ckpt.update_tgt(ctx, -1, flags); + if (ckpt.data_tgt_on_device) { + fprintf(stderr, "%s : fragmented checkpoint did not fall back to host storage\n", __func__); + ok = false; + } + } + if (ok) { + // The caller still requests ON_DEVICE here. load_tgt() must use the + // actual host mode recorded when the checkpoint was saved. + ckpt.load_tgt(ctx, -1, flags); + } + + llama_free(ctx); + return ok; +} + static bool decode_tokens(llama_context * ctx, const std::vector & tokens, uint32_t count) { llama_batch batch = llama_batch_init(count, 0, 1); for (uint32_t pos = 0; pos < count; ++pos) { @@ -35,6 +80,156 @@ static bool decode_one(llama_context * ctx, llama_token tok, llama_pos pos) { return ok; } +static bool decode_range( + llama_context * ctx, + const std::vector & tokens, + uint32_t first, + uint32_t count) { + if (count == 0) { + return true; + } + + llama_batch batch = llama_batch_init(count, 0, 1); + for (uint32_t i = 0; i < count; ++i) { + common_batch_add(batch, tokens[first + i], first + i, { 0 }, i + 1 == count); + } + const bool ok = llama_decode(ctx, batch) == 0; + llama_batch_free(batch); + return ok; +} + +static bool compare_logits( + llama_context * ctx_full, + llama_context * ctx_depth, + int n_vocab, + uint32_t n_accepted, + const char * stage, + float eps) { + const float * logits_full = llama_get_logits(ctx_full); + const float * logits_depth = llama_get_logits(ctx_depth); + if (logits_full == nullptr || logits_depth == nullptr) { + fprintf(stderr, "%s : missing %s logits after accepting %u draft tokens\n", __func__, stage, n_accepted); + return false; + } + + int argmax_full = 0; + int argmax_depth = 0; + float max_diff = 0.0f; + for (int token = 0; token < n_vocab; ++token) { + if (logits_full[token] > logits_full[argmax_full]) { + argmax_full = token; + } + if (logits_depth[token] > logits_depth[argmax_depth]) { + argmax_depth = token; + } + + const float diff = std::fabs(logits_full[token] - logits_depth[token]); + max_diff = std::max(max_diff, diff); + if (eps >= 0.0f && diff > eps) { + fprintf(stderr, "%s : %s logits mismatch after accepting %u draft tokens, token %d (%g != %g)\n", + __func__, stage, n_accepted, token, (double) logits_full[token], (double) logits_depth[token]); + return false; + } + } + if (eps < 0.0f && argmax_full != argmax_depth) { + fprintf(stderr, "%s : %s greedy token mismatch after accepting %u draft tokens (%d != %d, max logit diff %g)\n", + __func__, stage, n_accepted, argmax_full, argmax_depth, (double) max_diff); + return false; + } + if (eps < 0.0f) { + fprintf(stderr, "%s : %s greedy token %d preserved after accepting %u draft tokens (max logit diff %g)\n", + __func__, stage, argmax_full, n_accepted, (double) max_diff); + } + return true; +} + +static bool test_recompute_fallback( + const common_params & params, + llama_model * model, + const std::vector & input_tokens, + int n_vocab) { + constexpr uint32_t n_prefix = 4; + constexpr uint32_t n_draft = 5; + constexpr uint32_t n_verify = n_draft + 1; + + std::vector tokens = input_tokens; + tokens.resize(n_prefix + n_verify + 2, tokens.back()); + + for (uint32_t n_accepted = 0; n_accepted < n_draft; ++n_accepted) { + llama_context * ctx_full = make_ctx(params, model, n_draft); + llama_context * ctx_depth = make_ctx(params, model, 1); + if (ctx_full == nullptr || ctx_depth == nullptr) { + fprintf(stderr, "%s : failed to initialize fallback contexts\n", __func__); + llama_free(ctx_full); + llama_free(ctx_depth); + return false; + } + + bool ok = decode_range(ctx_full, tokens, 0, n_prefix) && + decode_range(ctx_depth, tokens, 0, n_prefix); + if (ok) { + ok = compare_logits(ctx_full, ctx_depth, n_vocab, n_accepted, "prefix", 1e-5f); + } + + constexpr llama_state_seq_flags partial_flags = + LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY | LLAMA_STATE_SEQ_FLAGS_ON_DEVICE; + common_prompt_checkpoint ckpt; + if (ok) { + ckpt.update_tgt(ctx_depth, 0, partial_flags); + if (!ckpt.data_tgt_on_device) { + fprintf(stderr, "%s : contiguous recurrent checkpoint unexpectedly fell back to host storage\n", __func__); + ok = false; + } + } + if (ok) { + ok = decode_range(ctx_full, tokens, n_prefix, n_verify) && + decode_range(ctx_depth, tokens, n_prefix, n_verify); + } + + const llama_pos rollback_pos = n_prefix + 1 + n_accepted; + const uint32_t n_rollback = n_draft - n_accepted; + if (ok) { + ok = llama_memory_seq_rm(llama_get_memory(ctx_full), 0, rollback_pos, -1); + } + if (ok && n_rollback <= 1) { + ok = llama_memory_seq_rm(llama_get_memory(ctx_depth), 0, rollback_pos, -1); + } else if (ok) { + ckpt.load_tgt(ctx_depth, 0, partial_flags); + ok = llama_memory_seq_rm(llama_get_memory(ctx_depth), 0, n_prefix, -1); + } + + const llama_pos replacement_pos = rollback_pos; + if (ok && n_rollback <= 1) { + ok = decode_one(ctx_full, tokens[n_prefix + n_verify], replacement_pos) && + decode_one(ctx_depth, tokens[n_prefix + n_verify], replacement_pos); + } else if (ok) { + std::vector replay_tokens = tokens; + replay_tokens[replacement_pos] = tokens[n_prefix + n_verify]; + ok = decode_one(ctx_full, tokens[n_prefix + n_verify], replacement_pos) && + decode_range(ctx_depth, replay_tokens, n_prefix, 2 + n_accepted); + } + if (ok) { + ok = compare_logits(ctx_full, ctx_depth, n_vocab, n_accepted, "replacement", -1.0f); + } + if (ok) { + ok = decode_one(ctx_full, tokens[n_prefix + n_verify + 1], replacement_pos + 1) && + decode_one(ctx_depth, tokens[n_prefix + n_verify + 1], replacement_pos + 1) && + compare_logits(ctx_full, ctx_depth, n_vocab, n_accepted, "continuation", -1.0f); + } + + llama_free(ctx_full); + llama_free(ctx_depth); + + if (!ok) { + fprintf(stderr, "%s : fallback validation failed after accepting %u draft tokens\n", __func__, n_accepted); + return false; + } + } + + fprintf(stderr, "%s : depth-1 recomputation preserves depth-5 greedy decisions at every rejection position\n", __func__); + return true; +} + int main(int argc, char ** argv) { std::setlocale(LC_NUMERIC, "C"); @@ -65,6 +260,23 @@ int main(int argc, char ** argv) { const llama_vocab * vocab = llama_model_get_vocab(model); const int n_vocab = llama_vocab_n_tokens(vocab); + std::vector tokens; + if (llama_vocab_type(vocab) == LLAMA_VOCAB_TYPE_NONE) { + tokens = { 1, 2, 3, 4, 5, 6, 7, 8, 9 }; + } else { + tokens = common_tokenize(vocab, "The quick brown fox jumps over the lazy dog", true); + } + if (tokens.empty()) { + fprintf(stderr, "%s : not enough prompt tokens\n", __func__); + return 1; + } + if (!test_fragmented_on_device_fallback(params, model, tokens)) { + return 1; + } + if (!test_recompute_fallback(params, model, tokens, n_vocab)) { + return 1; + } + llama_context * ctx_src = make_ctx(params, model); llama_context * ctx_dst = make_ctx(params, model); if (ctx_src == nullptr || ctx_dst == nullptr) { @@ -79,12 +291,6 @@ int main(int argc, char ** argv) { return 0; } - std::vector tokens; - if (llama_vocab_type(vocab) == LLAMA_VOCAB_TYPE_NONE) { - tokens = { 1, 2, 3, 4, 5, 6, 7, 8, 9 }; - } else { - tokens = common_tokenize(ctx_src, "The quick brown fox jumps over the lazy dog", true); - } const uint32_t n_rs_seq = llama_n_rs_seq(ctx_src); constexpr uint32_t n_rollback = 3; if (n_rs_seq < n_rollback) { @@ -93,10 +299,6 @@ int main(int argc, char ** argv) { llama_free(ctx_dst); return 0; } - if (tokens.empty()) { - fprintf(stderr, "%s : not enough prompt tokens\n", __func__); - return 1; - } tokens.resize(n_rs_seq + 1, tokens.back()); const uint32_t n_tokens = tokens.size(); @@ -161,6 +363,10 @@ int main(int argc, char ** argv) { constexpr llama_state_seq_flags partial_flags = LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY; common_prompt_checkpoint ckpt_partial; ckpt_partial.update_tgt(ctx_src, 0, partial_flags); + if (ckpt_partial.data_tgt_on_device) { + fprintf(stderr, "%s : host recurrent checkpoint recorded device storage\n", __func__); + return 1; + } ckpt_partial.load_tgt(ctx_dst, 0, partial_flags); if (!replay_and_compare("partial")) { diff --git a/tools/server/README.md b/tools/server/README.md index b63a0e6dac00334650536616b74d7c33d9c7f82a..77c1d33113329440e946b6f790d8168d7c6b4e33 100644 --- a/tools/server/README.md +++ b/tools/server/README.md @@ -254,6 +254,7 @@ For the full list of features, please refer to [server's changelog](https://gith | `--spec-draft-cpu-moe, -cmoed, --cpu-moe-draft` | keep all Mixture of Experts (MoE) weights in the CPU for the draft model
(env: LLAMA_ARG_SPEC_DRAFT_CPU_MOE) | | `--spec-draft-n-cpu-moe, --spec-draft-ncmoe, -ncmoed, --n-cpu-moe-draft N` | keep the Mixture of Experts (MoE) weights of the first N layers in the CPU for the draft model
(env: LLAMA_ARG_SPEC_DRAFT_N_CPU_MOE) | | `--spec-draft-n-max N` | number of tokens to draft for speculative decoding (default: 3)
(env: LLAMA_ARG_SPEC_DRAFT_N_MAX) | +| `--spec-mtp-rs-depth N` | number of target recurrent-state rollback snapshots for MTP; lower values save memory but replay accepted tokens after deep rejection (default: `--spec-draft-n-max`)
(env: LLAMA_ARG_SPEC_MTP_RS_DEPTH) | | `--spec-draft-n-min N` | minimum number of draft tokens to use for speculative decoding (default: 0)
(env: LLAMA_ARG_SPEC_DRAFT_N_MIN) | | `--spec-draft-p-split, --draft-p-split P` | speculative decoding split probability (default: 0.10)
(env: LLAMA_ARG_SPEC_DRAFT_P_SPLIT) | | `--spec-draft-p-min, --draft-p-min P` | minimum speculative decoding probability (greedy) (default: 0.00)
(env: LLAMA_ARG_SPEC_DRAFT_P_MIN) | diff --git a/tools/server/server-common.cpp b/tools/server/server-common.cpp index 585f65e83c655d3b8b7e398e8bf76552dc846f36..e57548fce8f7c9a3873e9556e0366c6b9d538214 100644 --- a/tools/server/server-common.cpp +++ b/tools/server/server-common.cpp @@ -83,6 +83,10 @@ json server_slot_stats::to_json() const { base["draft_n"] = n_draft_tokens; base["draft_n_accepted"] = n_draft_accepted; } + if (n_draft_replay_count > 0) { + base["draft_replay_count"] = n_draft_replay_count; + base["draft_replay_n"] = n_draft_replay_tokens; + } return base; } diff --git a/tools/server/server-common.h b/tools/server/server-common.h index 6488be344c6ad5d7100cae854b6047f5728f4989..5a3d135f5f514d8362f884b4aa858dac0b1725a0 100644 --- a/tools/server/server-common.h +++ b/tools/server/server-common.h @@ -353,6 +353,8 @@ struct server_slot_stats { uint64_t n_draft_tokens = 0; uint64_t n_draft_accepted = 0; uint64_t n_draft_verif_steps = 0; + uint64_t n_draft_replay_count = 0; + uint64_t n_draft_replay_tokens = 0; // these are absolute timestamps (in us) // note: must be signed - they are subtracted before the later ones are set diff --git a/tools/server/server-context.cpp b/tools/server/server-context.cpp index 2e7637c313f52588227c8bc2aa6f6c132c156a94..73bed4ff249876f5b23e322a89a32fa9e73d5ea7 100644 --- a/tools/server/server-context.cpp +++ b/tools/server/server-context.cpp @@ -636,6 +636,12 @@ struct server_slot { " acc per pos = (%s)\n", acceptance_rates_per_pos.c_str()); } + if (stats.n_draft_replay_count > 0) { + SLT_INF(*this, + " MTP replays = %10" PRIu64 " events / %5" PRIu64 " tokens\n", + stats.n_draft_replay_count, stats.n_draft_replay_tokens); + } + common_speculative_print_stats(spec); } @@ -1039,6 +1045,85 @@ private: } } + // The on-device recurrent speculative checkpoint is allocated lazily, so reserve + // the additional recurrent-state group before fitting the target model. + if (params_base.fit_params && spec_mtp) { + const uint32_t n_rs_seq = params_base.speculative.need_n_rs_seq(); + const int32_t draft_n_max = params_base.speculative.draft.n_max; + if (draft_n_max < 0) { + SRV_ERR("%s", "[spec] invalid negative MTP draft.n_max while fitting\n"); + return false; + } + if (n_rs_seq < static_cast(draft_n_max)) { + try { + common_params params_rs = params_base; + auto mparams_rs = common_model_params_to_llama(params_rs); + auto cparams_rs = common_context_params_to_llama(params_rs); + + std::vector devs_rs; + uint32_t hp_ngl_rs = 0; + uint32_t hp_nct_rs = 0; + uint32_t hp_nex_rs = 0; + cparams_rs.n_rs_seq = n_rs_seq; + const auto dmd_rs = common_get_device_memory_data( + params_base.model.path.c_str(), &mparams_rs, &cparams_rs, + devs_rs, hp_ngl_rs, hp_nct_rs, hp_nex_rs, GGML_LOG_LEVEL_ERROR); + + std::vector devs_rs_next; + uint32_t hp_ngl_rs_next = 0; + uint32_t hp_nct_rs_next = 0; + uint32_t hp_nex_rs_next = 0; + cparams_rs.n_rs_seq = n_rs_seq + 1; + const auto dmd_rs_next = common_get_device_memory_data( + params_base.model.path.c_str(), &mparams_rs, &cparams_rs, + devs_rs_next, hp_ngl_rs_next, hp_nct_rs_next, hp_nex_rs_next, GGML_LOG_LEVEL_ERROR); + + if (devs_rs.size() != devs_rs_next.size() || dmd_rs.size() < devs_rs.size() || + dmd_rs_next.size() < devs_rs_next.size()) { + throw std::runtime_error("target context device count changed during recurrent-state measurement"); + } + + // common_fit_params() consumes fit_params_target in the measured model-device + // order. The configured device list can contain a nullptr sentinel, CPU/ACCEL + // devices, or be reduced by split-mode none, so it is not a safe index map. + if (params_base.fit_params_target.size() < devs_rs.size()) { + throw std::runtime_error("fit_params_target has no entry for every target device"); + } + + size_t total = 0; + for (size_t j = 0; j < devs_rs.size(); ++j) { + const auto next_dev = std::find(devs_rs_next.begin(), devs_rs_next.end(), devs_rs[j]); + if (next_dev == devs_rs_next.end()) { + throw std::runtime_error("target context device mapping changed during recurrent-state measurement"); + } + const size_t next_index = next_dev - devs_rs_next.begin(); + if (next_index != j) { + throw std::runtime_error("target context device order changed during recurrent-state measurement"); + } + if (dmd_rs_next[j].context < dmd_rs[j].context) { + throw std::runtime_error("recurrent-state context measurement decreased"); + } + const size_t delta = dmd_rs_next[j].context - dmd_rs[j].context; + params_base.fit_params_target[j] += delta; + total += delta; + SRV_INF("[spec] recurrent checkpoint fit reservation: device %s, %.2f MiB " + "(context %.2f -> %.2f MiB)\n", + ggml_backend_dev_name(devs_rs[j]), delta / (1024.0 * 1024.0), + dmd_rs[j].context / (1024.0 * 1024.0), + dmd_rs_next[j].context / (1024.0 * 1024.0)); + } + if (total == 0) { + throw std::runtime_error("recurrent-state measurement produced no device reservation"); + } + SRV_INF("[spec] recurrent checkpoint fit reservation: %.2f MiB total\n", + total / (1024.0 * 1024.0)); + } catch (const std::exception & e) { + SRV_ERR("[spec] failed to reserve recurrent-state memory before fitting: %s\n", e.what()); + return false; + } + } + } + // optionally reserve VRAM for the draft / MTP context before fitting the target model if (params_base.fit_params) { if (has_spec) { @@ -3071,9 +3156,18 @@ private: (ctx_dft_seq_rm_type == COMMON_CONTEXT_SEQ_RM_TYPE_RS && draft.size() > llama_n_rs_seq(ctx_dft)); if (use_ckpt_tgt) { + llama_state_seq_flags ckpt_flags = LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY; + if (ctx_tgt_seq_rm_type == COMMON_CONTEXT_SEQ_RM_TYPE_RS) { + // Phase 1 proof of concept: avoid copying the recurrent fallback checkpoint + // through host memory on every speculative round. The device buffer is + // allocated lazily; load_model() reserves one measured recurrent-state + // group for it when automatic fitting and reduced RS depth are enabled. + ckpt_flags |= LLAMA_STATE_SEQ_FLAGS_ON_DEVICE; + } + //const int64_t t_start = ggml_time_us(); - ckpt.update_tgt(ctx_tgt, slot.id, LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY); + ckpt.update_tgt(ctx_tgt, slot.id, ckpt_flags); //const int64_t t_total = ggml_time_us() - t_start; //printf("checkpoint total: %f ms\n", t_total / 1000.0); @@ -3917,6 +4011,10 @@ private: ctx_tgt_seq_rm_type == COMMON_CONTEXT_SEQ_RM_TYPE_FULL || (ctx_tgt_seq_rm_type == COMMON_CONTEXT_SEQ_RM_TYPE_RS && n_rollback > llama_n_rs_seq(ctx_tgt)); + const bool use_rs_replay = + ctx_tgt_seq_rm_type == COMMON_CONTEXT_SEQ_RM_TYPE_RS && + n_rollback > llama_n_rs_seq(ctx_tgt); + // check for partial draft acceptance if (n_rollback > 0) { if (use_ckpt_tgt) { @@ -3928,11 +4026,20 @@ private: slot.spec_is_replay = true; slot.spec_draft = std::move(accepted); + if (use_rs_replay) { + slot.stats.n_draft_replay_count += 1; + slot.stats.n_draft_replay_tokens += slot.spec_draft.size() + 1; + } + const auto & ckpt = slot.spec_ckpt; SLT_DBG(slot, "restoring speculative checkpoint (pos_min = %d, pos_max = %d, size = %zu)\n", ckpt.pos_min, ckpt.pos_max, ckpt.size()); - ckpt.load_tgt(slot.ctx_tgt, slot.id, LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY); + llama_state_seq_flags ckpt_flags = LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY; + if (use_rs_replay) { + ckpt_flags |= LLAMA_STATE_SEQ_FLAGS_ON_DEVICE; + } + ckpt.load_tgt(slot.ctx_tgt, slot.id, ckpt_flags); if (slot.ctx_dft) { ckpt.load_dft(slot.ctx_dft, slot.id, LLAMA_STATE_SEQ_FLAGS_PARTIAL_ONLY); diff --git a/tools/server/tests/utils.py b/tools/server/tests/utils.py index 9171dbc02977e26ce5e2e9b876d73457583314e0..6e6ca5b82f399a0d3e30739a5de48fa4a5e2cbe1 100644 --- a/tools/server/tests/utils.py +++ b/tools/server/tests/utils.py @@ -99,6 +99,7 @@ class ServerProcess: spec_type: str | None = None spec_draft_n_min: int | None = None spec_draft_n_max: int | None = None + spec_mtp_rs_depth: int | None = None no_ui: bool | None = None jinja: bool | None = None reasoning_format: Literal['deepseek', 'none', 'nothink'] | None = None @@ -245,6 +246,8 @@ class ServerProcess: server_args.extend(["--spec-draft-n-max", self.spec_draft_n_max]) if self.spec_draft_n_min: server_args.extend(["--spec-draft-n-min", self.spec_draft_n_min]) + if self.spec_mtp_rs_depth: + server_args.extend(["--spec-mtp-rs-depth", self.spec_mtp_rs_depth]) if self.no_ui: server_args.append("--no-ui") if self.no_models_autoload: