From 4190a603d5db95da626c5ab5136fbdc2e16fb77a Mon Sep 17 00:00:00 2001 From: assouan <750048+assouan@users.noreply.github.com> Date: Sat, 22 Aug 2026 20:06:14 +0200 Subject: [PATCH 1/2] feat: async layer prefetching Adds async layer prefetching through `ModelManager`, allowing upcoming segments to be loaded ahead of execution. Layer streaming is now configurable end-to-end with resident-layer and prefetch-depth limits, exposed through the public API and the new `--resident-layers` and `--layer-prefetch-depth` CLI/server options. The runner derives the streaming policy from the graph cut, accounts for streaming allocations when enforcing VRAM budgets, and dynamically evicts or falls back when needed to stay within budget. Existing `new_sd_ctx` callers remain fully compatible. --- docs/performance.md | 13 +- examples/cli/main.cpp | 5 +- examples/common/common.cpp | 46 +++ examples/common/common.h | 8 +- examples/server/main.cpp | 5 +- include/stable-diffusion.h | 9 + src/core/ggml_extend.hpp | 696 +++++++++++++++++++++++++++++++++--- src/core/ggml_graph_cut.cpp | 270 +++++++++++--- src/core/ggml_graph_cut.h | 19 +- src/model_manager.cpp | 401 ++++++++++++++++++++- src/model_manager.h | 48 ++- src/stable-diffusion.cpp | 75 +++- src/weight_manager.h | 24 +- 13 files changed, 1485 insertions(+), 134 deletions(-) diff --git a/docs/performance.md b/docs/performance.md index ccfe4778a..61985038a 100644 --- a/docs/performance.md +++ b/docs/performance.md @@ -55,21 +55,28 @@ See [backend selection](./backend.md) for full syntax. ## Run models that don't fit in VRAM (CPU streaming). -`--offload-to-cpu` alone keeps every parameter in system RAM and stages it to the runtime backend on first use, then leaves it resident there. If the diffusion model is larger than the runtime backend's free memory (e.g. Flux dev at bf16 on an 8 GiB GPU), that residency stops fitting during the sampling loop and generation fails. Two additional flags make it fit by trading a small amount of speed for room: +`--offload-to-cpu` alone keeps every parameter in system RAM and stages it to the runtime backend on first use, then leaves it resident there. If the diffusion model is larger than the runtime backend's free memory (e.g. Flux dev at bf16 on an 8 GiB GPU), that residency stops fitting during the sampling loop and generation fails. The following flags make it fit by trading a small amount of speed for room: - `--max-vram ` sets a VRAM budget the graph-cut segmenter respects. It cuts each forward pass into segments sized to fit the budget, running them in sequence and freeing intermediate activations between them. Negative values auto-detect free VRAM and spare the given amount (`--max-vram -1` uses most of the free VRAM and keeps ~1 GiB headroom), a positive value caps the budget, `0` disables segmentation. - `--stream-layers` streams the diffusion model's transformer blocks one at a time. Each block's parameters are copied from the CPU to the runtime backend just before it runs and evicted when the residency budget is reached. Prefetching hides most of the copy latency behind compute. This flag only takes effect when the diffusion params backend is CPU, so it must be combined with `--offload-to-cpu` (or an explicit `--params-backend diffusion=cpu`); a warning is logged and the flag is ignored otherwise. +- `--resident-layers ` sets the maximum number of leading parameter-bearing graph-cut segments kept resident (default: `-1`; `auto` is an alias for `-1`). `-1` uses as many as the VRAM budget permits, `0` keeps none, and a positive `N` keeps up to `N`. +- `--layer-prefetch-depth ` sets the maximum number of future graph-cut segments copied through a separate transfer backend or queue when supported while the active segment computes (default: `0`). `0` disables asynchronous prefetching, `1` overlaps the next segment, and larger values provide deeper lookahead when the VRAM budget permits. -The three flags stack. The recommended shape for "biggest model my card can host": +Both controls require `--stream-layers` and a non-zero `--max-vram`. Prefetch is budgeted before residency; requested limits are reduced as needed, a positive VRAM cap is never exceeded to force either one, and the active segment is never evicted. Residents persist only across repeated sampling steps, and stale graph state is released automatically. An explicit non-negative residency limit uses an unmerged plan, which may add dispatch overhead. + +These flags stack. The recommended shape for "biggest model my card can host": ```shell sd-cli --diffusion-model flux1-dev.safetensors ... \ - --offload-to-cpu --max-vram -1 --stream-layers + --offload-to-cpu --max-vram -1 --stream-layers \ + --resident-layers auto --layer-prefetch-depth 1 ``` - `--offload-to-cpu`: params in RAM, staged as needed. - `--max-vram -1`: use most of the free VRAM as the compute budget, spare 1 GiB headroom, let the graph-cut segmenter split each forward pass to fit. - `--stream-layers`: on top of the segmenter, stream individual transformer blocks so their weights don't all need to be resident at once. +- `--resident-layers auto`: use the remaining budget for a leading resident prefix. +- `--layer-prefetch-depth 1`: prepare the next parameter-bearing segment concurrently with the current segment's computation. Ordered from fastest to smallest-VRAM: no flags → `--offload-to-cpu` → `--offload-to-cpu --max-vram ` → `--offload-to-cpu --max-vram --stream-layers`. Each step down costs a few percent of throughput to buy more room; combined they can run models roughly 3-4x larger than the raw VRAM would allow. diff --git a/examples/cli/main.cpp b/examples/cli/main.cpp index 1cc7a7af4..b43924a55 100644 --- a/examples/cli/main.cpp +++ b/examples/cli/main.cpp @@ -894,7 +894,8 @@ int main(int argc, const char* argv[]) { } } - sd_ctx_params_t sd_ctx_params = ctx_params.to_sd_ctx_params_t(cli_params.taesd_preview); + sd_ctx_params_t sd_ctx_params = ctx_params.to_sd_ctx_params_t(cli_params.taesd_preview); + sd_layer_stream_params_t layer_stream_params = ctx_params.to_sd_layer_stream_params_t(); SDImageVec results; int num_results = 0; @@ -904,7 +905,7 @@ int main(int argc, const char* argv[]) { num_results = 1; results.push_back(gen_params.init_image.release()); } else { - SDCtxPtr sd_ctx(new_sd_ctx(&sd_ctx_params)); + SDCtxPtr sd_ctx(new_sd_ctx_with_layer_stream(&sd_ctx_params, &layer_stream_params)); if (sd_ctx == nullptr) { LOG_INFO("new_sd_ctx_t failed"); diff --git a/examples/common/common.cpp b/examples/common/common.cpp index 35812157c..a1f8ce0a6 100644 --- a/examples/common/common.cpp +++ b/examples/common/common.cpp @@ -505,6 +505,11 @@ ArgOptions SDContextParams::get_options() { "maximum VRAM budget in GiB for graph-cut segmented execution. Accepts a single value or assignments by backend/device, e.g. 6 or cuda0=6,vulkan0=4. 0 disables graph splitting; a negative value auto-detects free VRAM, sparing the specified value", 0, &max_vram}, + {"", + "--resident-layers", + "maximum leading parameter-bearing graph-cut segments kept resident with --stream-layers: -1 selects automatically from the VRAM budget, auto is an alias for -1, 0 keeps none, N keeps up to N (default: -1)", + 0, + &resident_layers_spec}, }; options.int_options = { @@ -513,6 +518,10 @@ ArgOptions SDContextParams::get_options() { "number of threads to use during computation (default: -1). " "If threads <= 0, then threads will be set to the number of CPU physical cores", &n_threads}, + {"", + "--layer-prefetch-depth", + "number of future graph-cut segments to prefetch with --stream-layers (default: 0; 0 disables prefetching, 1 overlaps the next segment, 2+ enables deeper lookahead when VRAM permits)", + &layer_prefetch_depth}, }; options.bool_options = { @@ -721,6 +730,25 @@ bool SDContextParams::resolve(SDMode mode) { n_threads = sd_get_num_physical_cores(); } + std::string resident_spec = resident_layers_spec; + std::transform(resident_spec.begin(), resident_spec.end(), resident_spec.begin(), [](unsigned char c) { + return static_cast(std::tolower(c)); + }); + if (resident_spec == "auto") { + resident_layers = -1; + } else { + try { + size_t parsed = 0; + resident_layers = std::stoi(resident_layers_spec, &parsed); + if (parsed != resident_layers_spec.size()) { + LOG_ERROR("error: --resident-layers must be auto, -1, 0, or a positive integer"); + return false; + } + } catch (const std::exception&) { + LOG_ERROR("error: --resident-layers must be auto, -1, 0, or a positive integer"); + return false; + } + } build_embedding_map(); return true; @@ -754,6 +782,14 @@ bool SDContextParams::validate(SDMode mode) { LOG_ERROR("error: vae_format must be 'auto', 'flux', 'sd3', 'flux2', or 'wan'"); return false; } + if (resident_layers < -1) { + LOG_ERROR("error: --resident-layers must be auto, -1, 0, or a positive integer"); + return false; + } + if (layer_prefetch_depth < 0) { + LOG_ERROR("error: --layer-prefetch-depth must be >= 0"); + return false; + } return true; } @@ -832,6 +868,8 @@ std::string SDContextParams::to_string() const { << " offload_params_to_cpu: " << (offload_params_to_cpu ? "true" : "false") << ",\n" << " max_vram: \"" << max_vram << "\",\n" << " stream_layers: " << (stream_layers ? "true" : "false") << ",\n" + << " resident_layers: " << resident_layers << ",\n" + << " layer_prefetch_depth: " << layer_prefetch_depth << ",\n" << " eager_load: " << (eager_load ? "true" : "false") << ",\n" << " backend: \"" << backend << "\",\n" << " params_backend: \"" << params_backend << "\",\n" @@ -914,6 +952,14 @@ sd_ctx_params_t SDContextParams::to_sd_ctx_params_t(bool taesd_preview) { return sd_ctx_params; } +sd_layer_stream_params_t SDContextParams::to_sd_layer_stream_params_t() const { + sd_layer_stream_params_t params; + sd_layer_stream_params_init(¶ms); + params.resident_layers = resident_layers; + params.layer_prefetch_depth = layer_prefetch_depth; + return params; +} + SDGenerationParams::SDGenerationParams() { sd_sample_params_init(&sample_params); sd_sample_params_init(&high_noise_sample_params); diff --git a/examples/common/common.h b/examples/common/common.h index 34b4a013b..01f259ce0 100644 --- a/examples/common/common.h +++ b/examples/common/common.h @@ -150,8 +150,11 @@ struct SDContextParams { rng_type_t sampler_rng_type = RNG_TYPE_COUNT; bool offload_params_to_cpu = false; std::string max_vram = "0"; - bool stream_layers = false; - bool eager_load = false; + bool stream_layers = false; + std::string resident_layers_spec = "-1"; + int resident_layers = -1; + int layer_prefetch_depth = 0; + bool eager_load = false; std::string backend; std::string params_backend; std::string split_mode; @@ -183,6 +186,7 @@ struct SDContextParams { bool resolve_and_validate(SDMode mode); std::string to_string() const; sd_ctx_params_t to_sd_ctx_params_t(bool taesd_preview); + sd_layer_stream_params_t to_sd_layer_stream_params_t() const; }; struct SDGenerationParams { diff --git a/examples/server/main.cpp b/examples/server/main.cpp index dce35c118..eb4bea6c7 100644 --- a/examples/server/main.cpp +++ b/examples/server/main.cpp @@ -85,8 +85,9 @@ int main(int argc, const char** argv) { LOG_DEBUG("%s", ctx_params.to_string().c_str()); LOG_DEBUG("%s", default_gen_params.to_string().c_str()); - sd_ctx_params_t sd_ctx_params = ctx_params.to_sd_ctx_params_t(false); - SDCtxPtr sd_ctx(new_sd_ctx(&sd_ctx_params)); + sd_ctx_params_t sd_ctx_params = ctx_params.to_sd_ctx_params_t(false); + sd_layer_stream_params_t layer_stream_params = ctx_params.to_sd_layer_stream_params_t(); + SDCtxPtr sd_ctx(new_sd_ctx_with_layer_stream(&sd_ctx_params, &layer_stream_params)); if (sd_ctx == nullptr) { LOG_ERROR("new_sd_ctx_t failed"); diff --git a/include/stable-diffusion.h b/include/stable-diffusion.h index bab62bac9..8d4d00172 100644 --- a/include/stable-diffusion.h +++ b/include/stable-diffusion.h @@ -237,6 +237,12 @@ typedef struct { const char* model_args; } sd_ctx_params_t; +typedef struct { + uint32_t struct_size; // Set by sd_layer_stream_params_init; permits future extension + int resident_layers; // With stream_layers: maximum leading graph-cut segments kept resident (-1 = automatic, 0 = none) + int layer_prefetch_depth; // With stream_layers: future graph-cut segments prefetched during compute (0 = disabled) +} sd_layer_stream_params_t; + typedef struct { uint32_t sample_rate; uint32_t channels; @@ -477,8 +483,11 @@ SD_API void sd_hires_params_init(sd_hires_params_t* hires_params); SD_API void sd_ctx_params_init(sd_ctx_params_t* sd_ctx_params); SD_API char* sd_ctx_params_to_str(const sd_ctx_params_t* sd_ctx_params); +SD_API void sd_layer_stream_params_init(sd_layer_stream_params_t* params); SD_API sd_ctx_t* new_sd_ctx(const sd_ctx_params_t* sd_ctx_params); +SD_API sd_ctx_t* new_sd_ctx_with_layer_stream(const sd_ctx_params_t* sd_ctx_params, + const sd_layer_stream_params_t* layer_stream_params); SD_API void free_sd_ctx(sd_ctx_t* sd_ctx); SD_API void free_sd_audio(sd_audio_t* audio); diff --git a/src/core/ggml_extend.hpp b/src/core/ggml_extend.hpp index 9ed2875ad..d364c1dd2 100644 --- a/src/core/ggml_extend.hpp +++ b/src/core/ggml_extend.hpp @@ -12,6 +12,7 @@ #include #include #include +#include #include #include #include @@ -1788,9 +1789,22 @@ struct GGMLRunner { size_t max_graph_vram_bytes = 0; bool stream_layers_enabled = false; + bool stream_residency_enabled = false; + int resident_segment_limit = -1; + int segment_prefetch_depth = 0; + size_t runtime_resident_segment_cap = SIZE_MAX; + bool stream_next_forward_prefetch = false; + bool stream_policy_logged = false; + bool stream_prefetch_fallback_warned = false; + bool stream_prefetch_runtime_warned = false; + bool stream_prefetch_reduced_warned = false; + bool stream_shared_params_logged = false; + bool stream_limits_unavailable_warned = false; size_t observed_max_effective_budget_ = 0; bool graph_cut_layer_split_enabled = false; std::vector graph_cut_layer_split_backend_vram_limits_; + std::vector> stream_prefetch_plan_signature_; + size_t stream_prefetch_plan_depth_ = 0; std::vector extra_runtime_backends; // borrowed (SDBackendManager-owned) ggml_backend_sched_t sched = nullptr; // owned @@ -1835,6 +1849,10 @@ struct GGMLRunner { return std::move(*tensor); } + static size_t saturating_size_add(size_t lhs, size_t rhs) { + return rhs > SIZE_MAX - lhs ? SIZE_MAX : lhs + rhs; + } + template static sd::Tensor restore_trailing_singleton_dims(std::optional> tensor, size_t expected_dim) { @@ -1925,15 +1943,10 @@ struct GGMLRunner { } ggml_tensor* canonical_param_tensor(ggml_tensor* tensor) { - if (tensor == nullptr) { - return nullptr; - } - if (params_tensor_set_.find(tensor) != params_tensor_set_.end()) { - return tensor; - } - if (tensor->view_src != nullptr && - params_tensor_set_.find(tensor->view_src) != params_tensor_set_.end()) { - return tensor->view_src; + for (ggml_tensor* current = tensor; current != nullptr; current = current->view_src) { + if (params_tensor_set_.find(current) != params_tensor_set_.end()) { + return current; + } } return nullptr; } @@ -2000,14 +2013,15 @@ struct GGMLRunner { return true; } - void free_compute_backend_param_tensors(const std::vector& tensors) { + size_t free_compute_backend_param_tensors(const std::vector& tensors) { if (tensors.empty()) { - return; + return 0; } auto manager = weight_manager.lock(); if (manager != nullptr) { - manager->release_compute_backend_params(tensors); + return manager->release_compute_backend_params(tensors); } + return 0; } void free_params_backend_param_tensors(const std::vector& tensors) { @@ -2451,7 +2465,8 @@ struct GGMLRunner { bool resolve_graph_cut_plan(ggml_cgraph* gf, GraphCutPlan* plan_out, - size_t* effective_budget_out = nullptr) { + size_t* effective_budget_out = nullptr, + sd::ggml_graph_cut::StreamingPolicy* streaming_policy_out = nullptr) { GGML_ASSERT(plan_out != nullptr); GGML_ASSERT(gf != nullptr); @@ -2462,13 +2477,31 @@ struct GGMLRunner { if (dev != nullptr && ggml_backend_dev_type(dev) != GGML_BACKEND_DEVICE_TYPE_CPU) { size_t free_vram = 0, total_vram = 0; ggml_backend_dev_memory(dev, &free_vram, &total_vram); + size_t runner_owned_vram = 0; + if (auto manager = weight_manager.lock()) { + runner_owned_vram = manager->streaming_allocation_bytes( + stream_prefetch_owner_id(), + runtime_backend, + kept_compute_param_tensor_set); + } + // Resident and queued buffers reduce reported free VRAM, but + // this runner can reuse or release them. Other allocations + // remain charged through free_vram. constexpr size_t safety_margin = 512ull * 1024 * 1024; - free_clamp = (free_vram > safety_margin) ? (free_vram - safety_margin) : 0; + size_t reclaimable_free = saturating_size_add(free_vram, runner_owned_vram); + if (total_vram > 0) { + reclaimable_free = std::min(reclaimable_free, total_vram); + } + free_clamp = reclaimable_free > safety_margin + ? reclaimable_free - safety_margin + : 0; if (free_clamp < effective_budget) { - LOG_DEBUG("%s clamping streaming budget: actual free VRAM %.2f MB < user cap %.2f MB", - get_desc().c_str(), - free_clamp / (1024.0 * 1024.0), - effective_budget / (1024.0 * 1024.0)); + LOG_DEBUG( + "%s clamping streaming budget: free VRAM %.2f MB + runner-owned %.2f MB < user cap %.2f MB", + get_desc().c_str(), + free_vram / (1024.0 * 1024.0), + runner_owned_vram / (1024.0 * 1024.0), + effective_budget / (1024.0 * 1024.0)); effective_budget = free_clamp; } } @@ -2480,8 +2513,6 @@ struct GGMLRunner { observed_max_effective_budget_ = effective_budget; budget_increased = true; } else { - // Keep the plan cache stable, but never plan above what is free now: - // another model or process can take VRAM after the first measurement. effective_budget = std::min(observed_max_effective_budget_, free_clamp); } } @@ -2490,18 +2521,28 @@ struct GGMLRunner { *effective_budget_out = effective_budget; } - // When streaming and the model dwarfs the budget, cap the planner at - // a quarter so it builds smaller merged segments and chunk-K can fit - // alongside. Without streaming the cap only adds dispatch overhead. - size_t planner_budget = effective_budget; - if (stream_layers_enabled) { + // Explicit residency counts use the unmerged plan so their units do + // not change when the VRAM budget merges adjacent graph-cut segments. + const bool explicit_resident_limit = stream_layers_enabled && + resident_segment_limit >= 0; + size_t planner_budget = explicit_resident_limit ? 0 : effective_budget; + // For automatic residency, cap merged segments at a quarter of a tight + // budget so chunk-K and the compute workspace can fit alongside them. + if (stream_layers_enabled && !explicit_resident_limit) { size_t total_params_bytes = 0; for (const ggml_tensor* t : params_tensor_set_) { if (t != nullptr) { - total_params_bytes += ggml_nbytes(t); + const size_t tensor_bytes = ggml_nbytes(t); + if (tensor_bytes > SIZE_MAX - total_params_bytes) { + total_params_bytes = SIZE_MAX; + break; + } + total_params_bytes += tensor_bytes; } } - if (total_params_bytes * 4 > effective_budget * 3) { + const size_t quarter_ceil = effective_budget / 4 + + (effective_budget % 4 != 0 ? 1 : 0); + if (total_params_bytes > effective_budget - quarter_ceil) { planner_budget = effective_budget / 4; } } @@ -2512,8 +2553,25 @@ struct GGMLRunner { planner_budget, params_tensor_set_, get_desc().c_str()); + sd::ggml_graph_cut::StreamingPolicy streaming_policy; if (stream_layers_enabled) { - sd::ggml_graph_cut::annotate_residency(*plan_out, effective_budget); + int effective_resident_limit = stream_residency_enabled + ? resident_segment_limit + : 0; + if (runtime_resident_segment_cap != SIZE_MAX && + (effective_resident_limit < 0 || + static_cast(effective_resident_limit) > runtime_resident_segment_cap)) { + effective_resident_limit = static_cast(std::min( + runtime_resident_segment_cap, + static_cast(std::numeric_limits::max()))); + } + streaming_policy = sd::ggml_graph_cut::annotate_residency(*plan_out, + effective_budget, + effective_resident_limit, + segment_prefetch_depth); + } + if (streaming_policy_out != nullptr) { + *streaming_policy_out = streaming_policy; } if (stream_layers_enabled) { if (budget_increased) { @@ -2646,6 +2704,315 @@ struct GGMLRunner { return true; } + uintptr_t stream_prefetch_owner_id() const { + return reinterpret_cast(this); + } + + void clear_stream_prefetch_state() { + if (auto manager = weight_manager.lock()) { + manager->clear_param_prefetches(stream_prefetch_owner_id()); + } + stream_prefetch_plan_signature_.clear(); + stream_prefetch_plan_depth_ = 0; + } + + void update_stream_prefetch_plan_signature( + const std::vector>& segment_params, + size_t prefetch_depth) { + if (stream_prefetch_plan_signature_ == segment_params && + stream_prefetch_plan_depth_ == prefetch_depth) { + return; + } + if (auto manager = weight_manager.lock()) { + manager->clear_param_prefetches(stream_prefetch_owner_id()); + } + stream_prefetch_plan_signature_ = segment_params; + stream_prefetch_plan_depth_ = prefetch_depth; + } + + static uint64_t segment_prefetch_id(size_t segment_index) { + return static_cast(segment_index + 1); + } + + bool collect_prefetch_segment_params( + ggml_cgraph* gf, + const GraphCutPlan& plan, + std::vector>& segment_params, + std::vector& param_segments, + std::vector& segment_to_param_position) { + segment_params.assign(plan.segments.size(), {}); + param_segments.clear(); + segment_to_param_position.assign(plan.segments.size(), SIZE_MAX); + + std::unordered_map segment_occurrences; + for (size_t segment_index = 0; segment_index < plan.segments.size(); ++segment_index) { + std::unordered_set seen_in_segment; + for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors(gf, + plan.segments[segment_index])) { + ggml_tensor* canonical = canonical_param_tensor(tensor); + if (canonical == nullptr) { + if (!stream_prefetch_fallback_warned) { + LOG_WARN("%s cannot canonicalize a streamed parameter; disabling segment prefetch", + get_desc().c_str()); + stream_prefetch_fallback_warned = true; + } + return false; + } + if (!seen_in_segment.insert(canonical).second) { + continue; + } + segment_params[segment_index].push_back(canonical); + ++segment_occurrences[canonical]; + } + } + + size_t shared_params = 0; + for (const auto& occurrence : segment_occurrences) { + if (occurrence.second > 1) { + ++shared_params; + } + } + if (shared_params > 0) { + for (auto& params : segment_params) { + params.erase(std::remove_if(params.begin(), params.end(), [&](ggml_tensor* param) { + return segment_occurrences.at(param) > 1; + }), + params.end()); + } + if (!stream_shared_params_logged) { + LOG_INFO("%s segment prefetch excludes %zu cross-segment parameters from asynchronous transfers", + get_desc().c_str(), + shared_params); + stream_shared_params_logged = true; + } + } + + bool any_prefetch_params = false; + for (size_t segment_index = 0; segment_index < plan.segments.size(); ++segment_index) { + if (plan.segments[segment_index].input_param_bytes > 0) { + segment_to_param_position[segment_index] = param_segments.size(); + param_segments.push_back(segment_index); + any_prefetch_params = any_prefetch_params || !segment_params[segment_index].empty(); + } + } + return any_prefetch_params; + } + + std::unordered_set retained_stream_params( + const GraphCutPlan& plan, + ggml_cgraph* gf, + size_t protected_segment_index = SIZE_MAX) { + std::unordered_set retained_params; + auto collect_segment = [&](const GraphCutSegment& segment) { + for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors(gf, segment)) { + if (ggml_tensor* canonical = canonical_param_tensor(tensor)) { + retained_params.insert(canonical); + } + } + }; + + for (const GraphCutSegment& segment : plan.segments) { + if (segment.residency == sd::ggml_graph_cut::SegmentResidency::RESIDENT) { + collect_segment(segment); + } + } + if (protected_segment_index < plan.segments.size()) { + collect_segment(plan.segments[protected_segment_index]); + } + return retained_params; + } + + size_t release_unretained_resident_params(const GraphCutPlan& plan, + ggml_cgraph* gf, + size_t protected_segment_index = SIZE_MAX) { + const std::unordered_set retained_params = + retained_stream_params(plan, gf, protected_segment_index); + std::vector params_to_release; + params_to_release.reserve(kept_compute_param_tensor_set.size()); + for (const ggml_tensor* tensor : kept_compute_param_tensor_set) { + if (retained_params.count(tensor) == 0) { + params_to_release.push_back(const_cast(tensor)); + } + } + + const size_t released_bytes = free_compute_backend_param_tensors(params_to_release); + for (ggml_tensor* tensor : params_to_release) { + kept_compute_param_tensor_set.erase(tensor); + } + return released_bytes; + } + + size_t evict_resident_segment_for_prefetch(GraphCutPlan& plan, + ggml_cgraph* gf, + size_t active_segment_index, + bool& residency_changed) { + bool demoted = false; + std::vector resident_segments; + for (size_t segment_index = 0; segment_index < plan.segments.size(); ++segment_index) { + if (plan.segments[segment_index].residency == + sd::ggml_graph_cut::SegmentResidency::RESIDENT) { + resident_segments.push_back(segment_index); + } + } + + for (size_t resident_position = resident_segments.size(); resident_position-- > 0;) { + const size_t candidate_index = resident_segments[resident_position]; + if (candidate_index == active_segment_index) { + continue; + } + + std::unordered_set earlier_resident_params; + auto collect_earlier_resident_params = [&](size_t segment_index) { + for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors( + gf, plan.segments[segment_index])) { + if (ggml_tensor* canonical = canonical_param_tensor(tensor)) { + earlier_resident_params.insert(canonical); + } + } + }; + for (size_t position = 0; position < resident_position; ++position) { + collect_earlier_resident_params(resident_segments[position]); + } + + const std::vector candidate_params = + sd::ggml_graph_cut::param_tensors(gf, plan.segments[candidate_index]); + const bool has_demotable_params = std::any_of( + candidate_params.begin(), + candidate_params.end(), + [&](ggml_tensor* tensor) { + ggml_tensor* canonical = canonical_param_tensor(tensor); + return canonical != nullptr && + kept_compute_param_tensor_set.count(canonical) != 0 && + earlier_resident_params.count(canonical) == 0; + }); + if (!has_demotable_params) { + continue; + } + + runtime_resident_segment_cap = std::min(runtime_resident_segment_cap, + resident_position); + stream_policy_logged = false; + for (size_t position = resident_position; position < resident_segments.size(); ++position) { + plan.segments[resident_segments[position]].residency = + sd::ggml_graph_cut::SegmentResidency::STREAMED; + } + residency_changed = true; + demoted = true; + const size_t released_bytes = release_unretained_resident_params( + plan, gf, active_segment_index); + if (released_bytes > 0) { + LOG_WARN( + "%s evicted resident segment %zu to make room for prefetch; " + "resident cap reduced to %zu, released %.2f MB", + get_desc().c_str(), + candidate_index + 1, + runtime_resident_segment_cap, + released_bytes / (1024.0 * 1024.0)); + return released_bytes; + } + } + if (demoted) { + LOG_WARN( + "%s exhausted resident evictions for prefetch without releasing a complete buffer; " + "resident cap reduced to %zu", + get_desc().c_str(), + runtime_resident_segment_cap); + } + return 0; + } + + ParamPrefetchResult enqueue_segment_prefetch( + size_t segment_index, + const std::vector>& segment_params) { + if (segment_index >= segment_params.size()) { + return ParamPrefetchResult::FAILURE; + } + auto manager = weight_manager.lock(); + if (manager == nullptr) { + LOG_ERROR("%s segment prefetch requires a weight manager", get_desc().c_str()); + return ParamPrefetchResult::FAILURE; + } + return manager->enqueue_param_prefetch(stream_prefetch_owner_id(), + segment_prefetch_id(segment_index), + segment_params[segment_index]); + } + + ParamPrefetchResult enqueue_segment_prefetch_with_eviction( + size_t segment_index, + size_t active_segment_index, + const std::vector>& segment_params, + GraphCutPlan& plan, + ggml_cgraph* gf, + bool& residency_changed) { + ParamPrefetchResult result = enqueue_segment_prefetch(segment_index, segment_params); + while (result == ParamPrefetchResult::ALLOCATION_FAILURE) { + const size_t released_bytes = evict_resident_segment_for_prefetch( + plan, gf, active_segment_index, residency_changed); + if (released_bytes == 0) { + break; + } + result = enqueue_segment_prefetch(segment_index, segment_params); + } + return result; + } + + bool activate_segment_prefetch(size_t segment_index, + const std::vector>& segment_params) { + if (segment_index >= segment_params.size()) { + return false; + } + auto manager = weight_manager.lock(); + if (manager == nullptr) { + LOG_ERROR("%s segment prefetch requires a weight manager", get_desc().c_str()); + return false; + } + return manager->activate_param_prefetch(stream_prefetch_owner_id(), + segment_prefetch_id(segment_index), + segment_params[segment_index]); + } + + struct SegmentPrefetchQueueResult { + ParamPrefetchResult result = ParamPrefetchResult::SUCCESS; + size_t queued_depth = 0; + }; + + SegmentPrefetchQueueResult queue_future_segment_prefetches( + size_t active_position, + const std::vector& param_segments, + const std::vector>& segment_params, + size_t depth, + bool wrap, + GraphCutPlan& plan, + ggml_cgraph* gf, + bool& residency_changed) { + SegmentPrefetchQueueResult queue_result; + if (param_segments.empty() || active_position >= param_segments.size()) { + return queue_result; + } + depth = std::min(depth, param_segments.size() - 1); + for (size_t offset = 1; offset <= depth; ++offset) { + size_t position = active_position + offset; + if (position >= param_segments.size()) { + if (!wrap) { + break; + } + position %= param_segments.size(); + } + queue_result.result = enqueue_segment_prefetch_with_eviction( + param_segments[position], + param_segments[active_position], + segment_params, + plan, + gf, + residency_changed); + if (queue_result.result != ParamPrefetchResult::SUCCESS) { + return queue_result; + } + ++queue_result.queued_depth; + } + return queue_result; + } + struct PersistentExternalBinding { ggml_backend_buffer_t buffer = nullptr; void* data = nullptr; @@ -2771,7 +3138,8 @@ struct GGMLRunner { bool free_compute_params, bool preserve_backend_tensor_data_map, bool no_return = false, - const std::unordered_set* cache_keep_names = nullptr) { + const std::unordered_set* cache_keep_names = nullptr, + const std::function& before_compute = {}) { std::vector graph_param_tensors; std::vector params_to_prepare; if (!prepare_execute_graph_weights(gf, graph_param_tensors, params_to_prepare, !free_compute_params)) { @@ -2835,6 +3203,9 @@ struct GGMLRunner { } copy_data_to_backend_tensor(gf, !preserve_backend_tensor_data_map); + if (before_compute) { + before_compute(); + } if (sd_backend_is_cpu(runtime_backend)) { sd_backend_cpu_set_n_threads(runtime_backend, n_threads); } @@ -2931,15 +3302,108 @@ struct GGMLRunner { template std::optional> compute_graph_cut_segments(ggml_cgraph* gf, - const GraphCutPlan& plan, + GraphCutPlan& plan, int n_threads, bool log_residency, - bool no_return = false) { + bool no_return = false, + const sd::ggml_graph_cut::StreamingPolicy* prefetch_policy = nullptr) { GGML_ASSERT(gf != nullptr); free_compute_buffer(); free_cache_ctx_and_buffer(); + auto manager = weight_manager.lock(); + if (stream_layers_enabled && !kept_compute_param_tensor_set.empty()) { + std::unordered_set desired_resident_params; + for (const auto& segment : plan.segments) { + if (segment.residency != sd::ggml_graph_cut::SegmentResidency::RESIDENT) { + continue; + } + for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors(gf, segment)) { + if (ggml_tensor* canonical = canonical_param_tensor(tensor)) { + desired_resident_params.insert(canonical); + } + } + } + const bool same_resident_params = + desired_resident_params.size() == kept_compute_param_tensor_set.size() && + std::all_of(desired_resident_params.begin(), + desired_resident_params.end(), + [&](const ggml_tensor* tensor) { + return kept_compute_param_tensor_set.count(tensor) != 0; + }); + if (!same_resident_params) { + if (manager != nullptr) { + manager->clear_param_prefetches(stream_prefetch_owner_id()); + } + std::vector params_to_release; + for (const ggml_tensor* tensor : kept_compute_param_tensor_set) { + if (desired_resident_params.count(tensor) == 0) { + params_to_release.push_back(const_cast(tensor)); + } + } + free_compute_backend_param_tensors(params_to_release); + for (ggml_tensor* tensor : params_to_release) { + kept_compute_param_tensor_set.erase(tensor); + } + } + } + std::vector> segment_params; + std::vector param_segments; + std::vector segment_to_param_position; + bool prefetch_active = prefetch_policy != nullptr && + prefetch_policy->prefetch_depth > 0 && + collect_prefetch_segment_params(gf, + plan, + segment_params, + param_segments, + segment_to_param_position); + size_t runtime_prefetch_depth = prefetch_active + ? prefetch_policy->prefetch_depth + : 0; + if (prefetch_active) { + update_stream_prefetch_plan_signature(segment_params, runtime_prefetch_depth); + } else { + clear_stream_prefetch_state(); + } + bool residency_changed = false; + + struct PrefetchCleanupGuard { + std::shared_ptr manager; + uintptr_t owner_id = 0; + bool keep = false; + + ~PrefetchCleanupGuard() { + if (!keep && manager != nullptr) { + manager->clear_param_prefetches(owner_id); + } + } + } prefetch_cleanup{manager, stream_prefetch_owner_id()}; + + auto disable_prefetch = [&](const char* reason) { + if (!stream_prefetch_runtime_warned) { + LOG_WARN("%s segment prefetch failed while %s; continuing with synchronous streaming", + get_desc().c_str(), + reason); + stream_prefetch_runtime_warned = true; + } + clear_stream_prefetch_state(); + prefetch_active = false; + }; + + if (prefetch_active) { + ParamPrefetchResult result = enqueue_segment_prefetch_with_eviction( + param_segments.front(), + param_segments.front(), + segment_params, + plan, + gf, + residency_changed); + if (result != ParamPrefetchResult::SUCCESS) { + disable_prefetch("queueing the first segment"); + } + } + std::unordered_map persistent_externals; snapshot_persistent_externals(plan, gf, persistent_externals); @@ -2947,14 +3411,46 @@ struct GGMLRunner { for (size_t seg_idx = 0; seg_idx < plan.segments.size(); ++seg_idx) { const auto& segment = plan.segments[seg_idx]; const bool is_last = seg_idx + 1 == plan.segments.size(); - auto future_cut_names = sd::ggml_graph_cut::collect_future_input_names(gf, plan, seg_idx); + size_t param_position = SIZE_MAX; + if (!segment_to_param_position.empty()) { + param_position = segment_to_param_position[seg_idx]; + } + const bool has_prefetch_params = param_position != SIZE_MAX; + auto future_cut_names = sd::ggml_graph_cut::collect_future_input_names(gf, plan, seg_idx); + + if (has_prefetch_params && prefetch_active) { + ParamPrefetchResult result = enqueue_segment_prefetch_with_eviction( + seg_idx, + seg_idx, + segment_params, + plan, + gf, + residency_changed); + if (result != ParamPrefetchResult::SUCCESS || + !activate_segment_prefetch(seg_idx, segment_params)) { + disable_prefetch("activating a segment"); + } + } + if (log_residency) { - LOG_DEBUG("%s graph cut executing segment %zu/%zu: %s (residency=%s)", - get_desc().c_str(), - seg_idx + 1, - plan.segments.size(), - segment.group_name.c_str(), - segment.residency == sd::ggml_graph_cut::SegmentResidency::RESIDENT ? "RESIDENT" : "STREAMED"); + if (prefetch_policy != nullptr) { + LOG_DEBUG("%s graph cut executing segment %zu/%zu: %s (residency=%s, prefetch=%zu)", + get_desc().c_str(), + seg_idx + 1, + plan.segments.size(), + segment.group_name.c_str(), + segment.residency == sd::ggml_graph_cut::SegmentResidency::RESIDENT + ? "RESIDENT" + : "STREAMED", + prefetch_active ? runtime_prefetch_depth : 0); + } else { + LOG_DEBUG("%s graph cut executing segment %zu/%zu: %s (residency=%s)", + get_desc().c_str(), + seg_idx + 1, + plan.segments.size(), + segment.group_name.c_str(), + segment.residency == sd::ggml_graph_cut::SegmentResidency::RESIDENT ? "RESIDENT" : "STREAMED"); + } } else { LOG_DEBUG("%s graph cut executing segment %zu/%zu: %s", get_desc().c_str(), @@ -2985,14 +3481,60 @@ struct GGMLRunner { ggml_context* segment_graph_ctx = nullptr; ggml_cgraph* segment_graph = sd::ggml_graph_cut::build_segment_graph(gf, segment, &segment_graph_ctx); const bool keep_segment_params = segment.residency == sd::ggml_graph_cut::SegmentResidency::RESIDENT; - auto segment_output = execute_graph(segment_graph, + std::function before_compute; + if (has_prefetch_params && prefetch_active) { + before_compute = [this, + param_position, + ¶m_segments, + &segment_params, + &plan, + gf, + &runtime_prefetch_depth, + &prefetch_active, + &residency_changed, + &disable_prefetch]() { + if (!prefetch_active) { + return; + } + SegmentPrefetchQueueResult result = queue_future_segment_prefetches( + param_position, + param_segments, + segment_params, + runtime_prefetch_depth, + stream_next_forward_prefetch, + plan, + gf, + residency_changed); + if (result.result == ParamPrefetchResult::SUCCESS) { + return; + } + if (result.result == ParamPrefetchResult::ALLOCATION_FAILURE && + result.queued_depth > 0) { + runtime_prefetch_depth = result.queued_depth; + if (!stream_prefetch_reduced_warned) { + LOG_WARN("%s reduced segment prefetch depth to %zu after exhausting resident evictions", + get_desc().c_str(), + runtime_prefetch_depth); + stream_prefetch_reduced_warned = true; + } + return; + } + disable_prefetch("queueing future segments"); + }; + } + auto segment_output = execute_graph(segment_graph, n_threads, true, !keep_segment_params, true, !is_last || no_return, - &future_cut_names); + &future_cut_names, + before_compute); ggml_free(segment_graph_ctx); + if (residency_changed) { + release_unretained_resident_params(plan, gf); + residency_changed = false; + } if (!segment_output.has_value()) { free_cache_ctx_and_buffer(); free_compute_buffer(); @@ -3005,12 +3547,19 @@ struct GGMLRunner { backend_tensor_data_map.clear(); free_cache_ctx_and_buffer(); free_compute_ctx(); + prefetch_cleanup.keep = prefetch_active; return output; } public: void runner_done() { + stream_residency_enabled = false; + runtime_resident_segment_cap = SIZE_MAX; + stream_next_forward_prefetch = false; + stream_policy_logged = false; + stream_prefetch_reduced_warned = false; free_compute_buffer(); + clear_stream_prefetch_state(); std::vector tensors_to_release = std::move(this->runner_param_tensors); this->runner_param_tensors.clear(); runner_param_tensor_set.clear(); @@ -3031,6 +3580,7 @@ struct GGMLRunner { } virtual ~GGMLRunner() { + clear_stream_prefetch_state(); free_compute_buffer(); free_params_ctx(); free_compute_ctx(); @@ -3188,18 +3738,56 @@ struct GGMLRunner { if (can_attempt_graph_cut_segmented_compute()) { GraphCutPlan plan; - if (!resolve_graph_cut_plan(gf, &plan)) { + sd::ggml_graph_cut::StreamingPolicy streaming_policy; + if (!resolve_graph_cut_plan(gf, &plan, nullptr, &streaming_policy)) { free_compute_ctx(); return std::nullopt; } if (should_use_graph_cut_segmented_compute(plan)) { + if (stream_layers_enabled && !stream_policy_logged) { + const size_t parameter_segments = static_cast(std::count_if( + plan.segments.begin(), plan.segments.end(), [](const GraphCutSegment& segment) { + return segment.input_param_bytes > 0; + })); + LOG_INFO( + "%s layer stream: %zu parameter segments, resident requested=%d selected=%zu, " + "prefetch requested=%d selected=%zu", + get_desc().c_str(), + parameter_segments, + resident_segment_limit, + streaming_policy.resident_segments, + segment_prefetch_depth, + streaming_policy.prefetch_depth); + stream_policy_logged = true; + } + return compute_graph_cut_segments(gf, plan, n_threads, stream_layers_enabled, - no_return); + no_return, + stream_layers_enabled && streaming_policy.prefetch_depth > 0 + ? &streaming_policy + : nullptr); + } + if (stream_layers_enabled && + (resident_segment_limit >= 0 || segment_prefetch_depth > 0) && + !stream_limits_unavailable_warned) { + LOG_WARN("%s streaming limits require at least two valid graph-cut segments; using full-graph execution", + get_desc().c_str()); + stream_limits_unavailable_warned = true; } } + clear_stream_prefetch_state(); + if (!kept_compute_param_tensor_set.empty()) { + std::vector resident_params; + resident_params.reserve(kept_compute_param_tensor_set.size()); + for (const ggml_tensor* tensor : kept_compute_param_tensor_set) { + resident_params.push_back(const_cast(tensor)); + } + free_compute_backend_param_tensors(resident_params); + kept_compute_param_tensor_set.clear(); + } return execute_graph(gf, n_threads, free_compute_buffer, @@ -3227,7 +3815,8 @@ struct GGMLRunner { } void set_max_graph_vram_bytes(size_t max_vram_bytes) { - max_graph_vram_bytes = max_vram_bytes; + max_graph_vram_bytes = max_vram_bytes; + observed_max_effective_budget_ = 0; } void set_stream_layers_enabled(bool enabled) { @@ -3239,6 +3828,27 @@ struct GGMLRunner { stream_layers_enabled = enabled; } + void set_stream_segment_limits(int resident_segments, int prefetch_depth) { + resident_segment_limit = std::max(-1, resident_segments); + segment_prefetch_depth = std::max(0, prefetch_depth); + runtime_resident_segment_cap = SIZE_MAX; + stream_policy_logged = false; + stream_limits_unavailable_warned = false; + } + + void set_stream_residency_enabled(bool enabled) { + if (stream_residency_enabled == enabled) { + return; + } + stream_residency_enabled = enabled; + runtime_resident_segment_cap = SIZE_MAX; + stream_policy_logged = false; + } + + void set_stream_next_forward_prefetch(bool enabled) { + stream_next_forward_prefetch = enabled; + } + void set_graph_cut_layer_split_enabled(bool enabled) { graph_cut_layer_split_enabled = enabled; if (!enabled) { diff --git a/src/core/ggml_graph_cut.cpp b/src/core/ggml_graph_cut.cpp index 542b7fe16..0d5e0de7c 100644 --- a/src/core/ggml_graph_cut.cpp +++ b/src/core/ggml_graph_cut.cpp @@ -21,6 +21,46 @@ namespace sd::ggml_graph_cut { static constexpr double MAX_VRAM_BYTES_PER_GIB = 1024.0 * 1024.0 * 1024.0; + static size_t saturating_add(size_t lhs, size_t rhs) { + return rhs > SIZE_MAX - lhs ? SIZE_MAX : lhs + rhs; + } + + static bool sum_fits(size_t lhs, size_t rhs, size_t limit) { + return lhs <= limit && rhs <= limit - lhs; + } + + static const ggml_tensor* canonical_params_tensor( + const std::unordered_set& params_tensor_set, + const ggml_tensor* tensor) { + for (const ggml_tensor* current = tensor; current != nullptr; current = current->view_src) { + if (params_tensor_set.find(current) != params_tensor_set.end()) { + return current; + } + } + return nullptr; + } + + static size_t tensor_backend_allocation_bytes(ggml_backend_t backend, + const ggml_tensor* tensor) { + if (tensor == nullptr) { + return 0; + } + ggml_backend_buffer_type_t buft = backend != nullptr + ? ggml_backend_get_default_buffer_type(backend) + : nullptr; + if (buft == nullptr) { + return ggml_nbytes(tensor); + } + size_t bytes = ggml_backend_buft_get_alloc_size(buft, tensor); + size_t alignment = ggml_backend_buft_get_alignment(buft); + if (alignment <= 1) { + return bytes; + } + return bytes > SIZE_MAX - (alignment - 1) + ? SIZE_MAX + : GGML_PAD(bytes, alignment); + } + static std::string graph_cut_tensor_display_name(const ggml_tensor* tensor) { if (tensor == nullptr) { return ""; @@ -44,12 +84,7 @@ namespace sd::ggml_graph_cut { static bool is_params_tensor(const std::unordered_set& params_tensor_set, const ggml_tensor* tensor) { - if (tensor == nullptr) { - return false; - } - return params_tensor_set.find(tensor) != params_tensor_set.end() || - (tensor->view_src != nullptr && - params_tensor_set.find(tensor->view_src) != params_tensor_set.end()); + return canonical_params_tensor(params_tensor_set, tensor) != nullptr; } static int graph_node_index_by_name(ggml_cgraph* gf, const char* name) { @@ -79,11 +114,18 @@ namespace sd::ggml_graph_cut { return shape; } - static size_t graph_cut_segment_vram_bytes(const Segment& segment) { - return segment.compute_buffer_size + - segment.input_param_bytes + - segment.input_previous_cut_bytes + - segment.output_bytes; + static bool graph_cut_segment_fits(const Segment& segment, size_t budget) { + size_t bytes = 0; + for (size_t allocation : {segment.compute_buffer_size, + segment.input_param_bytes, + segment.input_previous_cut_bytes, + segment.output_bytes}) { + if (!sum_fits(bytes, allocation, budget)) { + return false; + } + bytes += allocation; + } + return true; } static std::string lower_ascii_copy(std::string value) { @@ -350,6 +392,7 @@ namespace sd::ggml_graph_cut { const char* log_desc) { std::set internal_nodes; std::unordered_set input_seen; + std::unordered_set param_seen; std::vector input_refs; std::stack work_stack; @@ -426,20 +469,33 @@ namespace sd::ggml_graph_cut { : ggml_nbytes(current_input)); switch (input.type) { case Segment::INPUT_PREVIOUS_CUT: - segment.input_previous_cut_bytes += tensor_bytes; + segment.input_previous_cut_bytes = saturating_add(segment.input_previous_cut_bytes, + tensor_bytes); break; - case Segment::INPUT_PARAM: - segment.input_param_bytes += tensor_bytes; + case Segment::INPUT_PARAM: { + const ggml_tensor* canonical = canonical_params_tensor(params_tensor_set, + current_input); + if (canonical != nullptr && param_seen.insert(canonical).second) { + const size_t allocation_bytes = tensor_backend_allocation_bytes(backend, + canonical); + segment.input_param_bytes = saturating_add( + segment.input_param_bytes, + allocation_bytes); + segment.input_param_allocations.push_back({canonical, allocation_bytes}); + } break; + } case Segment::INPUT_EXTERNAL: default: - segment.input_external_bytes += tensor_bytes; + segment.input_external_bytes = saturating_add(segment.input_external_bytes, + tensor_bytes); break; } } for (int output_node_index : segment.output_node_indices) { - ggml_tensor* output = ggml_graph_node(gf, output_node_index); - segment.output_bytes += cache_tensor_bytes(output); + ggml_tensor* output = ggml_graph_node(gf, output_node_index); + segment.output_bytes = saturating_add(segment.output_bytes, + cache_tensor_bytes(output)); } segment.compute_buffer_size = measure_segment_compute_buffer(backend, gf, segment, log_desc); @@ -890,7 +946,7 @@ namespace sd::ggml_graph_cut { GGML_ASSERT(!single_plan.segments.empty()); size_t best_end_segment_index = start_segment_index; - bool can_merge_next_segment = graph_cut_segment_vram_bytes(single_plan.segments.back()) <= max_graph_vram_bytes; + bool can_merge_next_segment = graph_cut_segment_fits(single_plan.segments.back(), max_graph_vram_bytes); while (can_merge_next_segment && best_end_segment_index + 1 < base_plan.segments.size()) { const size_t next_end_segment_index = best_end_segment_index + 1; @@ -910,8 +966,7 @@ namespace sd::ggml_graph_cut { GGML_ASSERT(!candidate_plan.segments.empty()); const auto& candidate_segment = candidate_plan.segments.back(); - const size_t candidate_bytes = graph_cut_segment_vram_bytes(candidate_segment); - if (candidate_bytes > max_graph_vram_bytes) { + if (!graph_cut_segment_fits(candidate_segment, max_graph_vram_bytes)) { break; } @@ -995,54 +1050,171 @@ namespace sd::ggml_graph_cut { return resolved_plan; } - void annotate_residency(Plan& plan, size_t max_graph_vram_bytes) { + static void add_segment_param_bytes( + const Plan& plan, + size_t segment_index, + bool prefetch_only, + const std::unordered_map& param_occurrences, + std::unordered_set& seen_params, + std::unordered_set& seen_fallback_segments, + size_t& bytes) { + const Segment& segment = plan.segments[segment_index]; + size_t described_bytes = 0; + for (const Segment::ParamAllocation& allocation : segment.input_param_allocations) { + described_bytes = saturating_add(described_bytes, allocation.bytes); + if (allocation.tensor == nullptr) { + continue; + } + auto occurrence = param_occurrences.find(allocation.tensor); + if (prefetch_only && occurrence != param_occurrences.end() && occurrence->second > 1) { + continue; + } + if (seen_params.insert(allocation.tensor).second) { + bytes = saturating_add(bytes, allocation.bytes); + } + } + + if (described_bytes < segment.input_param_bytes && + seen_fallback_segments.insert(segment_index).second) { + bytes = saturating_add(bytes, segment.input_param_bytes - described_bytes); + } + } + + static size_t peak_streaming_param_bytes( + const Plan& plan, + const std::vector& param_segments, + const std::unordered_map& param_occurrences, + size_t resident_count, + size_t prefetch_depth) { + size_t peak = 0; + prefetch_depth = param_segments.size() < 2 + ? 0 + : std::min(prefetch_depth, param_segments.size() - 1); + + for (size_t active_position = 0; active_position < param_segments.size(); ++active_position) { + std::unordered_set seen_params; + std::unordered_set seen_fallback_segments; + size_t bytes = 0; + + for (size_t resident_position = 0; + resident_position < resident_count; + ++resident_position) { + add_segment_param_bytes(plan, + param_segments[resident_position], + false, + param_occurrences, + seen_params, + seen_fallback_segments, + bytes); + } + add_segment_param_bytes(plan, + param_segments[active_position], + false, + param_occurrences, + seen_params, + seen_fallback_segments, + bytes); + for (size_t offset = 1; offset <= prefetch_depth; ++offset) { + const size_t future_position = (active_position + offset) % param_segments.size(); + add_segment_param_bytes(plan, + param_segments[future_position], + true, + param_occurrences, + seen_params, + seen_fallback_segments, + bytes); + } + peak = std::max(peak, bytes); + } + return peak; + } + + StreamingPolicy annotate_residency(Plan& plan, + size_t max_graph_vram_bytes, + int resident_segment_limit, + int segment_prefetch_depth) { + StreamingPolicy policy; // Cached plans may be reused with a smaller live budget. for (auto& seg : plan.segments) { seg.residency = SegmentResidency::STREAMED; } - if (max_graph_vram_bytes == 0 || plan.segments.size() < 2) { - return; + + if (max_graph_vram_bytes == 0) { + return policy; } - bool any_param_bearing = false; - for (const auto& seg : plan.segments) { - if (seg.input_param_bytes > 0) { - any_param_bearing = true; - break; + std::vector param_segments; + std::unordered_map param_occurrences; + param_segments.reserve(plan.segments.size()); + for (size_t i = 0; i < plan.segments.size(); ++i) { + if (plan.segments[i].input_param_bytes > 0) { + param_segments.push_back(i); + std::unordered_set seen_in_segment; + for (const Segment::ParamAllocation& allocation : + plan.segments[i].input_param_allocations) { + if (allocation.tensor != nullptr && + seen_in_segment.insert(allocation.tensor).second) { + ++param_occurrences[allocation.tensor]; + } + } } } - if (!any_param_bearing) { - return; + if (param_segments.empty()) { + return policy; } - // Leave room for the largest active streamed segment. - size_t worst_streamed_footprint = 0; + const size_t resident_cap = resident_segment_limit < 0 + ? param_segments.size() + : std::min(static_cast(resident_segment_limit), + param_segments.size()); + const size_t prefetch_cap = param_segments.size() < 2 + ? 0 + : std::min(static_cast(std::max(0, segment_prefetch_depth)), + param_segments.size() - 1); + + size_t worst_non_param_footprint = 0; for (const auto& seg : plan.segments) { - const size_t seg_footprint = seg.input_param_bytes + - seg.compute_buffer_size + - seg.output_bytes + - seg.input_previous_cut_bytes + - seg.input_external_bytes; - if (seg_footprint > worst_streamed_footprint) { - worst_streamed_footprint = seg_footprint; + size_t seg_footprint = seg.compute_buffer_size; + seg_footprint = saturating_add(seg_footprint, seg.output_bytes); + seg_footprint = saturating_add(seg_footprint, seg.input_previous_cut_bytes); + seg_footprint = saturating_add(seg_footprint, seg.input_external_bytes); + if (seg_footprint > worst_non_param_footprint) { + worst_non_param_footprint = seg_footprint; } } constexpr size_t safety = 512ull * 1024 * 1024; - const size_t reserved = safety + worst_streamed_footprint; - - if (max_graph_vram_bytes <= reserved) { - return; + if (!sum_fits(safety, worst_non_param_footprint, max_graph_vram_bytes)) { + return policy; + } + const size_t base_reserved = safety + worst_non_param_footprint; + + for (size_t candidate = 1; candidate <= prefetch_cap; ++candidate) { + const size_t peak_param_bytes = peak_streaming_param_bytes(plan, + param_segments, + param_occurrences, + 0, + candidate); + if (!sum_fits(base_reserved, peak_param_bytes, max_graph_vram_bytes)) { + break; + } + policy.prefetch_depth = candidate; } - const size_t available = max_graph_vram_bytes - reserved; - size_t cumulative = 0; - for (auto& seg : plan.segments) { - if (cumulative + seg.input_param_bytes > available) { + for (size_t candidate = 1; candidate <= resident_cap; ++candidate) { + const size_t peak_param_bytes = peak_streaming_param_bytes(plan, + param_segments, + param_occurrences, + candidate, + policy.prefetch_depth); + if (!sum_fits(base_reserved, peak_param_bytes, max_graph_vram_bytes)) { break; } - seg.residency = SegmentResidency::RESIDENT; - cumulative += seg.input_param_bytes; + policy.resident_segments = candidate; + } + for (size_t i = 0; i < policy.resident_segments; ++i) { + plan.segments[param_segments[i]].residency = SegmentResidency::RESIDENT; } + return policy; } } // namespace sd::ggml_graph_cut diff --git a/src/core/ggml_graph_cut.h b/src/core/ggml_graph_cut.h index d955049d8..32486d2bc 100644 --- a/src/core/ggml_graph_cut.h +++ b/src/core/ggml_graph_cut.h @@ -33,11 +33,17 @@ namespace sd::ggml_graph_cut { int node_index = -1; }; + struct ParamAllocation { + const ggml_tensor* tensor = nullptr; + size_t bytes = 0; + }; + size_t compute_buffer_size = 0; size_t output_bytes = 0; size_t input_external_bytes = 0; size_t input_previous_cut_bytes = 0; size_t input_param_bytes = 0; + std::vector input_param_allocations; std::string group_name; std::vector internal_node_indices; std::vector output_node_indices; @@ -67,6 +73,11 @@ namespace sd::ggml_graph_cut { size_t budgeted_graph_cut_plan_max_vram_bytes = 0; }; + struct StreamingPolicy { + size_t resident_segments = 0; + size_t prefetch_depth = 0; + }; + static constexpr const char* GGML_RUNNER_CUT_PREFIX = "ggml_runner_cut:"; static constexpr const char* GGML_RUNNER_CUT_SUFFIX = "|"; @@ -122,8 +133,12 @@ namespace sd::ggml_graph_cut { const std::unordered_set& params_tensor_set, const char* log_desc); - // Mark leading segments resident when they fit after streamed-segment headroom. - void annotate_residency(Plan& plan, size_t max_graph_vram_bytes); + // Mark leading parameter-bearing segments resident and select a prefetch + // depth that fits after the active segment's headroom. + StreamingPolicy annotate_residency(Plan& plan, + size_t max_graph_vram_bytes, + int resident_segment_limit, + int segment_prefetch_depth); } // namespace sd::ggml_graph_cut #endif // __SD_CORE_GGML_GRAPH_CUT_H__ diff --git a/src/model_manager.cpp b/src/model_manager.cpp index 8db303dda..038337340 100644 --- a/src/model_manager.cpp +++ b/src/model_manager.cpp @@ -1,6 +1,7 @@ #include "model_manager.h" #include +#include #include #include #include @@ -16,6 +17,10 @@ static size_t aligned_offset(const void* buffer, size_t offset, size_t alignment return offset + align; } +static size_t saturating_add(size_t lhs, size_t rhs) { + return rhs > SIZE_MAX - lhs ? SIZE_MAX : lhs + rhs; +} + static bool lora_specs_equal(const std::vector& lhs, const std::vector& rhs) { if (lhs.size() != rhs.size()) { @@ -95,8 +100,352 @@ static bool device_supports_param_op(ggml_backend_dev_t device, return supported; } +ggml_backend_t ModelManager::prefetch_backend_for(ggml_backend_t compute_backend) { + auto existing = prefetch_backends_.find(compute_backend); + if (existing != prefetch_backends_.end()) { + return existing->second; + } + if (compute_backend == nullptr) { + return nullptr; + } + ggml_backend_dev_t device = ggml_backend_get_device(compute_backend); + if (device == nullptr || ggml_backend_dev_type(device) == GGML_BACKEND_DEVICE_TYPE_CPU) { + return nullptr; + } + ggml_backend_t transfer_backend = ggml_backend_dev_init(device, nullptr); + if (transfer_backend == nullptr) { + LOG_WARN("model manager failed to create a prefetch backend for %s", + ggml_backend_name(compute_backend)); + prefetch_backends_[compute_backend] = nullptr; + return nullptr; + } + prefetch_backends_[compute_backend] = transfer_backend; + return transfer_backend; +} + +void ModelManager::synchronize_prefetch_block(PrefetchBlock& block) { + if (block.event != nullptr) { + ggml_backend_event_synchronize(block.event); + ggml_backend_event_free(block.event); + block.event = nullptr; + } else if (block.transfer_backend != nullptr) { + ggml_backend_synchronize(block.transfer_backend); + } + block.transfer_backend = nullptr; +} + +void ModelManager::free_prefetch_block(PrefetchBlock& block) { + synchronize_prefetch_block(block); + block.staged_tensors.clear(); + if (block.buffer != nullptr) { + ggml_backend_buffer_free(block.buffer); + block.buffer = nullptr; + } + if (block.staging_ctx != nullptr) { + ggml_free(block.staging_ctx); + block.staging_ctx = nullptr; + } +} + +ParamPrefetchResult ModelManager::populate_prefetch_block(PrefetchBlock& block) { + if (block.states.empty() || block.compute_backend == nullptr) { + return ParamPrefetchResult::FAILURE; + } + + block.transfer_backend = prefetch_backend_for(block.compute_backend); + if (block.transfer_backend == nullptr) { + return ParamPrefetchResult::FAILURE; + } + + ggml_init_params init_params; + init_params.mem_size = std::max(1, block.states.size()) * ggml_tensor_overhead(); + init_params.mem_buffer = nullptr; + init_params.no_alloc = true; + block.staging_ctx = ggml_init(init_params); + if (block.staging_ctx == nullptr) { + LOG_WARN("model manager failed to create the segment prefetch tensor context"); + return ParamPrefetchResult::FAILURE; + } + + block.staged_tensors.reserve(block.states.size()); + for (TensorState* state : block.states) { + if (state == nullptr || state->tensor == nullptr || state->tensor->buffer == nullptr || + state->tensor->data == nullptr || + state->params_backend == nullptr || state->staged_to_compute_backend || + state->active_prepare_count > 0) { + LOG_WARN("model manager segment prefetch source state changed before transfer"); + return ParamPrefetchResult::FAILURE; + } + ggml_tensor* staging_tensor = ggml_dup_tensor(block.staging_ctx, state->tensor); + ggml_set_name(staging_tensor, state->tensor->name); + block.staged_tensors.push_back({state, staging_tensor}); + } + + ggml_backend_buffer_type_t buffer_type = ggml_backend_get_default_buffer_type(block.compute_backend); + if (buffer_type == nullptr) { + LOG_WARN("model manager failed to resolve the segment prefetch buffer type"); + return ParamPrefetchResult::FAILURE; + } + block.buffer = ggml_backend_alloc_ctx_tensors_from_buft(block.staging_ctx, buffer_type); + if (block.buffer == nullptr) { + LOG_DEBUG("model manager failed to allocate the segment prefetch weight buffer"); + return ParamPrefetchResult::ALLOCATION_FAILURE; + } + ggml_backend_buffer_set_usage(block.buffer, GGML_BACKEND_BUFFER_USAGE_WEIGHTS); + + for (const auto& pair : block.staged_tensors) { + TensorState* state = pair.first; + ggml_tensor* staging_tensor = pair.second; + const bool host_source = state->tensor->buffer != nullptr && + ggml_backend_buffer_is_host(state->tensor->buffer); + if (host_source && + (!ggml_is_contiguous(state->tensor) || !ggml_is_contiguous(staging_tensor) || + ggml_nbytes(state->tensor) != ggml_nbytes(staging_tensor))) { + LOG_WARN("model manager segment prefetch requires contiguous host parameter tensors"); + return ParamPrefetchResult::FAILURE; + } + } + + for (const auto& pair : block.staged_tensors) { + TensorState* state = pair.first; + ggml_tensor* staging_tensor = pair.second; + const bool host_source = state->tensor->buffer != nullptr && + ggml_backend_buffer_is_host(state->tensor->buffer); + if (host_source) { + ggml_backend_tensor_set_async(block.transfer_backend, + staging_tensor, + state->tensor->data, + 0, + ggml_nbytes(state->tensor)); + } else { + ggml_backend_tensor_copy_async(state->params_backend, + block.transfer_backend, + state->tensor, + staging_tensor); + } + } + + ggml_backend_dev_t device = ggml_backend_get_device(block.transfer_backend); + block.event = ggml_backend_event_new(device); + if (block.event != nullptr) { + ggml_backend_event_record(block.event, block.transfer_backend); + } + + LOG_DEBUG("model manager queued segment %" PRIu64 + " prefetch (%6.2f MB, %zu tensors) to %s", + block.key.segment_id, + ggml_backend_buffer_get_size(block.buffer) / (1024.f * 1024.f), + block.states.size(), + ggml_backend_name(block.compute_backend)); + return ParamPrefetchResult::SUCCESS; +} + +ParamPrefetchResult ModelManager::enqueue_param_prefetch( + uintptr_t owner_id, + uint64_t segment_id, + const std::vector& tensors) { + if (tensors.empty()) { + return ParamPrefetchResult::SUCCESS; + } + + std::vector required_states; + if (!resolve_required_tensor_states(tensors, required_states) || + !load_tensors_to_params_backend(required_states)) { + return ParamPrefetchResult::FAILURE; + } + + std::vector states; + states.reserve(required_states.size()); + ggml_backend_t compute_backend = nullptr; + for (TensorState* state : required_states) { + if (state == nullptr || should_ignore(*state) || is_optional_missing_tensor(state->name) || + state->compute_backend == state->params_backend || state->staged_to_compute_backend) { + continue; + } + if (state->active_prepare_count > 0) { + LOG_WARN("cannot prefetch active tensor '%s'", state->name.c_str()); + return ParamPrefetchResult::FAILURE; + } + if (compute_backend == nullptr) { + compute_backend = state->compute_backend; + } else if (compute_backend != state->compute_backend) { + LOG_WARN("segment prefetch cannot span multiple compute backends"); + return ParamPrefetchResult::FAILURE; + } + states.push_back(state); + } + PrefetchKey key{owner_id, segment_id}; + if (states.empty()) { + auto existing = prefetch_blocks_.find(key); + if (existing != prefetch_blocks_.end()) { + std::unique_ptr stale = std::move(existing->second); + prefetch_blocks_.erase(existing); + free_prefetch_block(*stale); + } + return ParamPrefetchResult::SUCCESS; + } + if (compute_backend == nullptr || sd_backend_is_cpu(compute_backend)) { + LOG_WARN("segment prefetch requires a non-CPU compute backend"); + return ParamPrefetchResult::FAILURE; + } + auto block = std::make_unique(); + block->key = key; + block->states = std::move(states); + block->compute_backend = compute_backend; + + auto existing = prefetch_blocks_.find(key); + if (existing != prefetch_blocks_.end()) { + const auto& existing_states = existing->second->states; + const bool same_states = existing->second->compute_backend == compute_backend && + existing_states.size() == block->states.size() && + std::is_permutation(existing_states.begin(), + existing_states.end(), + block->states.begin()); + if (same_states) { + return ParamPrefetchResult::SUCCESS; + } + clear_param_prefetches(owner_id); + } + + ParamPrefetchResult result = populate_prefetch_block(*block); + if (result != ParamPrefetchResult::SUCCESS) { + free_prefetch_block(*block); + return result; + } + prefetch_blocks_.emplace(key, std::move(block)); + return ParamPrefetchResult::SUCCESS; +} + +bool ModelManager::activate_param_prefetch(uintptr_t owner_id, + uint64_t segment_id, + const std::vector& tensors) { + std::vector required_states; + if (!resolve_required_tensor_states(tensors, required_states)) { + return false; + } + PrefetchKey key{owner_id, segment_id}; + const bool already_staged = std::all_of(required_states.begin(), required_states.end(), + [&](TensorState* state) { + return state == nullptr || should_ignore(*state) || + is_optional_missing_tensor(state->name) || + state->compute_backend == state->params_backend || + state->staged_to_compute_backend; + }); + if (already_staged) { + auto existing = prefetch_blocks_.find(key); + if (existing != prefetch_blocks_.end()) { + std::unique_ptr stale = std::move(existing->second); + prefetch_blocks_.erase(existing); + free_prefetch_block(*stale); + } + return true; + } + + auto existing = prefetch_blocks_.find(key); + if (existing == prefetch_blocks_.end()) { + LOG_WARN("segment %" PRIu64 " was not queued for prefetch", segment_id); + return false; + } + std::unique_ptr block = std::move(existing->second); + prefetch_blocks_.erase(existing); + synchronize_prefetch_block(*block); + + LOG_DEBUG("model manager activated prefetched segment %" PRIu64 + " (%6.2f MB, %zu tensors) on %s", + segment_id, + ggml_backend_buffer_get_size(block->buffer) / (1024.f * 1024.f), + block->states.size(), + ggml_backend_name(block->compute_backend)); + + for (const auto& pair : block->staged_tensors) { + TensorState* state = pair.first; + ggml_tensor* staging_tensor = pair.second; + if (state == nullptr || state->tensor == nullptr || state->staged_to_compute_backend || + state->active_prepare_count > 0 || staging_tensor == nullptr) { + LOG_WARN("segment %" PRIu64 " cannot be activated because tensor state changed", segment_id); + free_prefetch_block(*block); + return false; + } + } + for (auto& pair : block->staged_tensors) { + TensorState* state = pair.first; + ggml_tensor* staging_tensor = pair.second; + std::swap(state->tensor->buffer, staging_tensor->buffer); + std::swap(state->tensor->data, staging_tensor->data); + std::swap(state->tensor->extra, staging_tensor->extra); + state->staged_to_compute_backend = true; + } + + auto staging_block = std::make_unique(); + staging_block->compute_backend = block->compute_backend; + staging_block->buffer = block->buffer; + staging_block->staging_ctx = block->staging_ctx; + staging_block->staged_tensors = std::move(block->staged_tensors); + block->buffer = nullptr; + block->staging_ctx = nullptr; + compute_staging_blocks_.push_back(std::move(staging_block)); + return true; +} + +void ModelManager::clear_param_prefetches(uintptr_t owner_id) { + for (auto it = prefetch_blocks_.begin(); it != prefetch_blocks_.end();) { + if (it->first.owner_id == owner_id) { + free_prefetch_block(*it->second); + it = prefetch_blocks_.erase(it); + } else { + ++it; + } + } +} + +size_t ModelManager::streaming_allocation_bytes( + uintptr_t owner_id, + ggml_backend_t compute_backend, + const std::unordered_set& resident_tensors) const { + size_t bytes = 0; + for (const auto& block : compute_staging_blocks_) { + if (block == nullptr || block->buffer == nullptr || + block->compute_backend != compute_backend) { + continue; + } + const bool contains_resident = std::any_of( + block->staged_tensors.begin(), + block->staged_tensors.end(), + [&](const std::pair& pair) { + return pair.first != nullptr && + resident_tensors.count(pair.first->tensor) != 0; + }); + if (contains_resident) { + bytes = saturating_add(bytes, ggml_backend_buffer_get_size(block->buffer)); + } + } + for (const auto& entry : prefetch_blocks_) { + const PrefetchBlock* block = entry.second.get(); + if (entry.first.owner_id != owner_id || block == nullptr || block->buffer == nullptr || + block->compute_backend != compute_backend) { + continue; + } + bytes = saturating_add(bytes, ggml_backend_buffer_get_size(block->buffer)); + } + return bytes; +} + +void ModelManager::clear_all_param_prefetches() { + for (auto& entry : prefetch_blocks_) { + free_prefetch_block(*entry.second); + } + prefetch_blocks_.clear(); +} + ModelManager::~ModelManager() { + clear_all_param_prefetches(); release_all(); + for (auto& entry : prefetch_backends_) { + if (entry.second != nullptr) { + ggml_backend_free(entry.second); + } + } + prefetch_backends_.clear(); } void ModelManager::set_common_ignore_tensors(std::set ignore_tensors) { @@ -255,6 +604,7 @@ bool ModelManager::unregister_param_tensors(const std::string& desc, size_t* reg return true; } + clear_all_param_prefetches(); release_compute_staging_blocks(false); std::vector storage_blocks_to_release; @@ -608,6 +958,7 @@ bool ModelManager::apply_loras_to_params(const std::vector& states } void ModelManager::reset_lora_applied_params() { + clear_all_param_prefetches(); release_compute_staging_blocks(true); release_params_storage_blocks(true); for (auto& state : tensor_states_) { @@ -946,34 +1297,36 @@ void ModelManager::free_compute_staging_block(ComputeStagingBlock& block) { block.staged_tensors.clear(); } -void ModelManager::release_compute_staging_blocks(bool force, - const std::unordered_set* target_states) { +size_t ModelManager::release_compute_staging_blocks(bool force) { + size_t released_bytes = 0; for (auto it = compute_staging_blocks_.begin(); it != compute_staging_blocks_.end();) { ComputeStagingBlock* block = it->get(); bool can_release = force; if (!can_release) { can_release = std::all_of(block->staged_tensors.begin(), block->staged_tensors.end(), - [target_states](const std::pair& pair) { + [](const std::pair& pair) { TensorState* state = pair.first; if (state == nullptr) { return true; } - if (target_states != nullptr && - target_states->find(state) == target_states->end()) { - return false; - } return state->active_prepare_count == 0; }); } if (can_release) { + if (block->buffer != nullptr) { + released_bytes = saturating_add( + released_bytes, + ggml_backend_buffer_get_size(block->buffer)); + } free_compute_staging_block(*block); it = compute_staging_blocks_.erase(it); } else { ++it; } } + return released_bytes; } void ModelManager::free_params_storage_block(ParamsStorageBlock& block) { @@ -1097,6 +1450,7 @@ bool ModelManager::assign_compute_backend(const std::vector& tenso return false; } + bool any_change = false; for (TensorState* state : required_states) { if (state == nullptr || state->tensor == nullptr) { continue; @@ -1121,7 +1475,20 @@ bool ModelManager::assign_compute_backend(const std::vector& tenso return false; } - state->compute_backend = compute_backend; + any_change = true; + } + + if (any_change) { + clear_all_param_prefetches(); + } + for (TensorState* state : required_states) { + if (state == nullptr || state->tensor == nullptr) { + continue; + } + + const bool params_follow_compute = state->params_follow_compute_backend || + state->residency_mode == ResidencyMode::Disk; + state->compute_backend = compute_backend; if (params_follow_compute) { state->params_backend = compute_backend; } @@ -1165,32 +1532,32 @@ bool ModelManager::prepare_params(const std::vector& tensors) { return true; } -void ModelManager::finish_compute_backend_usage(const std::vector& states) { +size_t ModelManager::finish_compute_backend_usage(const std::vector& states) { if (states.empty()) { - return; + return 0; } - std::unordered_set target_states; + std::unordered_set unique_states; for (TensorState* state : states) { - if (state == nullptr || !target_states.insert(state).second) { + if (state == nullptr || !unique_states.insert(state).second) { continue; } if (state->active_prepare_count > 0) { state->active_prepare_count--; } } - release_compute_staging_blocks(false, &target_states); + return release_compute_staging_blocks(false); } -void ModelManager::release_compute_backend_params(const std::vector& tensors) { +size_t ModelManager::release_compute_backend_params(const std::vector& tensors) { if (tensors.empty()) { - return; + return 0; } std::vector required_states; if (!resolve_required_tensor_states(tensors, required_states)) { - return; + return 0; } - finish_compute_backend_usage(required_states); + return finish_compute_backend_usage(required_states); } void ModelManager::release_params_backend_params(const std::vector& tensors) { diff --git a/src/model_manager.h b/src/model_manager.h index 3dc633864..24d8ea2e0 100644 --- a/src/model_manager.h +++ b/src/model_manager.h @@ -61,12 +61,35 @@ class ModelManager : public RunnerWeightManager { std::vector> staged_tensors; }; + struct PrefetchKey { + uintptr_t owner_id = 0; + uint64_t segment_id = 0; + + bool operator<(const PrefetchKey& other) const { + return owner_id < other.owner_id || + (owner_id == other.owner_id && segment_id < other.segment_id); + } + }; + + struct PrefetchBlock { + PrefetchKey key; + std::vector states; + ggml_backend_t compute_backend = nullptr; + ggml_backend_t transfer_backend = nullptr; + ggml_backend_event_t event = nullptr; + ggml_context* staging_ctx = nullptr; + ggml_backend_buffer_t buffer = nullptr; + std::vector> staged_tensors; + }; + ModelLoader model_loader_; std::vector> tensor_states_; std::map tensor_states_by_name_; std::vector> params_storage_blocks_; std::vector> compute_staging_blocks_; std::map split_buffer_types_; + std::map> prefetch_blocks_; + std::map prefetch_backends_; bool warned_split_lora_skip_ = false; std::set common_ignore_tensors_; std::vector loras_; @@ -76,9 +99,15 @@ class ModelManager : public RunnerWeightManager { bool enable_mmap_ = false; bool writable_mmap_ = false; - void finish_compute_backend_usage(const std::vector& states); + size_t finish_compute_backend_usage(const std::vector& states); void release_all(); + ParamPrefetchResult populate_prefetch_block(PrefetchBlock& block); + ggml_backend_t prefetch_backend_for(ggml_backend_t compute_backend); + void synchronize_prefetch_block(PrefetchBlock& block); + void free_prefetch_block(PrefetchBlock& block); + void clear_all_param_prefetches(); + bool resolve_required_tensor_states(const std::vector& tensors, std::vector& required_states) const; bool should_ignore(const TensorState& state) const; @@ -97,8 +126,7 @@ class ModelManager : public RunnerWeightManager { ggml_backend_buffer_type_t params_buffer_type_for(const TensorState& state) const; ggml_backend_buffer_type_t split_buffer_type_for(const TensorState& state) const; - void release_compute_staging_blocks(bool force = false, - const std::unordered_set* target_states = nullptr); + size_t release_compute_staging_blocks(bool force = false); void release_params_storage_blocks(bool force = false, const std::unordered_set* target_states = nullptr); void free_compute_staging_block(ComputeStagingBlock& block); @@ -180,8 +208,20 @@ class ModelManager : public RunnerWeightManager { bool assign_compute_backend(const std::vector& tensors, ggml_backend_t compute_backend) override; bool prepare_params(const std::vector& tensors) override; - void release_compute_backend_params(const std::vector& tensors) override; + size_t release_compute_backend_params(const std::vector& tensors) override; void release_params_backend_params(const std::vector& tensors) override; + ParamPrefetchResult enqueue_param_prefetch( + uintptr_t owner_id, + uint64_t segment_id, + const std::vector& tensors) override; + bool activate_param_prefetch(uintptr_t owner_id, + uint64_t segment_id, + const std::vector& tensors) override; + void clear_param_prefetches(uintptr_t owner_id) override; + size_t streaming_allocation_bytes( + uintptr_t owner_id, + ggml_backend_t compute_backend, + const std::unordered_set& resident_tensors) const override; }; #endif // __MODEL_MANAGER_H__ diff --git a/src/stable-diffusion.cpp b/src/stable-diffusion.cpp index 109a2a483..c1a144ce1 100644 --- a/src/stable-diffusion.cpp +++ b/src/stable-diffusion.cpp @@ -247,8 +247,10 @@ class StableDiffusionGGML { sd_tiling_params_t vae_tiling_params = {false, false, 0, 0, 0.5f, 0, 0, nullptr}; bool enable_mmap = false; sd::ggml_graph_cut::MaxVramAssignment max_vram_assignment; - bool stream_layers = false; - bool eager_load = false; + bool stream_layers = false; + int resident_layers = -1; + int layer_prefetch_depth = 0; + bool eager_load = false; std::string backend_spec; std::string params_backend_spec; std::string split_mode_spec; @@ -857,10 +859,35 @@ class StableDiffusionGGML { return true; } - bool init(const sd_ctx_params_t* sd_ctx_params) { - n_threads = sd_ctx_params->n_threads; - enable_mmap = sd_ctx_params->enable_mmap; - stream_layers = sd_ctx_params->stream_layers; + bool init(const sd_ctx_params_t* sd_ctx_params, + const sd_layer_stream_params_t* layer_stream_params) { + n_threads = sd_ctx_params->n_threads; + enable_mmap = sd_ctx_params->enable_mmap; + stream_layers = sd_ctx_params->stream_layers; + resident_layers = -1; + layer_prefetch_depth = 0; + if (layer_stream_params != nullptr) { + constexpr size_t min_struct_size = offsetof(sd_layer_stream_params_t, + layer_prefetch_depth) + + sizeof(layer_stream_params->layer_prefetch_depth); + if (layer_stream_params->struct_size < min_struct_size) { + LOG_ERROR("layer stream params struct is too small: %u < %zu", + layer_stream_params->struct_size, + min_struct_size); + return false; + } + resident_layers = layer_stream_params->resident_layers; + layer_prefetch_depth = layer_stream_params->layer_prefetch_depth; + } + if (resident_layers < -1 || layer_prefetch_depth < 0) { + LOG_ERROR("layer stream limits require resident_layers >= -1 and layer_prefetch_depth >= 0"); + return false; + } + if (!stream_layers && (resident_layers >= 0 || layer_prefetch_depth > 0)) { + LOG_WARN("resident_layers and layer_prefetch_depth require stream_layers; ignoring layer stream limits"); + resident_layers = -1; + layer_prefetch_depth = 0; + } eager_load = sd_ctx_params->eager_load; backend_spec = SAFE_STR(sd_ctx_params->backend); params_backend_spec = SAFE_STR(sd_ctx_params->params_backend); @@ -932,6 +959,17 @@ class StableDiffusionGGML { LOG_WARN("--stream-layers has no effect unless diffusion params backend is cpu; ignoring"); stream_layers = false; } + if (stream_layers && + backend_manager.runtime_backends(SDBackendModule::DIFFUSION).size() > 1) { + LOG_WARN("--stream-layers is not supported when the diffusion model uses multiple runtime backends; ignoring"); + stream_layers = false; + } + if (stream_layers && + max_graph_vram_bytes_for_module(SDBackendModule::DIFFUSION) == 0) { + LOG_WARN("--stream-layers has no effect because diffusion --max-vram is 0; " + "residency and prefetch controls are ignored"); + stream_layers = false; + } if (eager_load && graph_cut_layer_split_active()) { LOG_WARN("--eager-load is not supported with graph-cut layer split; weights will be prepared lazily"); eager_load = false; @@ -1354,23 +1392,25 @@ class StableDiffusionGGML { } diffusion_model->set_max_graph_vram_bytes(max_graph_vram_bytes_for_module(SDBackendModule::DIFFUSION)); - diffusion_model->set_stream_layers_enabled(stream_layers); if (!register_runner_params("Diffusion model", diffusion_model, SDBackendModule::DIFFUSION, &unet_params_mem_size)) { return false; } + diffusion_model->set_stream_segment_limits(resident_layers, layer_prefetch_depth); + diffusion_model->set_stream_layers_enabled(stream_layers); if (high_noise_diffusion_model) { high_noise_diffusion_model->set_max_graph_vram_bytes(max_graph_vram_bytes_for_module(SDBackendModule::DIFFUSION)); - high_noise_diffusion_model->set_stream_layers_enabled(stream_layers); if (!register_runner_params("High noise diffusion model", high_noise_diffusion_model, SDBackendModule::DIFFUSION, &unet_params_mem_size)) { return false; } + high_noise_diffusion_model->set_stream_segment_limits(resident_layers, layer_prefetch_depth); + high_noise_diffusion_model->set_stream_layers_enabled(stream_layers); } if (strlen(SAFE_STR(sd_ctx_params->ip_adapter_path)) > 0 && clip_vision == nullptr) { @@ -2540,6 +2580,8 @@ class StableDiffusionGGML { } }; RunnerDoneOnExit sample_diffusion_runner_done{work_diffusion_model.get()}; + // Residency only pays off across the repeated denoising forwards. + work_diffusion_model->set_stream_residency_enabled(true); RunnerDoneOnExit sample_control_runner_done{!control_image.empty() && control_net != nullptr ? control_net.get() : nullptr}; @@ -2628,6 +2670,9 @@ class StableDiffusionGGML { return {}; } + work_diffusion_model->set_stream_next_forward_prefetch( + static_cast(std::abs(step)) < steps); + if (step == 1 || step == -1) { pretty_progress(0, (int)steps, 0); last_progress_us = ggml_time_us(); @@ -3568,6 +3613,13 @@ void sd_ctx_params_init(sd_ctx_params_t* sd_ctx_params) { sd_ctx_params->pulid_weights_path = nullptr; } +void sd_layer_stream_params_init(sd_layer_stream_params_t* params) { + *params = {}; + params->struct_size = sizeof(*params); + params->resident_layers = -1; + params->layer_prefetch_depth = 0; +} + char* sd_ctx_params_to_str(const sd_ctx_params_t* sd_ctx_params) { char* buf = (char*)malloc(8192); if (!buf) @@ -3845,6 +3897,11 @@ static bool sd_version_supports_image_generation(SDVersion version) { } sd_ctx_t* new_sd_ctx(const sd_ctx_params_t* sd_ctx_params) { + return new_sd_ctx_with_layer_stream(sd_ctx_params, nullptr); +} + +sd_ctx_t* new_sd_ctx_with_layer_stream(const sd_ctx_params_t* sd_ctx_params, + const sd_layer_stream_params_t* layer_stream_params) { sd_ctx_t* sd_ctx = (sd_ctx_t*)malloc(sizeof(sd_ctx_t)); if (sd_ctx == nullptr) { return nullptr; @@ -3856,7 +3913,7 @@ sd_ctx_t* new_sd_ctx(const sd_ctx_params_t* sd_ctx_params) { return nullptr; } - if (!sd_ctx->sd->init(sd_ctx_params)) { + if (!sd_ctx->sd->init(sd_ctx_params, layer_stream_params)) { delete sd_ctx->sd; sd_ctx->sd = nullptr; free(sd_ctx); diff --git a/src/weight_manager.h b/src/weight_manager.h index 82f6d03e4..3f80794a2 100644 --- a/src/weight_manager.h +++ b/src/weight_manager.h @@ -1,19 +1,41 @@ #ifndef __WEIGHT_MANAGER_H__ #define __WEIGHT_MANAGER_H__ +#include +#include +#include #include #include "ggml-backend.h" struct ggml_tensor; +enum class ParamPrefetchResult : uint8_t { + SUCCESS = 0, + ALLOCATION_FAILURE, + FAILURE, +}; + struct RunnerWeightManager { virtual ~RunnerWeightManager() = default; virtual bool assign_compute_backend(const std::vector& tensors, ggml_backend_t compute_backend) = 0; virtual bool prepare_params(const std::vector& tensors) = 0; - virtual void release_compute_backend_params(const std::vector& tensors) = 0; + virtual size_t release_compute_backend_params(const std::vector& tensors) = 0; virtual void release_params_backend_params(const std::vector& tensors) = 0; + + virtual ParamPrefetchResult enqueue_param_prefetch( + uintptr_t owner_id, + uint64_t segment_id, + const std::vector& tensors) = 0; + virtual bool activate_param_prefetch(uintptr_t owner_id, + uint64_t segment_id, + const std::vector& tensors) = 0; + virtual void clear_param_prefetches(uintptr_t owner_id) = 0; + virtual size_t streaming_allocation_bytes( + uintptr_t owner_id, + ggml_backend_t compute_backend, + const std::unordered_set& resident_tensors) const = 0; }; #endif // __WEIGHT_MANAGER_H__ From 469af44a2a64ae7fac559ce3dcf7d3bc06c784e8 Mon Sep 17 00:00:00 2001 From: assouan <750048+assouan@users.noreply.github.com> Date: Sat, 22 Aug 2026 20:53:05 +0200 Subject: [PATCH 2/2] refactor: reuse canonical layer stream metadata Reworks stream-layer prefetch bookkeeping to use graph-cut segment parameter allocations directly, with a new shared segment-parameter map used by prefetching, residency retention, and eviction. Prefetch scheduling now explicitly targets parameter-bearing segments and keeps async transfers limited to non-shared cross-segment params. The docs, CLI option text, and public header comments were updated to match this behavior. --- docs/performance.md | 2 +- examples/common/common.cpp | 2 +- examples/common/common.h | 18 +-- include/stable-diffusion.h | 2 +- src/core/ggml_extend.hpp | 265 +++++++++++++++++------------------- src/core/ggml_graph_cut.cpp | 4 +- src/stable-diffusion.cpp | 5 +- src/weight_manager.h | 10 +- 8 files changed, 145 insertions(+), 163 deletions(-) diff --git a/docs/performance.md b/docs/performance.md index 61985038a..07742bc7b 100644 --- a/docs/performance.md +++ b/docs/performance.md @@ -60,7 +60,7 @@ See [backend selection](./backend.md) for full syntax. - `--max-vram ` sets a VRAM budget the graph-cut segmenter respects. It cuts each forward pass into segments sized to fit the budget, running them in sequence and freeing intermediate activations between them. Negative values auto-detect free VRAM and spare the given amount (`--max-vram -1` uses most of the free VRAM and keeps ~1 GiB headroom), a positive value caps the budget, `0` disables segmentation. - `--stream-layers` streams the diffusion model's transformer blocks one at a time. Each block's parameters are copied from the CPU to the runtime backend just before it runs and evicted when the residency budget is reached. Prefetching hides most of the copy latency behind compute. This flag only takes effect when the diffusion params backend is CPU, so it must be combined with `--offload-to-cpu` (or an explicit `--params-backend diffusion=cpu`); a warning is logged and the flag is ignored otherwise. - `--resident-layers ` sets the maximum number of leading parameter-bearing graph-cut segments kept resident (default: `-1`; `auto` is an alias for `-1`). `-1` uses as many as the VRAM budget permits, `0` keeps none, and a positive `N` keeps up to `N`. -- `--layer-prefetch-depth ` sets the maximum number of future graph-cut segments copied through a separate transfer backend or queue when supported while the active segment computes (default: `0`). `0` disables asynchronous prefetching, `1` overlaps the next segment, and larger values provide deeper lookahead when the VRAM budget permits. +- `--layer-prefetch-depth ` sets the maximum number of future parameter-bearing graph-cut segments copied through a separate transfer backend or queue when supported while the active segment computes (default: `0`). `0` disables asynchronous prefetching, `1` overlaps the next segment, and larger values provide deeper lookahead when the VRAM budget permits. Both controls require `--stream-layers` and a non-zero `--max-vram`. Prefetch is budgeted before residency; requested limits are reduced as needed, a positive VRAM cap is never exceeded to force either one, and the active segment is never evicted. Residents persist only across repeated sampling steps, and stale graph state is released automatically. An explicit non-negative residency limit uses an unmerged plan, which may add dispatch overhead. diff --git a/examples/common/common.cpp b/examples/common/common.cpp index a1f8ce0a6..1dc6b63b1 100644 --- a/examples/common/common.cpp +++ b/examples/common/common.cpp @@ -520,7 +520,7 @@ ArgOptions SDContextParams::get_options() { &n_threads}, {"", "--layer-prefetch-depth", - "number of future graph-cut segments to prefetch with --stream-layers (default: 0; 0 disables prefetching, 1 overlaps the next segment, 2+ enables deeper lookahead when VRAM permits)", + "number of future parameter-bearing graph-cut segments to prefetch with --stream-layers (default: 0; 0 disables prefetching, 1 overlaps the next segment, 2+ enables deeper lookahead when VRAM permits)", &layer_prefetch_depth}, }; diff --git a/examples/common/common.h b/examples/common/common.h index 01f259ce0..9f4ef93a6 100644 --- a/examples/common/common.h +++ b/examples/common/common.h @@ -146,15 +146,15 @@ struct SDContextParams { std::map embedding_map; std::vector embedding_vec; - rng_type_t rng_type = CUDA_RNG; - rng_type_t sampler_rng_type = RNG_TYPE_COUNT; - bool offload_params_to_cpu = false; - std::string max_vram = "0"; - bool stream_layers = false; - std::string resident_layers_spec = "-1"; - int resident_layers = -1; - int layer_prefetch_depth = 0; - bool eager_load = false; + rng_type_t rng_type = CUDA_RNG; + rng_type_t sampler_rng_type = RNG_TYPE_COUNT; + bool offload_params_to_cpu = false; + std::string max_vram = "0"; + bool stream_layers = false; + std::string resident_layers_spec = "-1"; + int resident_layers = -1; + int layer_prefetch_depth = 0; + bool eager_load = false; std::string backend; std::string params_backend; std::string split_mode; diff --git a/include/stable-diffusion.h b/include/stable-diffusion.h index 8d4d00172..9f739c035 100644 --- a/include/stable-diffusion.h +++ b/include/stable-diffusion.h @@ -240,7 +240,7 @@ typedef struct { typedef struct { uint32_t struct_size; // Set by sd_layer_stream_params_init; permits future extension int resident_layers; // With stream_layers: maximum leading graph-cut segments kept resident (-1 = automatic, 0 = none) - int layer_prefetch_depth; // With stream_layers: future graph-cut segments prefetched during compute (0 = disabled) + int layer_prefetch_depth; // With stream_layers: future parameter-bearing graph-cut segments prefetched during compute (0 = disabled) } sd_layer_stream_params_t; typedef struct { diff --git a/src/core/ggml_extend.hpp b/src/core/ggml_extend.hpp index d364c1dd2..ebe6d21a2 100644 --- a/src/core/ggml_extend.hpp +++ b/src/core/ggml_extend.hpp @@ -1795,7 +1795,6 @@ struct GGMLRunner { size_t runtime_resident_segment_cap = SIZE_MAX; bool stream_next_forward_prefetch = false; bool stream_policy_logged = false; - bool stream_prefetch_fallback_warned = false; bool stream_prefetch_runtime_warned = false; bool stream_prefetch_reduced_warned = false; bool stream_shared_params_logged = false; @@ -2488,7 +2487,7 @@ struct GGMLRunner { // this runner can reuse or release them. Other allocations // remain charged through free_vram. constexpr size_t safety_margin = 512ull * 1024 * 1024; - size_t reclaimable_free = saturating_size_add(free_vram, runner_owned_vram); + size_t reclaimable_free = saturating_size_add(free_vram, runner_owned_vram); if (total_vram > 0) { reclaimable_free = std::min(reclaimable_free, total_vram); } @@ -2734,99 +2733,99 @@ struct GGMLRunner { return static_cast(segment_index + 1); } - bool collect_prefetch_segment_params( - ggml_cgraph* gf, - const GraphCutPlan& plan, - std::vector>& segment_params, - std::vector& param_segments, - std::vector& segment_to_param_position) { - segment_params.assign(plan.segments.size(), {}); - param_segments.clear(); - segment_to_param_position.assign(plan.segments.size(), SIZE_MAX); + struct SegmentStreamParams { + std::vector> all_by_segment; + std::vector> prefetch_by_segment; + std::vector param_segments; + std::vector segment_to_param_position; + bool has_prefetch_params = false; + }; + + SegmentStreamParams collect_stream_segment_params(const GraphCutPlan& plan, + bool include_prefetch_params) { + SegmentStreamParams stream_params; + stream_params.all_by_segment.resize(plan.segments.size()); + stream_params.prefetch_by_segment.resize(plan.segments.size()); + stream_params.segment_to_param_position.assign(plan.segments.size(), SIZE_MAX); std::unordered_map segment_occurrences; for (size_t segment_index = 0; segment_index < plan.segments.size(); ++segment_index) { - std::unordered_set seen_in_segment; - for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors(gf, - plan.segments[segment_index])) { - ggml_tensor* canonical = canonical_param_tensor(tensor); - if (canonical == nullptr) { - if (!stream_prefetch_fallback_warned) { - LOG_WARN("%s cannot canonicalize a streamed parameter; disabling segment prefetch", - get_desc().c_str()); - stream_prefetch_fallback_warned = true; - } - return false; - } - if (!seen_in_segment.insert(canonical).second) { + const GraphCutSegment& segment = plan.segments[segment_index]; + auto& params = stream_params.all_by_segment[segment_index]; + params.reserve(segment.input_param_allocations.size()); + for (const GraphCutSegment::ParamAllocation& allocation : + segment.input_param_allocations) { + if (allocation.tensor == nullptr) { continue; } - segment_params[segment_index].push_back(canonical); - ++segment_occurrences[canonical]; + ggml_tensor* tensor = const_cast(allocation.tensor); + params.push_back(tensor); + if (include_prefetch_params) { + ++segment_occurrences[tensor]; + } } - } - - size_t shared_params = 0; - for (const auto& occurrence : segment_occurrences) { - if (occurrence.second > 1) { - ++shared_params; + if (segment.input_param_bytes > 0) { + stream_params.segment_to_param_position[segment_index] = + stream_params.param_segments.size(); + stream_params.param_segments.push_back(segment_index); } } - if (shared_params > 0) { - for (auto& params : segment_params) { - params.erase(std::remove_if(params.begin(), params.end(), [&](ggml_tensor* param) { - return segment_occurrences.at(param) > 1; - }), - params.end()); - } - if (!stream_shared_params_logged) { - LOG_INFO("%s segment prefetch excludes %zu cross-segment parameters from asynchronous transfers", - get_desc().c_str(), - shared_params); - stream_shared_params_logged = true; - } + + if (!include_prefetch_params) { + return stream_params; } - bool any_prefetch_params = false; - for (size_t segment_index = 0; segment_index < plan.segments.size(); ++segment_index) { - if (plan.segments[segment_index].input_param_bytes > 0) { - segment_to_param_position[segment_index] = param_segments.size(); - param_segments.push_back(segment_index); - any_prefetch_params = any_prefetch_params || !segment_params[segment_index].empty(); + const size_t shared_params = static_cast(std::count_if( + segment_occurrences.begin(), + segment_occurrences.end(), + [](const auto& occurrence) { return occurrence.second > 1; })); + for (size_t segment_index : stream_params.param_segments) { + auto& prefetch_params = stream_params.prefetch_by_segment[segment_index]; + for (ggml_tensor* tensor : stream_params.all_by_segment[segment_index]) { + if (segment_occurrences.at(tensor) == 1) { + prefetch_params.push_back(tensor); + } } + stream_params.has_prefetch_params = + stream_params.has_prefetch_params || !prefetch_params.empty(); + } + if (shared_params > 0 && !stream_shared_params_logged) { + LOG_INFO("%s segment prefetch excludes %zu cross-segment parameters from asynchronous transfers", + get_desc().c_str(), + shared_params); + stream_shared_params_logged = true; } - return any_prefetch_params; + return stream_params; } std::unordered_set retained_stream_params( const GraphCutPlan& plan, - ggml_cgraph* gf, + const SegmentStreamParams& stream_params, size_t protected_segment_index = SIZE_MAX) { std::unordered_set retained_params; - auto collect_segment = [&](const GraphCutSegment& segment) { - for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors(gf, segment)) { - if (ggml_tensor* canonical = canonical_param_tensor(tensor)) { - retained_params.insert(canonical); - } + auto collect_segment = [&](size_t segment_index) { + for (ggml_tensor* tensor : stream_params.all_by_segment[segment_index]) { + retained_params.insert(tensor); } }; - for (const GraphCutSegment& segment : plan.segments) { - if (segment.residency == sd::ggml_graph_cut::SegmentResidency::RESIDENT) { - collect_segment(segment); + for (size_t segment_index = 0; segment_index < plan.segments.size(); ++segment_index) { + if (plan.segments[segment_index].residency == + sd::ggml_graph_cut::SegmentResidency::RESIDENT) { + collect_segment(segment_index); } } if (protected_segment_index < plan.segments.size()) { - collect_segment(plan.segments[protected_segment_index]); + collect_segment(protected_segment_index); } return retained_params; } size_t release_unretained_resident_params(const GraphCutPlan& plan, - ggml_cgraph* gf, + const SegmentStreamParams& stream_params, size_t protected_segment_index = SIZE_MAX) { const std::unordered_set retained_params = - retained_stream_params(plan, gf, protected_segment_index); + retained_stream_params(plan, stream_params, protected_segment_index); std::vector params_to_release; params_to_release.reserve(kept_compute_param_tensor_set.size()); for (const ggml_tensor* tensor : kept_compute_param_tensor_set) { @@ -2843,8 +2842,8 @@ struct GGMLRunner { } size_t evict_resident_segment_for_prefetch(GraphCutPlan& plan, - ggml_cgraph* gf, size_t active_segment_index, + const SegmentStreamParams& stream_params, bool& residency_changed) { bool demoted = false; std::vector resident_segments; @@ -2863,27 +2862,22 @@ struct GGMLRunner { std::unordered_set earlier_resident_params; auto collect_earlier_resident_params = [&](size_t segment_index) { - for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors( - gf, plan.segments[segment_index])) { - if (ggml_tensor* canonical = canonical_param_tensor(tensor)) { - earlier_resident_params.insert(canonical); - } + for (ggml_tensor* tensor : stream_params.all_by_segment[segment_index]) { + earlier_resident_params.insert(tensor); } }; for (size_t position = 0; position < resident_position; ++position) { collect_earlier_resident_params(resident_segments[position]); } - const std::vector candidate_params = - sd::ggml_graph_cut::param_tensors(gf, plan.segments[candidate_index]); + const std::vector& candidate_params = + stream_params.all_by_segment[candidate_index]; const bool has_demotable_params = std::any_of( candidate_params.begin(), candidate_params.end(), [&](ggml_tensor* tensor) { - ggml_tensor* canonical = canonical_param_tensor(tensor); - return canonical != nullptr && - kept_compute_param_tensor_set.count(canonical) != 0 && - earlier_resident_params.count(canonical) == 0; + return kept_compute_param_tensor_set.count(tensor) != 0 && + earlier_resident_params.count(tensor) == 0; }); if (!has_demotable_params) { continue; @@ -2899,7 +2893,7 @@ struct GGMLRunner { residency_changed = true; demoted = true; const size_t released_bytes = release_unretained_resident_params( - plan, gf, active_segment_index); + plan, stream_params, active_segment_index); if (released_bytes > 0) { LOG_WARN( "%s evicted resident segment %zu to make room for prefetch; " @@ -2921,10 +2915,9 @@ struct GGMLRunner { return 0; } - ParamPrefetchResult enqueue_segment_prefetch( - size_t segment_index, - const std::vector>& segment_params) { - if (segment_index >= segment_params.size()) { + ParamPrefetchResult enqueue_segment_prefetch(size_t segment_index, + const SegmentStreamParams& stream_params) { + if (segment_index >= stream_params.prefetch_by_segment.size()) { return ParamPrefetchResult::FAILURE; } auto manager = weight_manager.lock(); @@ -2934,31 +2927,30 @@ struct GGMLRunner { } return manager->enqueue_param_prefetch(stream_prefetch_owner_id(), segment_prefetch_id(segment_index), - segment_params[segment_index]); + stream_params.prefetch_by_segment[segment_index]); } ParamPrefetchResult enqueue_segment_prefetch_with_eviction( size_t segment_index, size_t active_segment_index, - const std::vector>& segment_params, + const SegmentStreamParams& stream_params, GraphCutPlan& plan, - ggml_cgraph* gf, bool& residency_changed) { - ParamPrefetchResult result = enqueue_segment_prefetch(segment_index, segment_params); + ParamPrefetchResult result = enqueue_segment_prefetch(segment_index, stream_params); while (result == ParamPrefetchResult::ALLOCATION_FAILURE) { const size_t released_bytes = evict_resident_segment_for_prefetch( - plan, gf, active_segment_index, residency_changed); + plan, active_segment_index, stream_params, residency_changed); if (released_bytes == 0) { break; } - result = enqueue_segment_prefetch(segment_index, segment_params); + result = enqueue_segment_prefetch(segment_index, stream_params); } return result; } bool activate_segment_prefetch(size_t segment_index, - const std::vector>& segment_params) { - if (segment_index >= segment_params.size()) { + const SegmentStreamParams& stream_params) { + if (segment_index >= stream_params.prefetch_by_segment.size()) { return false; } auto manager = weight_manager.lock(); @@ -2968,7 +2960,7 @@ struct GGMLRunner { } return manager->activate_param_prefetch(stream_prefetch_owner_id(), segment_prefetch_id(segment_index), - segment_params[segment_index]); + stream_params.prefetch_by_segment[segment_index]); } struct SegmentPrefetchQueueResult { @@ -2978,32 +2970,30 @@ struct GGMLRunner { SegmentPrefetchQueueResult queue_future_segment_prefetches( size_t active_position, - const std::vector& param_segments, - const std::vector>& segment_params, + const SegmentStreamParams& stream_params, size_t depth, bool wrap, GraphCutPlan& plan, - ggml_cgraph* gf, bool& residency_changed) { SegmentPrefetchQueueResult queue_result; - if (param_segments.empty() || active_position >= param_segments.size()) { + if (stream_params.param_segments.empty() || + active_position >= stream_params.param_segments.size()) { return queue_result; } - depth = std::min(depth, param_segments.size() - 1); + depth = std::min(depth, stream_params.param_segments.size() - 1); for (size_t offset = 1; offset <= depth; ++offset) { size_t position = active_position + offset; - if (position >= param_segments.size()) { + if (position >= stream_params.param_segments.size()) { if (!wrap) { break; } - position %= param_segments.size(); + position %= stream_params.param_segments.size(); } queue_result.result = enqueue_segment_prefetch_with_eviction( - param_segments[position], - param_segments[active_position], - segment_params, + stream_params.param_segments[position], + stream_params.param_segments[active_position], + stream_params, plan, - gf, residency_changed); if (queue_result.result != ParamPrefetchResult::SUCCESS) { return queue_result; @@ -3312,17 +3302,22 @@ struct GGMLRunner { free_compute_buffer(); free_cache_ctx_and_buffer(); - auto manager = weight_manager.lock(); + auto manager = weight_manager.lock(); + const bool prefetch_requested = prefetch_policy != nullptr && + prefetch_policy->prefetch_depth > 0; + SegmentStreamParams stream_params; + if (stream_layers_enabled) { + stream_params = collect_stream_segment_params(plan, prefetch_requested); + } if (stream_layers_enabled && !kept_compute_param_tensor_set.empty()) { std::unordered_set desired_resident_params; - for (const auto& segment : plan.segments) { - if (segment.residency != sd::ggml_graph_cut::SegmentResidency::RESIDENT) { + for (size_t segment_index = 0; segment_index < plan.segments.size(); ++segment_index) { + if (plan.segments[segment_index].residency != + sd::ggml_graph_cut::SegmentResidency::RESIDENT) { continue; } - for (ggml_tensor* tensor : sd::ggml_graph_cut::param_tensors(gf, segment)) { - if (ggml_tensor* canonical = canonical_param_tensor(tensor)) { - desired_resident_params.insert(canonical); - } + for (ggml_tensor* tensor : stream_params.all_by_segment[segment_index]) { + desired_resident_params.insert(tensor); } } const bool same_resident_params = @@ -3348,21 +3343,13 @@ struct GGMLRunner { } } } - std::vector> segment_params; - std::vector param_segments; - std::vector segment_to_param_position; - bool prefetch_active = prefetch_policy != nullptr && - prefetch_policy->prefetch_depth > 0 && - collect_prefetch_segment_params(gf, - plan, - segment_params, - param_segments, - segment_to_param_position); + bool prefetch_active = prefetch_requested && stream_params.has_prefetch_params; size_t runtime_prefetch_depth = prefetch_active ? prefetch_policy->prefetch_depth : 0; if (prefetch_active) { - update_stream_prefetch_plan_signature(segment_params, runtime_prefetch_depth); + update_stream_prefetch_plan_signature(stream_params.prefetch_by_segment, + runtime_prefetch_depth); } else { clear_stream_prefetch_state(); } @@ -3393,11 +3380,10 @@ struct GGMLRunner { if (prefetch_active) { ParamPrefetchResult result = enqueue_segment_prefetch_with_eviction( - param_segments.front(), - param_segments.front(), - segment_params, + stream_params.param_segments.front(), + stream_params.param_segments.front(), + stream_params, plan, - gf, residency_changed); if (result != ParamPrefetchResult::SUCCESS) { disable_prefetch("queueing the first segment"); @@ -3412,22 +3398,21 @@ struct GGMLRunner { const auto& segment = plan.segments[seg_idx]; const bool is_last = seg_idx + 1 == plan.segments.size(); size_t param_position = SIZE_MAX; - if (!segment_to_param_position.empty()) { - param_position = segment_to_param_position[seg_idx]; + if (!stream_params.segment_to_param_position.empty()) { + param_position = stream_params.segment_to_param_position[seg_idx]; } - const bool has_prefetch_params = param_position != SIZE_MAX; - auto future_cut_names = sd::ggml_graph_cut::collect_future_input_names(gf, plan, seg_idx); + const bool is_param_segment = param_position != SIZE_MAX; + auto future_cut_names = sd::ggml_graph_cut::collect_future_input_names(gf, plan, seg_idx); - if (has_prefetch_params && prefetch_active) { + if (is_param_segment && prefetch_active) { ParamPrefetchResult result = enqueue_segment_prefetch_with_eviction( seg_idx, seg_idx, - segment_params, + stream_params, plan, - gf, residency_changed); if (result != ParamPrefetchResult::SUCCESS || - !activate_segment_prefetch(seg_idx, segment_params)) { + !activate_segment_prefetch(seg_idx, stream_params)) { disable_prefetch("activating a segment"); } } @@ -3482,13 +3467,11 @@ struct GGMLRunner { ggml_cgraph* segment_graph = sd::ggml_graph_cut::build_segment_graph(gf, segment, &segment_graph_ctx); const bool keep_segment_params = segment.residency == sd::ggml_graph_cut::SegmentResidency::RESIDENT; std::function before_compute; - if (has_prefetch_params && prefetch_active) { + if (is_param_segment && prefetch_active) { before_compute = [this, param_position, - ¶m_segments, - &segment_params, + &stream_params, &plan, - gf, &runtime_prefetch_depth, &prefetch_active, &residency_changed, @@ -3498,12 +3481,10 @@ struct GGMLRunner { } SegmentPrefetchQueueResult result = queue_future_segment_prefetches( param_position, - param_segments, - segment_params, + stream_params, runtime_prefetch_depth, stream_next_forward_prefetch, plan, - gf, residency_changed); if (result.result == ParamPrefetchResult::SUCCESS) { return; @@ -3532,7 +3513,7 @@ struct GGMLRunner { before_compute); ggml_free(segment_graph_ctx); if (residency_changed) { - release_unretained_resident_params(plan, gf); + release_unretained_resident_params(plan, stream_params); residency_changed = false; } if (!segment_output.has_value()) { @@ -3840,9 +3821,9 @@ struct GGMLRunner { if (stream_residency_enabled == enabled) { return; } - stream_residency_enabled = enabled; - runtime_resident_segment_cap = SIZE_MAX; - stream_policy_logged = false; + stream_residency_enabled = enabled; + runtime_resident_segment_cap = SIZE_MAX; + stream_policy_logged = false; } void set_stream_next_forward_prefetch(bool enabled) { diff --git a/src/core/ggml_graph_cut.cpp b/src/core/ggml_graph_cut.cpp index 0d5e0de7c..0e117a1fb 100644 --- a/src/core/ggml_graph_cut.cpp +++ b/src/core/ggml_graph_cut.cpp @@ -478,7 +478,7 @@ namespace sd::ggml_graph_cut { if (canonical != nullptr && param_seen.insert(canonical).second) { const size_t allocation_bytes = tensor_backend_allocation_bytes(backend, canonical); - segment.input_param_bytes = saturating_add( + segment.input_param_bytes = saturating_add( segment.input_param_bytes, allocation_bytes); segment.input_param_allocations.push_back({canonical, allocation_bytes}); @@ -1086,7 +1086,7 @@ namespace sd::ggml_graph_cut { const std::unordered_map& param_occurrences, size_t resident_count, size_t prefetch_depth) { - size_t peak = 0; + size_t peak = 0; prefetch_depth = param_segments.size() < 2 ? 0 : std::min(prefetch_depth, param_segments.size() - 1); diff --git a/src/stable-diffusion.cpp b/src/stable-diffusion.cpp index c1a144ce1..2205202ad 100644 --- a/src/stable-diffusion.cpp +++ b/src/stable-diffusion.cpp @@ -966,8 +966,9 @@ class StableDiffusionGGML { } if (stream_layers && max_graph_vram_bytes_for_module(SDBackendModule::DIFFUSION) == 0) { - LOG_WARN("--stream-layers has no effect because diffusion --max-vram is 0; " - "residency and prefetch controls are ignored"); + LOG_WARN( + "--stream-layers has no effect because diffusion --max-vram is 0; " + "residency and prefetch controls are ignored"); stream_layers = false; } if (eager_load && graph_cut_layer_split_active()) { diff --git a/src/weight_manager.h b/src/weight_manager.h index 3f80794a2..7300336c4 100644 --- a/src/weight_manager.h +++ b/src/weight_manager.h @@ -17,17 +17,17 @@ enum class ParamPrefetchResult : uint8_t { }; struct RunnerWeightManager { - virtual ~RunnerWeightManager() = default; + virtual ~RunnerWeightManager() = default; virtual bool assign_compute_backend(const std::vector& tensors, - ggml_backend_t compute_backend) = 0; - virtual bool prepare_params(const std::vector& tensors) = 0; + ggml_backend_t compute_backend) = 0; + virtual bool prepare_params(const std::vector& tensors) = 0; virtual size_t release_compute_backend_params(const std::vector& tensors) = 0; - virtual void release_params_backend_params(const std::vector& tensors) = 0; + virtual void release_params_backend_params(const std::vector& tensors) = 0; virtual ParamPrefetchResult enqueue_param_prefetch( uintptr_t owner_id, uint64_t segment_id, - const std::vector& tensors) = 0; + const std::vector& tensors) = 0; virtual bool activate_param_prefetch(uintptr_t owner_id, uint64_t segment_id, const std::vector& tensors) = 0;