42 size_t input_data_size = 0,
44 float pool_multiplier = 3.0f
54 explicit Pipeline(
const std::string& config_path);
72 void setDims(
size_t x,
size_t y = 1,
size_t z = 1) { dims_ = {x, y, z}; }
73 void setDims(std::array<size_t, 3> dims) { dims_ = dims; }
74 std::array<size_t, 3> getDims()
const {
return dims_; }
82 template<
typename StageT,
typename... Args>
95 int connect(
Stage* dependent,
const std::vector<Stage*>& producers);
113 bool isWarmupOnFinalizeEnabled()
const {
return warmup_on_finalize_; }
122 bool isPoolManagedDecompOutput()
const {
return pool_managed_decomp_; }
159 return original_input_size_ > 0 ? original_input_size_ : input_size_;
180 cudaStream_t stream = 0
213 size_t output_buf_capacity,
214 size_t* actual_output_size,
215 cudaStream_t stream = 0
237 cudaStream_t stream = 0
264 size_t output_buf_capacity,
265 size_t* actual_output_size,
266 cudaStream_t stream = 0
310 size_t output_buf_capacity,
311 size_t* actual_output_size,
312 cudaStream_t stream = 0
342 void reset(cudaStream_t stream = 0);
351 bool isProfilingEnabled()
const {
return profiling_enabled_; }
379 bool isBoundsCheckEnabled()
const {
return dag_->isBoundsCheckEnabled(); }
387 bool isColoringEnabled()
const {
return dag_->isColoringEnabled(); }
388 size_t getColorRegionCount()
const {
return dag_->getColorRegionCount(); }
400 bool isGraphModeEnabled()
const {
return graph_mode_enabled_; }
413 bool isGraphCaptured()
const {
return graph_captured_; }
415 size_t getPeakMemoryUsage()
const;
416 size_t getCurrentMemoryUsage()
const;
417 void printPipeline()
const;
424 std::vector<FZMStageInfo> stages;
425 std::vector<FZMBufferEntry> buffers;
429 void writeToFile(
const std::string& filename, cudaStream_t stream = 0);
491 const std::string& filename,
494 cudaStream_t stream = 0,
496 size_t pool_override_bytes = 0
517 const std::string& filename,
520 cudaStream_t stream = 0,
549 const void* header_bytes,
555 cudaStream_t stream = 0
594 ~PoolBuffer() { free(0); }
595 PoolBuffer() =
default;
596 PoolBuffer(
const PoolBuffer&) =
delete;
597 PoolBuffer& operator=(
const PoolBuffer&) =
delete;
599 void free(cudaStream_t s) {
600 if (ptr && pool) { pool->
free(ptr, s); ptr =
nullptr; capacity = 0; }
602 bool allocate(MemoryPool* p,
size_t bytes, cudaStream_t s,
603 const char* tag,
bool persistent =
false) {
606 ptr = pool->
allocate(bytes, s, tag, persistent);
607 if (ptr) capacity = bytes;
608 return ptr !=
nullptr;
613 struct PinnedBuffer {
617 ~PinnedBuffer() {
if (ptr) cudaFreeHost(ptr); }
618 PinnedBuffer() =
default;
619 PinnedBuffer(
const PinnedBuffer&) =
delete;
620 PinnedBuffer& operator=(
const PinnedBuffer&) =
delete;
623 bool ensureCapacity(
size_t bytes) {
624 if (capacity >= bytes)
return true;
625 if (ptr) { cudaFreeHost(ptr); ptr =
nullptr; capacity = 0; }
626 if (cudaHostAlloc(&ptr, bytes, cudaHostAllocDefault) != cudaSuccess)
return false;
633 struct DeviceBuffer {
637 ~DeviceBuffer() {
if (ptr) cudaFree(ptr); }
638 DeviceBuffer() =
default;
639 DeviceBuffer(
const DeviceBuffer&) =
delete;
640 DeviceBuffer& operator=(
const DeviceBuffer&) =
delete;
643 bool ensureCapacity(
size_t bytes) {
644 if (capacity >= bytes)
return true;
645 if (ptr) { cudaFree(ptr); ptr =
nullptr; capacity = 0; }
646 if (cudaMalloc(&ptr, bytes) != cudaSuccess)
return false;
654 Stage* addRawStage(Stage* stage);
656 struct OutputBuffer {
659 size_t allocated_size;
663 std::vector<OutputBuffer> getOutputBuffers()
const;
665 static void* loadCompressedData(
666 const std::string& filename,
667 const FZMFileHeader& header,
668 cudaStream_t stream = 0,
669 MemoryPool* pool =
nullptr
673 std::pair<std::vector<Stage*>, std::vector<Stage*>> identifyTopology();
674 void setupInputBuffers(
const std::vector<Stage*>& sources);
675 int autoDetectUnconnectedOutputs();
676 void detectMultiOutputScenario(
int pipeline_outputs);
677 void configureStreamsIfNeeded();
680 void typeCheckConnections();
681 void computeInputAlignment();
682 void notifyStagesFinalizeHooks();
683 void refinePoolSize();
684 void setupGraphModeInput();
685 void preallocatePadBuffer();
686 void preallocateConcatBuffers();
690 std::pair<const void*, size_t> prepareInputSource(
691 const void* d_input,
size_t input_size, cudaStream_t stream);
698 void propagateBufferSizes(
bool force_from_current_inputs =
false);
700 std::vector<Stage*> getSourceStages()
const;
701 std::vector<Stage*> getSinkStages()
const;
706 struct FwdStageDesc {
708 std::vector<int> output_buf_ids;
709 std::vector<int> input_buf_ids;
713 using PipelineOutputMap = std::unordered_map<int, std::pair<void*, size_t>>;
716 void buildOrReuseInvCache(
717 const PipelineOutputMap& po_map,
720 cudaStream_t stream);
738 size_t caller_capacity,
742 cudaStream_t stream);
751 void buildStaticBufferMetadata();
758 std::vector<size_t> readConcatSegmentSizes(
759 const void* d_blob,
size_t n, cudaStream_t stream)
const;
763 static FZMFileHeader parseHeaderFromMemory(
const void* data,
size_t size);
764 static size_t computeFilePoolSize(
const FZMFileHeader& fh,
size_t pool_override_bytes);
765 static std::pair<std::vector<std::unique_ptr<Stage>>, std::vector<FwdStageDesc>>
766 reconstructForwardTopology(
const FZMFileHeader& fh);
767 static std::unordered_map<Stage*, size_t> buildSourceSizesFromHeader(
768 const FZMFileHeader& fh,
const std::vector<FwdStageDesc>& fwd_topology);
775 static std::pair<std::unique_ptr<CompressionDAG>,
776 std::unordered_map<Stage*, int>>
778 const std::vector<FwdStageDesc>& fwd_stages,
779 const PipelineOutputMap& pipeline_outputs,
782 const std::unordered_map<Stage*, size_t>& source_sizes,
783 bool enable_profiling
788 struct OutputBufferInfo {
792 std::string stage_name;
793 std::string output_name;
796 std::vector<OutputBufferInfo> collectOutputBuffers()
const;
799 size_t calculateConcatSize(
const std::vector<OutputBufferInfo>& outputs)
const;
801 size_t writeConcatBuffer(
802 const std::vector<OutputBufferInfo>& outputs,
803 uint8_t* d_concat_bytes,
807 void concatOutputs(
void** d_output,
size_t* output_size, cudaStream_t stream);
811 std::unique_ptr<MemoryPool> mem_pool_;
812 std::unique_ptr<CompressionDAG> dag_;
815 std::vector<std::unique_ptr<Stage>> stages_;
816 std::unordered_map<Stage*, DAGNode*> stage_to_node_;
818 struct ConnectionInfo {
821 std::string output_name;
824 std::vector<ConnectionInfo> connections_;
828 bool warmup_on_finalize_;
829 bool pool_managed_decomp_;
834 bool was_compressed_;
836 bool profiling_enabled_;
837 PipelinePerfResult last_perf_result_;
839 std::vector<DAGNode*> input_nodes_;
840 std::vector<DAGNode*> output_nodes_;
841 std::vector<int> input_buffer_ids_;
842 std::vector<int> output_buffer_ids_;
844 PoolBuffer d_concat_buffer_;
849 std::vector<void*> d_decomp_outputs_;
852 PinnedBuffer h_concat_header_;
854 PinnedBuffer h_copy_descs_;
855 DeviceBuffer d_copy_descs_;
861 std::vector<size_t> source_input_sizes_;
865 size_t input_alignment_bytes_;
866 PoolBuffer d_pad_buf_;
870 size_t original_input_size_;
872 size_t input_size_hint_;
873 float pool_multiplier_;
877 std::array<size_t, 3> dims_;
888 std::unique_ptr<CompressionDAG> inv_dag;
889 std::unordered_map<Stage*, int> inv_result_map;
890 std::unordered_map<int, int> fwd_to_inv_ext_buf;
891 std::unordered_map<Stage*, size_t> source_sizes;
893 std::unique_ptr<InvDAGCache> inv_cache_;
895 struct BufferMetadata {
898 size_t allocated_size;
903 std::vector<BufferMetadata> buffer_metadata_;
905 bool graph_mode_enabled_;
906 bool graph_captured_;
910 PoolBuffer d_graph_input_;
911 size_t d_graph_input_size_;
913 cudaGraph_t captured_graph_;
914 cudaGraphExec_t graph_exec_;