21#include <unordered_map>
52 std::string implementation;
53 std::vector<std::string> stages;
59 size_t legal_group_count = 0;
60 std::vector<FusionGroupInfo> installed_groups;
73 size_t input_data_size = 0,
75 float pool_multiplier = 3.0f
85 explicit Pipeline(
const std::string& config_path);
95 void setFusionPolicy(
FusionPolicy mode) { fusion_policy_ = mode; }
96 FusionPolicy getFusionPolicy()
const {
return fusion_policy_; }
98 size_t getFusedGroupCount()
const {
return dag_ ? dag_->getFusedGroupCount() : 0; }
100 const FusionInfo& getFusionInfo()
const {
return fusion_info_; }
103 void setNumStreams(
int num_streams);
111 void setDims(
size_t x,
size_t y = 1,
size_t z = 1) { dims_ = {x, y, z}; }
112 void setDims(std::array<size_t, 3> dims) { dims_ = dims; }
113 std::array<size_t, 3> getDims()
const {
return dims_; }
121 template<
typename StageT,
typename... Args>
122 StageT* addStage(Args&&... args);
131 int connect(Stage* dependent, Stage* producer,
const std::string& output_name =
"output");
134 int connect(Stage* dependent,
const std::vector<Stage*>& producers);
159 void bindExternalInput(Stage* stage);
173 void warmup(fz::stream_t stream = 0);
176 void setWarmupOnFinalize(
bool enable) { warmup_on_finalize_ = enable; }
177 bool isWarmupOnFinalizeEnabled()
const {
return warmup_on_finalize_; }
185 void setPoolManagedDecompOutput(
bool enable) { pool_managed_decomp_ = enable; }
186 bool isPoolManagedDecompOutput()
const {
return pool_managed_decomp_; }
220 void setPrimarySource(Stage* stage) { primary_source_stage_ = stage; }
235 size_t getMaxCompressedSize(
size_t input_bytes)
const;
256 size_t getLastUncompressedSize()
const {
257 return original_input_size_ > 0 ? original_input_size_ : input_size_;
278 fz::stream_t stream = 0
311 size_t output_buf_capacity,
312 size_t* actual_output_size,
313 fz::stream_t stream = 0
335 fz::stream_t stream = 0
362 size_t output_buf_capacity,
363 size_t* actual_output_size,
364 fz::stream_t stream = 0
408 size_t output_buf_capacity,
409 size_t* actual_output_size,
410 fz::stream_t stream = 0
425 BorrowedDeviceBuffer compress(ConstDeviceSpan input, fz::stream_t stream = 0);
432 size_t compressInto(ConstDeviceSpan input, DeviceSpan output, fz::stream_t stream = 0);
442 BorrowedDeviceBuffer decompressBorrowed(ConstDeviceSpan input, fz::stream_t stream = 0);
449 OwnedDeviceBuffer decompressOwned(ConstDeviceSpan input, fz::stream_t stream = 0);
456 size_t decompressInto(ConstDeviceSpan input, DeviceSpan output, fz::stream_t stream = 0);
464 size_t decompressIntoAsync(ConstDeviceSpan input, DeviceSpan output, fz::stream_t stream = 0);
490 void prepareInverse(
size_t uncompressed_size);
493 void reset(fz::stream_t stream = 0);
501 void enableProfiling(
bool enable);
502 bool isProfilingEnabled()
const {
return profiling_enabled_; }
505 const PipelinePerfResult& getLastPerfResult()
const {
return last_perf_result_; }
508 CompressionDAG* getDAG() {
return dag_.get(); }
511 size_t getPoolThreshold()
const;
522 bool isMemPoolFallbackMode()
const;
529 void enableBoundsCheck(
bool enable) { dag_->enableBoundsCheck(enable); }
530 bool isBoundsCheckEnabled()
const {
return dag_->isBoundsCheckEnabled(); }
537 void setColoringEnabled(
bool enable) {
538 coloring_enabled_ = enable;
539 dag_->setColoringEnabled(enable);
541 bool isColoringEnabled()
const {
return dag_->isColoringEnabled(); }
542 size_t getColorRegionCount()
const {
return dag_->getColorRegionCount(); }
551 bool isColoringRequested()
const {
return coloring_enabled_; }
562 std::unordered_map<std::string, std::vector<std::string>> collectRunNotes()
const {
563 std::unordered_map<std::string, std::vector<std::string>> notes;
564 for (
const auto& s : stages_) {
566 auto n = s->getRunNotes();
567 if (!n.empty()) notes.emplace(s->getName(), std::move(n));
581 void enableGraphMode(
bool enable);
582 bool isGraphModeEnabled()
const {
return graph_mode_enabled_; }
594 void captureGraph(fz::stream_t stream = 0);
595 bool isGraphCaptured()
const {
return graph_captured_; }
597 size_t getPeakMemoryUsage()
const;
598 size_t getCurrentMemoryUsage()
const;
599 void printPipeline()
const;
606 std::vector<FZMStageInfo> stages;
607 std::vector<FZMBufferEntry> buffers;
611 void writeToFile(
const std::string& filename, fz::stream_t stream = 0);
614 static FZMFileHeader readHeader(
const std::string& filename);
635 std::vector<uint8_t> serializeHeaderToMemory()
const;
656 void primeInverseFromHeader(
const void* header_bytes,
size_t header_size);
672 static void decompressFromFile(
673 const std::string& filename,
676 fz::stream_t stream = 0,
678 size_t pool_override_bytes = 0
698 void decompressFromFileInstance(
699 const std::string& filename,
702 fz::stream_t stream = 0,
730 void decompressFromMemory(
731 const void* header_bytes,
737 fz::stream_t stream = 0
755 void loadConfig(
const std::string& path);
765 void saveConfig(
const std::string& path)
const;
776 ~PoolBuffer() { free(0); }
777 PoolBuffer() =
default;
778 PoolBuffer(
const PoolBuffer&) =
delete;
779 PoolBuffer& operator=(
const PoolBuffer&) =
delete;
781 void free(fz::stream_t s) {
782 if (ptr && pool) { pool->
free(ptr, s); ptr =
nullptr; capacity = 0; }
784 bool allocate(MemoryPool* p,
size_t bytes, fz::stream_t s,
785 const char* tag,
bool persistent =
false) {
788 ptr = pool->
allocate(bytes, s, tag, persistent);
789 if (ptr) capacity = bytes;
790 return ptr !=
nullptr;
795 struct PinnedBuffer {
799 ~PinnedBuffer() {
if (ptr) cudaFreeHost(ptr); }
800 PinnedBuffer() =
default;
801 PinnedBuffer(
const PinnedBuffer&) =
delete;
802 PinnedBuffer& operator=(
const PinnedBuffer&) =
delete;
805 bool ensureCapacity(
size_t bytes) {
806 if (capacity >= bytes)
return true;
807 if (ptr) { cudaFreeHost(ptr); ptr =
nullptr; capacity = 0; }
808 if (cudaHostAlloc(&ptr, bytes, cudaHostAllocDefault) != cudaSuccess)
return false;
815 struct DeviceBuffer {
819 ~DeviceBuffer() {
if (ptr) cudaFree(ptr); }
820 DeviceBuffer() =
default;
821 DeviceBuffer(
const DeviceBuffer&) =
delete;
822 DeviceBuffer& operator=(
const DeviceBuffer&) =
delete;
825 bool ensureCapacity(
size_t bytes) {
826 if (capacity >= bytes)
return true;
827 if (ptr) { cudaFree(ptr); ptr =
nullptr; capacity = 0; }
828 if (cudaMalloc(&ptr, bytes) != cudaSuccess)
return false;
836 Stage* addRawStage(Stage* stage);
838 struct OutputBuffer {
841 size_t allocated_size;
845 std::vector<OutputBuffer> getOutputBuffers()
const;
847 static void* loadCompressedData(
848 const std::string& filename,
849 const FZMFileHeader& header,
850 fz::stream_t stream = 0,
851 MemoryPool* pool =
nullptr
855 std::pair<std::vector<Stage*>, std::vector<Stage*>> identifyTopology();
856 void setupInputBuffers(
const std::vector<Stage*>& sources);
857 int autoDetectUnconnectedOutputs();
858 void detectMultiOutputScenario(
int pipeline_outputs);
859 void configureStreamsIfNeeded();
862 void typeCheckConnections();
863 void computeInputAlignment();
866 void bindSemanticContracts();
870 void planAndInstallFusion();
871 void notifyStagesFinalizeHooks();
872 void refinePoolSize();
873 void setupGraphModeInput();
874 void preallocatePadBuffer();
875 void preallocateConcatBuffers();
879 std::pair<const void*, size_t> prepareInputSource(
880 const void* d_input,
size_t input_size, fz::stream_t stream);
887 void propagateBufferSizes(
bool force_from_current_inputs =
false);
889 std::vector<Stage*> getSourceStages()
const;
890 std::vector<Stage*> getSinkStages()
const;
895 struct FwdStageDesc {
897 std::vector<int> output_buf_ids;
898 std::vector<int> input_buf_ids;
902 using PipelineOutputMap = std::unordered_map<int, std::pair<void*, size_t>>;
905 void buildOrReuseInvCache(
906 const PipelineOutputMap& po_map,
909 fz::stream_t stream);
927 size_t caller_capacity,
931 fz::stream_t stream);
940 void buildStaticBufferMetadata();
947 std::vector<size_t> readConcatSegmentSizes(
948 const void* d_blob,
size_t n, fz::stream_t stream)
const;
952 static FZMFileHeader parseHeaderFromMemory(
const void* data,
size_t size);
953 static size_t computeFilePoolSize(
const FZMFileHeader& fh,
size_t pool_override_bytes);
954 static std::pair<std::vector<std::unique_ptr<Stage>>, std::vector<FwdStageDesc>>
955 reconstructForwardTopology(
const FZMFileHeader& fh);
956 static std::unordered_map<Stage*, size_t> buildSourceSizesFromHeader(
957 const FZMFileHeader& fh,
const std::vector<FwdStageDesc>& fwd_topology);
964 static std::pair<std::unique_ptr<CompressionDAG>,
965 std::unordered_map<Stage*, int>>
967 const std::vector<FwdStageDesc>& fwd_stages,
968 const PipelineOutputMap& pipeline_outputs,
971 const std::unordered_map<Stage*, size_t>& source_sizes,
972 bool enable_profiling
977 struct OutputBufferInfo {
981 std::string stage_name;
982 std::string output_name;
985 std::vector<OutputBufferInfo> collectOutputBuffers()
const;
988 size_t calculateConcatSize(
const std::vector<OutputBufferInfo>& outputs)
const;
990 size_t writeConcatBuffer(
991 const std::vector<OutputBufferInfo>& outputs,
992 uint8_t* d_concat_bytes,
996 void concatOutputs(
void** d_output,
size_t* output_size, fz::stream_t stream);
1000 std::unique_ptr<MemoryPool> mem_pool_;
1001 std::unique_ptr<CompressionDAG> dag_;
1004 std::vector<std::unique_ptr<Stage>> stages_;
1005 std::unordered_map<Stage*, DAGNode*> stage_to_node_;
1007 struct ConnectionInfo {
1010 std::string output_name;
1013 std::vector<ConnectionInfo> connections_;
1017 bool warmup_on_finalize_;
1018 bool pool_managed_decomp_;
1022 bool is_compressed_;
1023 bool was_compressed_;
1025 bool profiling_enabled_;
1028 bool coloring_enabled_ =
true;
1029 PipelinePerfResult last_perf_result_;
1031 std::vector<DAGNode*> input_nodes_;
1036 Stage* primary_source_stage_ =
nullptr;
1043 std::vector<std::pair<DAGNode*, int>> explicit_external_bindings_;
1044 std::vector<DAGNode*> output_nodes_;
1045 std::vector<int> input_buffer_ids_;
1046 std::vector<int> output_buffer_ids_;
1048 PoolBuffer d_concat_buffer_;
1053 std::vector<void*> d_decomp_outputs_;
1056 PinnedBuffer h_concat_header_;
1058 PinnedBuffer h_copy_descs_;
1059 DeviceBuffer d_copy_descs_;
1065 std::vector<size_t> source_input_sizes_;
1069 size_t input_alignment_bytes_;
1070 PoolBuffer d_pad_buf_;
1074 size_t original_input_size_;
1076 size_t input_size_hint_;
1077 float pool_multiplier_;
1081 std::array<size_t, 3> dims_;
1091 struct InvDAGCache {
1092 std::unique_ptr<CompressionDAG> inv_dag;
1093 std::unordered_map<Stage*, int> inv_result_map;
1094 std::unordered_map<int, int> fwd_to_inv_ext_buf;
1095 std::unordered_map<Stage*, size_t> source_sizes;
1097 std::unique_ptr<InvDAGCache> inv_cache_;
1099 struct BufferMetadata {
1102 size_t allocated_size;
1107 std::vector<BufferMetadata> buffer_metadata_;
1109 bool graph_mode_enabled_;
1110 bool graph_captured_;
1112 FusionInfo fusion_info_;
1116 PoolBuffer d_graph_input_;
1117 size_t d_graph_input_size_;
1119 fz::graph_t captured_graph_;
1120 fz::graph_exec_t graph_exec_;
1125template<
typename StageT,
typename... Args>
1126StageT* Pipeline::addStage(Args&&... args) {
1127 if (is_finalized_) {
1128 throw std::runtime_error(
"Cannot add stages after finalization");
1138 if constexpr (!StageT::isSupportedOnBackend()) {
1139 throw std::runtime_error(
1140 "addStage(): this stage type is not supported on the current "
1141 "GPU backend (FZGMOD_BACKEND) this library was built for");
1143 auto stage_ptr = std::make_unique<StageT>(std::forward<Args>(args)...);
1144 StageT* stage = stage_ptr.get();
1146 stage->setDims(dims_);
1148 DAGNode* node = dag_->addStage(stage, stage->getName());
1149 size_t num_outputs = stage->getNumOutputs();
1150 auto output_names = stage->getOutputNames();
1154 for (
size_t i = 0; i < num_outputs; i++) {
1155 std::string out_name = i < output_names.size() ? output_names[i] : std::to_string(i);
1156 dag_->addUnconnectedOutput(node, 1, i, stage->getName() +
"." + out_name +
"_unconnected");
1159 stage_to_node_[stage] = node;
1160 stages_.push_back(std::move(stage_ptr));
void free(void *ptr, fz::stream_t stream)
void * allocate(size_t size, fz::stream_t stream, const std::string &tag="", bool persistent=false)
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.
MemoryStrategy
Definition dag.h:32
@ MINIMAL
Allocate on-demand, free at last consumer. Lowest peak memory.
FusionPolicy
Definition compressor.h:48
Pipeline and per-stage profiling result types.
Base class interface for all compression stages.
Backward-compatible shim.
One finalize-time fusion specialization selected for compress execution.
Definition compressor.h:51
Resolved fusion decision for diagnostics and benchmark provenance.
Definition compressor.h:57
std::string fallback_reason
policy_off, no_legal_group, or no_profitable_implementation; empty on a hit.
Definition compressor.h:62
Backend-neutral GPU type aliases.