feat: stream model conversion (#1581)

This commit is contained in:
Shikaku2
2026-07-05 05:46:07 -04:00
committed by GitHub
parent e9dee542c4
commit da6db07c5c
12 changed files with 676 additions and 82 deletions

View File

@@ -1,14 +1,33 @@
#include <algorithm>
#include <condition_variable>
#include <cstdint>
#include <cstring>
#include <exception>
#include <fstream>
#include <memory>
#include <mutex>
#include <regex>
#include <string>
#include <thread>
#include <vector>
#include "core/util.h"
#include "model_io/gguf_io.h"
#include "model_io/safetensors_io.h"
#include "model_io/streaming_writer.h"
#include "model_loader.h"
#include "util.h"
#include "ggml_extend_backend.h"
struct TensorExportInfo {
TensorStorage storage;
ggml_type type;
};
struct TensorExportJob {
TensorExportInfo info;
std::vector<uint8_t> data;
std::string error;
bool success = false;
};
static ggml_type get_export_tensor_type(ModelLoader& model_loader,
const TensorStorage& tensor_storage,
@@ -33,47 +52,262 @@ static ggml_type get_export_tensor_type(ModelLoader& model_loader,
return tensor_type;
}
static bool load_tensors_for_export(ModelLoader& model_loader,
ggml_context* ggml_ctx,
ggml_type type,
const TensorTypeRules& tensor_type_rules,
std::vector<TensorWriteInfo>& tensors) {
std::mutex tensor_mutex;
auto on_new_tensor_cb = [&](const TensorStorage& tensor_storage, ggml_tensor** dst_tensor) -> bool {
const std::string& name = tensor_storage.name;
ggml_type tensor_type = get_export_tensor_type(model_loader, tensor_storage, type, tensor_type_rules);
static bool collect_tensors_for_export(ModelLoader& model_loader,
ggml_type type,
const TensorTypeRules& tensor_type_rules,
std::vector<TensorExportInfo>& tensors) {
tensors.clear();
tensors.reserve(model_loader.get_tensor_storage_map().size());
for (const auto& kv : model_loader.get_tensor_storage_map()) {
const TensorStorage& tensor_storage = kv.second;
TensorExportInfo info;
info.storage = tensor_storage;
info.type = get_export_tensor_type(model_loader, tensor_storage, type, tensor_type_rules);
tensors.push_back(std::move(info));
}
LOG_INFO("collected %zu tensors for export", tensors.size());
return true;
}
std::lock_guard<std::mutex> lock(tensor_mutex);
ggml_tensor* tensor = ggml_new_tensor(ggml_ctx, tensor_type, tensor_storage.n_dims, tensor_storage.ne);
if (tensor == nullptr) {
LOG_ERROR("ggml_new_tensor failed");
static size_t export_tensor_nbytes(const TensorExportInfo& info) {
TensorStorage output_storage = info.storage;
output_storage.type = info.type;
return static_cast<size_t>(output_storage.nbytes());
}
static TensorWritePlan tensor_write_plan_from_export_info(const TensorExportInfo& info) {
TensorWritePlan plan;
plan.name = info.storage.name;
plan.type = info.type;
plan.n_dims = info.storage.n_dims;
for (int i = 0; i < SD_MAX_DIMS; i++) {
plan.ne[i] = info.storage.ne[i];
}
return plan;
}
static std::vector<TensorWritePlan> tensor_write_plans_from_export_infos(const std::vector<TensorExportInfo>& tensors) {
std::vector<TensorWritePlan> plans;
plans.reserve(tensors.size());
for (const TensorExportInfo& info : tensors) {
plans.push_back(tensor_write_plan_from_export_info(info));
}
return plans;
}
static bool preallocate_output_file(const std::string& output_path, uint64_t file_size, std::string* error) {
if (file_size == 0) {
return true;
}
std::fstream file(output_path, std::ios::binary | std::ios::in | std::ios::out);
if (!file.is_open()) {
if (error != nullptr) {
*error = "failed to open output file '" + output_path + "' for preallocation";
}
return false;
}
// This portable fallback sets the final file size. A platform-specific
// posix_fallocate/ftruncate path can replace it later.
file.seekp(static_cast<std::streamoff>(file_size - 1), std::ios::beg);
file.put('\0');
file.flush();
if (!file) {
if (error != nullptr) {
*error = "failed to preallocate output file '" + output_path + "'";
}
return false;
}
return true;
}
static bool load_tensor_for_export(ModelLoader& model_loader, TensorExportJob& job) {
size_t mem_size = 1 * 1024 * 1024;
mem_size += ggml_tensor_overhead();
TensorStorage output_storage = job.info.storage;
output_storage.type = job.info.type;
mem_size += static_cast<size_t>(output_storage.nbytes());
ggml_context* ggml_ctx = ggml_init({mem_size, nullptr, false});
if (ggml_ctx == nullptr) {
job.error = "ggml_init failed for tensor '" + job.info.storage.name + "'";
return false;
}
ggml_tensor* tensor = ggml_new_tensor(ggml_ctx, job.info.type, job.info.storage.n_dims, job.info.storage.ne);
if (tensor == nullptr) {
ggml_free(ggml_ctx);
job.error = "ggml_new_tensor failed for tensor '" + job.info.storage.name + "'";
return false;
}
ggml_set_name(tensor, job.info.storage.name.c_str());
const size_t tensor_nbytes = ggml_nbytes(tensor);
if (tensor_nbytes > 0 && !model_loader.load_tensor(job.info.storage, tensor)) {
ggml_free(ggml_ctx);
job.error = "failed to load tensor '" + job.info.storage.name + "'";
return false;
}
job.data.resize(tensor_nbytes);
if (tensor_nbytes > 0) {
memcpy(job.data.data(), tensor->data, tensor_nbytes);
}
ggml_free(ggml_ctx);
return true;
}
static bool stream_tensor_data(ModelLoader& model_loader,
const std::string& output_path,
const std::vector<TensorExportInfo>& tensors,
const StreamingModelWriter& writer,
int n_threads,
std::string* error) {
n_threads = n_threads > 0 ? n_threads : sd_get_num_physical_cores();
n_threads = std::max(1, n_threads);
LOG_INFO("streaming convert with %d threads", n_threads);
int64_t start_time = ggml_time_ms();
uint64_t bytes_written = 0;
size_t tensors_written = 0;
size_t next_tensor_index = 0;
bool failed = false;
std::string failure;
const size_t memory_budget = 1024ull * 1024ull * 1024ull;
size_t reserved_bytes = 0;
std::mutex work_mutex;
std::mutex progress_mutex;
std::condition_variable memory_cv;
std::vector<std::thread> workers;
workers.reserve(n_threads);
auto reserve_memory = [&](size_t bytes) -> bool {
std::unique_lock<std::mutex> lock(work_mutex);
memory_cv.wait(lock, [&]() {
return failed || reserved_bytes == 0 || reserved_bytes + bytes <= memory_budget;
});
if (failed) {
return false;
}
ggml_set_name(tensor, name.c_str());
if (!tensor->data) {
GGML_ASSERT(ggml_nelements(tensor) == 0);
// Avoid crashing writers by setting a dummy pointer for zero-sized tensors.
LOG_DEBUG("setting dummy pointer for zero-sized tensor %s", name.c_str());
tensor->data = ggml_get_mem_buffer(ggml_ctx);
}
TensorWriteInfo write_info;
write_info.tensor = tensor;
write_info.n_dims = tensor_storage.n_dims;
for (int i = 0; i < tensor_storage.n_dims; ++i) {
write_info.ne[i] = tensor_storage.ne[i];
}
*dst_tensor = tensor;
tensors.push_back(std::move(write_info));
reserved_bytes += bytes;
return true;
};
bool success = model_loader.load_tensors(on_new_tensor_cb);
LOG_INFO("load tensors done");
return success;
auto release_memory = [&](size_t bytes) {
{
std::lock_guard<std::mutex> lock(work_mutex);
reserved_bytes -= std::min(reserved_bytes, bytes);
}
memory_cv.notify_all();
};
auto fail = [&](const std::string& message) {
{
std::lock_guard<std::mutex> lock(work_mutex);
if (!failed) {
failed = true;
failure = message;
}
}
memory_cv.notify_all();
};
for (int worker = 0; worker < n_threads; worker++) {
workers.emplace_back([&]() {
std::fstream output_file(output_path, std::ios::binary | std::ios::in | std::ios::out);
if (!output_file.is_open()) {
fail("failed to open output file '" + output_path + "' for tensor writing");
return;
}
while (true) {
size_t tensor_index = 0;
{
std::lock_guard<std::mutex> lock(work_mutex);
if (failed || next_tensor_index >= tensors.size()) {
return;
}
tensor_index = next_tensor_index++;
}
const size_t tensor_bytes = export_tensor_nbytes(tensors[tensor_index]);
if (!reserve_memory(tensor_bytes)) {
return;
}
TensorExportJob job;
job.info = tensors[tensor_index];
try {
job.success = load_tensor_for_export(model_loader, job);
} catch (const std::exception& e) {
job.error = e.what();
job.success = false;
}
if (!job.success) {
release_memory(tensor_bytes);
fail(job.error.empty() ? "streaming conversion failed" : job.error);
return;
}
std::string write_error;
if (!writer.write_tensor(output_file,
tensor_index,
job.data.empty() ? nullptr : job.data.data(),
job.data.size(),
&write_error)) {
release_memory(tensor_bytes);
fail(write_error.empty() ? "streaming conversion write failed" : write_error);
return;
}
{
std::lock_guard<std::mutex> lock(progress_mutex);
bytes_written += job.data.size();
tensors_written++;
float elapsed_seconds = (ggml_time_ms() - start_time) / 1000.0f;
pretty_bytes_progress(static_cast<int>(tensors_written),
static_cast<int>(tensors.size()),
bytes_written,
elapsed_seconds);
}
release_memory(tensor_bytes);
}
});
}
for (auto& worker : workers) {
worker.join();
}
printf("\n");
if (failed) {
if (error != nullptr) {
*error = failure;
}
return false;
}
LOG_INFO("streaming conversion completed, taking %.2fs", (ggml_time_ms() - start_time) / 1000.f);
return true;
}
static bool write_model_file_streaming(ModelLoader& model_loader,
const std::string& output_path,
const std::vector<TensorExportInfo>& tensors,
StreamingModelWriter& writer,
int n_threads,
std::string* error) {
std::vector<TensorWritePlan> plans = tensor_write_plans_from_export_infos(tensors);
if (!writer.write_metadata(output_path, plans, error)) {
return false;
}
if (!preallocate_output_file(output_path, writer.file_size(), error)) {
return false;
}
model_loader.process_model_files(false, false);
return stream_tensor_data(model_loader, output_path, tensors, writer, n_threads, error);
}
static bool init_convert_path(ModelLoader& model_loader, const char* path, const char* prefix, bool& loaded_any) {
@@ -91,42 +325,29 @@ static bool init_convert_path(ModelLoader& model_loader, const char* path, const
static bool export_loaded_model(ModelLoader& model_loader,
const char* output_path,
sd_type_t output_type,
const char* tensor_type_rules) {
const char* tensor_type_rules,
int n_threads) {
ggml_type type = sd_type_to_ggml_type(output_type);
bool output_is_safetensors = ends_with(output_path, ".safetensors");
TensorTypeRules type_rules = parse_tensor_type_rules(tensor_type_rules);
auto backend = sd_backend_cpu_init();
size_t mem_size = 1 * 1024 * 1024; // for padding
mem_size += model_loader.get_tensor_storage_map().size() * ggml_tensor_overhead();
mem_size += model_loader.get_params_mem_size(backend, type);
LOG_INFO("model tensors mem size: %.2fMB", mem_size / 1024.f / 1024.f);
ggml_context* ggml_ctx = ggml_init({mem_size, nullptr, false});
if (ggml_ctx == nullptr) {
LOG_ERROR("ggml_init failed for converter");
ggml_backend_free(backend);
return false;
}
std::vector<TensorWriteInfo> tensors;
bool success = load_tensors_for_export(model_loader, ggml_ctx, type, type_rules, tensors);
ggml_backend_free(backend);
std::vector<TensorExportInfo> tensors;
bool success = collect_tensors_for_export(model_loader, type, type_rules, tensors);
std::string error;
if (success) {
std::unique_ptr<StreamingModelWriter> writer;
if (output_is_safetensors) {
success = write_safetensors_file(output_path, tensors, &error);
writer = std::make_unique<SafetensorsStreamingWriter>();
} else {
success = write_gguf_file(output_path, tensors, &error);
writer = std::make_unique<GGUFStreamingWriter>();
}
success = write_model_file_streaming(model_loader, output_path, tensors, *writer, n_threads, &error);
}
if (!success && !error.empty()) {
LOG_ERROR("%s", error.c_str());
}
ggml_free(ggml_ctx);
return success;
}
@@ -139,7 +360,8 @@ bool convert_with_components(const char* model_path,
const char* output_path,
sd_type_t output_type,
const char* tensor_type_rules,
bool convert_name) {
bool convert_name,
int n_threads) {
ModelLoader model_loader;
bool loaded_any = false;
@@ -161,7 +383,7 @@ bool convert_with_components(const char* model_path,
model_loader.convert_tensors_name();
}
return export_loaded_model(model_loader, output_path, output_type, tensor_type_rules);
return export_loaded_model(model_loader, output_path, output_type, tensor_type_rules, n_threads);
}
bool convert(const char* input_path,
@@ -179,5 +401,6 @@ bool convert(const char* input_path,
output_path,
output_type,
tensor_type_rules,
convert_name);
convert_name,
0);
}