21#include <unordered_map>
79 size_t legal_group_count = 0;
80 std::vector<SpecializationGroupInfo> installed_groups;
96 size_t input_data_size = 0,
98 float pool_multiplier = 3.0f
108 explicit Pipeline(
const std::string& config_path);
121 size_t getSpecializedGroupCount()
const {
return dag_ ? dag_->getFusedGroupCount() : 0; }
128 size_t getFusedGroupCount()
const {
return getSpecializedGroupCount(); }
132 void setNumStreams(
int num_streams);
140 void setDims(
size_t x,
size_t y = 1,
size_t z = 1) { dims_ = {x, y, z}; }
141 void setDims(std::array<size_t, 3> dims) { dims_ = dims; }
142 std::array<size_t, 3> getDims()
const {
return dims_; }
150 template<
typename StageT,
typename... Args>
151 StageT* addStage(Args&&... args);
160 int connect(
Stage* dependent,
Stage* producer,
const std::string& output_name =
"output");
163 int connect(
Stage* dependent,
const std::vector<Stage*>& producers);
188 void bindExternalInput(
Stage* stage);
202 void warmup(fz::stream_t stream = 0);
205 void setWarmupOnFinalize(
bool enable) { warmup_on_finalize_ = enable; }
206 bool isWarmupOnFinalizeEnabled()
const {
return warmup_on_finalize_; }
214 void setPoolManagedDecompOutput(
bool enable) { pool_managed_decomp_ = enable; }
215 bool isPoolManagedDecompOutput()
const {
return pool_managed_decomp_; }
249 void setPrimarySource(
Stage* stage) { primary_source_stage_ = stage; }
264 size_t getMaxCompressedSize(
size_t input_bytes)
const;
285 size_t getLastUncompressedSize()
const {
286 return original_input_size_ > 0 ? original_input_size_ : input_size_;
307 fz::stream_t stream = 0
340 size_t output_buf_capacity,
341 size_t* actual_output_size,
342 fz::stream_t stream = 0
364 fz::stream_t stream = 0
391 size_t output_buf_capacity,
392 size_t* actual_output_size,
393 fz::stream_t stream = 0
437 size_t output_buf_capacity,
438 size_t* actual_output_size,
439 fz::stream_t stream = 0
519 void prepareInverse(
size_t uncompressed_size);
522 void reset(fz::stream_t stream = 0);
530 void enableProfiling(
bool enable);
531 bool isProfilingEnabled()
const {
return profiling_enabled_; }
540 size_t getPoolThreshold()
const;
551 bool isMemPoolFallbackMode()
const;
558 void enableBoundsCheck(
bool enable) { dag_->enableBoundsCheck(enable); }
559 bool isBoundsCheckEnabled()
const {
return dag_->isBoundsCheckEnabled(); }
566 void setColoringEnabled(
bool enable) {
567 coloring_enabled_ = enable;
568 dag_->setColoringEnabled(enable);
570 bool isColoringEnabled()
const {
return dag_->isColoringEnabled(); }
571 size_t getColorRegionCount()
const {
return dag_->getColorRegionCount(); }
580 bool isColoringRequested()
const {
return coloring_enabled_; }
591 std::unordered_map<std::string, std::vector<std::string>> collectRunNotes()
const {
592 std::unordered_map<std::string, std::vector<std::string>> notes;
593 for (
const auto& s : stages_) {
595 auto n = s->getRunNotes();
596 if (!n.empty()) notes.emplace(s->getName(), std::move(n));
610 void enableGraphMode(
bool enable);
611 bool isGraphModeEnabled()
const {
return graph_mode_enabled_; }
623 void captureGraph(fz::stream_t stream = 0);
624 bool isGraphCaptured()
const {
return graph_captured_; }
626 size_t getPeakMemoryUsage()
const;
627 size_t getCurrentMemoryUsage()
const;
628 void printPipeline()
const;
635 std::vector<FZMStageInfo> stages;
636 std::vector<FZMBufferEntry> buffers;
640 void writeToFile(
const std::string& filename, fz::stream_t stream = 0);
643 static FZMFileHeader readHeader(
const std::string& filename);
664 std::vector<uint8_t> serializeHeaderToMemory()
const;
685 void primeInverseFromHeader(
const void* header_bytes,
size_t header_size);
701 static void decompressFromFile(
702 const std::string& filename,
705 fz::stream_t stream = 0,
707 size_t pool_override_bytes = 0
727 void decompressFromFileInstance(
728 const std::string& filename,
731 fz::stream_t stream = 0,
759 void decompressFromMemory(
760 const void* header_bytes,
766 fz::stream_t stream = 0
784 void loadConfig(
const std::string& path);
794 void saveConfig(
const std::string& path)
const;
805 ~PoolBuffer() { free(0); }
806 PoolBuffer() =
default;
807 PoolBuffer(
const PoolBuffer&) =
delete;
808 PoolBuffer& operator=(
const PoolBuffer&) =
delete;
810 void free(fz::stream_t s) {
811 if (ptr && pool) { pool->
free(ptr, s); ptr =
nullptr; capacity = 0; }
813 bool allocate(MemoryPool* p,
size_t bytes, fz::stream_t s,
814 const char* tag,
bool persistent =
false) {
817 ptr = pool->
allocate(bytes, s, tag, persistent);
818 if (ptr) capacity = bytes;
819 return ptr !=
nullptr;
824 struct PinnedBuffer {
828 ~PinnedBuffer() {
if (ptr) cudaFreeHost(ptr); }
829 PinnedBuffer() =
default;
830 PinnedBuffer(
const PinnedBuffer&) =
delete;
831 PinnedBuffer& operator=(
const PinnedBuffer&) =
delete;
834 bool ensureCapacity(
size_t bytes) {
835 if (capacity >= bytes)
return true;
836 if (ptr) { cudaFreeHost(ptr); ptr =
nullptr; capacity = 0; }
837 if (cudaHostAlloc(&ptr, bytes, cudaHostAllocDefault) != cudaSuccess)
return false;
844 struct DeviceBuffer {
848 ~DeviceBuffer() {
if (ptr) cudaFree(ptr); }
849 DeviceBuffer() =
default;
850 DeviceBuffer(
const DeviceBuffer&) =
delete;
851 DeviceBuffer& operator=(
const DeviceBuffer&) =
delete;
854 bool ensureCapacity(
size_t bytes) {
855 if (capacity >= bytes)
return true;
856 if (ptr) { cudaFree(ptr); ptr =
nullptr; capacity = 0; }
857 if (cudaMalloc(&ptr, bytes) != cudaSuccess)
return false;
865 Stage* addRawStage(Stage* stage);
867 struct OutputBuffer {
870 size_t allocated_size;
874 std::vector<OutputBuffer> getOutputBuffers()
const;
876 static void* loadCompressedData(
877 const std::string& filename,
878 const FZMFileHeader& header,
879 fz::stream_t stream = 0,
880 MemoryPool* pool =
nullptr
884 std::pair<std::vector<Stage*>, std::vector<Stage*>> identifyTopology();
885 void setupInputBuffers(
const std::vector<Stage*>& sources);
886 int autoDetectUnconnectedOutputs();
887 void detectMultiOutputScenario(
int pipeline_outputs);
888 void configureStreamsIfNeeded();
891 void typeCheckConnections();
892 void computeInputAlignment();
895 void bindSemanticContracts();
899 void planAndInstallFusion();
900 void notifyStagesFinalizeHooks();
901 void refinePoolSize();
902 void setupGraphModeInput();
903 void preallocatePadBuffer();
904 void preallocateConcatBuffers();
908 std::pair<const void*, size_t> prepareInputSource(
909 const void* d_input,
size_t input_size, fz::stream_t stream);
916 void propagateBufferSizes(
bool force_from_current_inputs =
false);
918 std::vector<Stage*> getSourceStages()
const;
919 std::vector<Stage*> getSinkStages()
const;
924 struct FwdStageDesc {
926 std::vector<int> output_buf_ids;
927 std::vector<int> input_buf_ids;
931 using PipelineOutputMap = std::unordered_map<int, std::pair<void*, size_t>>;
934 void buildOrReuseInvCache(
935 const PipelineOutputMap& po_map,
938 fz::stream_t stream);
948 std::vector<size_t> resolveDecompressSegmentSizes(
949 const void* d_input,
size_t input_size, fz::stream_t stream)
const;
957 PipelineOutputMap mapCompressedBufferPointers(
958 const void* d_input,
const std::vector<size_t>& seg_sizes)
const;
966 std::pair<Stage*, size_t> resolvePrimarySource()
const;
975 void* allocateDecompressOutput(
976 size_t actual_size,
void* caller_output,
size_t caller_capacity,
977 fz::stream_t stream);
995 size_t caller_capacity,
999 fz::stream_t stream);
1008 void buildStaticBufferMetadata();
1015 std::vector<size_t> readConcatSegmentSizes(
1016 const void* d_blob,
size_t n, fz::stream_t stream)
const;
1020 static FZMFileHeader parseHeaderFromMemory(
const void* data,
size_t size);
1021 static size_t computeFilePoolSize(
const FZMFileHeader& fh,
size_t pool_override_bytes);
1022 static std::pair<std::vector<std::unique_ptr<Stage>>, std::vector<FwdStageDesc>>
1023 reconstructForwardTopology(
const FZMFileHeader& fh);
1024 static std::unordered_map<Stage*, size_t> buildSourceSizesFromHeader(
1025 const FZMFileHeader& fh,
const std::vector<FwdStageDesc>& fwd_topology);
1032 static std::pair<std::unique_ptr<CompressionDAG>,
1033 std::unordered_map<Stage*, int>>
1035 const std::vector<FwdStageDesc>& fwd_stages,
1036 const PipelineOutputMap& pipeline_outputs,
1039 const std::unordered_map<Stage*, size_t>& source_sizes,
1040 bool enable_profiling,
1041 bool enable_inverse_fusion =
false
1046 struct OutputBufferInfo {
1050 std::string stage_name;
1051 std::string output_name;
1054 std::vector<OutputBufferInfo> collectOutputBuffers()
const;
1057 size_t calculateConcatSize(
const std::vector<OutputBufferInfo>& outputs)
const;
1059 size_t writeConcatBuffer(
1060 const std::vector<OutputBufferInfo>& outputs,
1061 uint8_t* d_concat_bytes,
1065 void concatOutputs(
void** d_output,
size_t* output_size, fz::stream_t stream);
1069 std::unique_ptr<MemoryPool> mem_pool_;
1070 std::unique_ptr<CompressionDAG> dag_;
1073 std::vector<std::unique_ptr<Stage>> stages_;
1074 std::unordered_map<Stage*, DAGNode*> stage_to_node_;
1076 struct ConnectionInfo {
1079 std::string output_name;
1082 std::vector<ConnectionInfo> connections_;
1086 bool warmup_on_finalize_;
1087 bool pool_managed_decomp_;
1091 bool is_compressed_;
1092 bool was_compressed_;
1094 bool profiling_enabled_;
1097 bool coloring_enabled_ =
true;
1098 PipelinePerfResult last_perf_result_;
1100 std::vector<DAGNode*> input_nodes_;
1105 Stage* primary_source_stage_ =
nullptr;
1112 std::vector<std::pair<DAGNode*, int>> explicit_external_bindings_;
1113 std::vector<DAGNode*> output_nodes_;
1114 std::vector<int> input_buffer_ids_;
1115 std::vector<int> output_buffer_ids_;
1117 PoolBuffer d_concat_buffer_;
1122 std::vector<void*> d_decomp_outputs_;
1125 PinnedBuffer h_concat_header_;
1127 PinnedBuffer h_copy_descs_;
1128 DeviceBuffer d_copy_descs_;
1134 std::vector<size_t> source_input_sizes_;
1138 size_t input_alignment_bytes_;
1139 PoolBuffer d_pad_buf_;
1143 size_t original_input_size_;
1145 size_t input_size_hint_;
1146 float pool_multiplier_;
1150 std::array<size_t, 3> dims_;
1160 struct InvDAGCache {
1161 std::unique_ptr<CompressionDAG> inv_dag;
1162 std::unordered_map<Stage*, int> inv_result_map;
1163 std::unordered_map<int, int> fwd_to_inv_ext_buf;
1164 std::unordered_map<Stage*, size_t> source_sizes;
1166 std::unique_ptr<InvDAGCache> inv_cache_;
1168 struct BufferMetadata {
1171 size_t allocated_size;
1176 std::vector<BufferMetadata> buffer_metadata_;
1178 bool graph_mode_enabled_;
1179 bool graph_captured_;
1185 PoolBuffer d_graph_input_;
1186 size_t d_graph_input_size_;
1188 fz::graph_t captured_graph_;
1189 fz::graph_exec_t graph_exec_;
1194template<
typename StageT,
typename... Args>
1195StageT* Pipeline::addStage(Args&&... args) {
1196 if (is_finalized_) {
1197 throw std::runtime_error(
"Cannot add stages after finalization");
1207 if constexpr (!StageT::isSupportedOnBackend()) {
1208 throw std::runtime_error(
1209 "addStage(): this stage type is not supported on the current "
1210 "GPU backend (FZGMOD_BACKEND) this library was built for");
1212 auto stage_ptr = std::make_unique<StageT>(std::forward<Args>(args)...);
1213 StageT* stage = stage_ptr.get();
1217 DAGNode* node = dag_->addStage(stage, stage->getName());
1218 size_t num_outputs = stage->getNumOutputs();
1219 auto output_names = stage->getOutputNames();
1223 for (
size_t i = 0; i < num_outputs; i++) {
1224 std::string out_name = i < output_names.size() ? output_names[i] : std::to_string(i);
1225 dag_->addUnconnectedOutput(node, 1, i, stage->getName() +
"." + out_name +
"_unconnected");
1228 stage_to_node_[stage] = node;
1229 stages_.push_back(std::move(stage_ptr));
Definition device_buffer.h:54
void free(void *ptr, fz::stream_t stream)
void * allocate(size_t size, fz::stream_t stream, const std::string &tag="", bool persistent=false)
Definition device_buffer.h:84
virtual void setDims(const std::array< size_t, 3 > &dims)
Definition stage.h:207
TOML-based pipeline configuration file support.
Compression DAG wiring, execution, and memory strategy types.
Backend-neutral device span and buffer value types.
Stream-ordered CUDA memory pool for pipeline buffer management.
SpecializationPolicy
Pipeline Specialization policy.
Definition compressor.h:62
MemoryStrategy
Definition dag.h:32
@ MINIMAL
Allocate on-demand, free at last consumer. Lowest peak memory.
SpecializationInfo FusionInfo
Definition compressor.h:86
SpecializationPolicy FusionPolicy
Definition compressor.h:66
Pipeline and per-stage profiling result types.
Base class interface for all compression stages.
Backward-compatible shim.
Definition device_buffer.h:35
Definition device_buffer.h:24
One finalize-time specialization selected for execution (compress or inverse).
Definition compressor.h:69
std::vector< std::string > stages
the stages it replaced
Definition compressor.h:71
std::string execution_path
last runtime subpath; empty if not applicable
Definition compressor.h:72
std::string implementation
strategy impl name, e.g. "warp-register"
Definition compressor.h:70
Resolved specialization decision for diagnostics and benchmark provenance.
Definition compressor.h:77
std::string fallback_reason
policy_off, no_legal_group, or legacy no_profitable_implementation; empty on a hit.
Definition compressor.h:84
std::vector< SpecializationGroupInfo > installed_inverse_groups
Lazily populated after the first decompress builds its inverse DAG.
Definition compressor.h:82
Backend-neutral GPU type aliases.