45 size_t input_data_size = 0,
47 float pool_multiplier = 3.0f
57 explicit Pipeline(
const std::string& config_path);
75 void setDims(
size_t x,
size_t y = 1,
size_t z = 1) { dims_ = {x, y, z}; }
76 void setDims(std::array<size_t, 3> dims) { dims_ = dims; }
77 std::array<size_t, 3> getDims()
const {
return dims_; }
85 template<
typename StageT,
typename... Args>
98 int connect(
Stage* dependent,
const std::vector<Stage*>& producers);
116 bool isWarmupOnFinalizeEnabled()
const {
return warmup_on_finalize_; }
125 bool isPoolManagedDecompOutput()
const {
return pool_managed_decomp_; }
162 return original_input_size_ > 0 ? original_input_size_ : input_size_;
183 fz::stream_t stream = 0
216 size_t output_buf_capacity,
217 size_t* actual_output_size,
218 fz::stream_t stream = 0
240 fz::stream_t stream = 0
267 size_t output_buf_capacity,
268 size_t* actual_output_size,
269 fz::stream_t stream = 0
313 size_t output_buf_capacity,
314 size_t* actual_output_size,
315 fz::stream_t stream = 0
345 void reset(fz::stream_t stream = 0);
354 bool isProfilingEnabled()
const {
return profiling_enabled_; }
382 bool isBoundsCheckEnabled()
const {
return dag_->isBoundsCheckEnabled(); }
390 coloring_enabled_ = enable;
391 dag_->setColoringEnabled(enable);
393 bool isColoringEnabled()
const {
return dag_->isColoringEnabled(); }
394 size_t getColorRegionCount()
const {
return dag_->getColorRegionCount(); }
415 std::unordered_map<std::string, std::vector<std::string>> notes;
416 for (
const auto& s : stages_) {
418 auto n = s->getRunNotes();
419 if (!n.empty()) notes.emplace(s->getName(), std::move(n));
434 bool isGraphModeEnabled()
const {
return graph_mode_enabled_; }
447 bool isGraphCaptured()
const {
return graph_captured_; }
449 size_t getPeakMemoryUsage()
const;
450 size_t getCurrentMemoryUsage()
const;
451 void printPipeline()
const;
458 std::vector<FZMStageInfo> stages;
459 std::vector<FZMBufferEntry> buffers;
463 void writeToFile(
const std::string& filename, fz::stream_t stream = 0);
525 const std::string& filename,
528 fz::stream_t stream = 0,
530 size_t pool_override_bytes = 0
551 const std::string& filename,
554 fz::stream_t stream = 0,
583 const void* header_bytes,
589 fz::stream_t stream = 0
628 ~PoolBuffer() { free(0); }
629 PoolBuffer() =
default;
630 PoolBuffer(
const PoolBuffer&) =
delete;
631 PoolBuffer& operator=(
const PoolBuffer&) =
delete;
633 void free(fz::stream_t s) {
634 if (ptr && pool) { pool->
free(ptr, s); ptr =
nullptr; capacity = 0; }
636 bool allocate(MemoryPool* p,
size_t bytes, fz::stream_t s,
637 const char* tag,
bool persistent =
false) {
640 ptr = pool->
allocate(bytes, s, tag, persistent);
641 if (ptr) capacity = bytes;
642 return ptr !=
nullptr;
647 struct PinnedBuffer {
651 ~PinnedBuffer() {
if (ptr) cudaFreeHost(ptr); }
652 PinnedBuffer() =
default;
653 PinnedBuffer(
const PinnedBuffer&) =
delete;
654 PinnedBuffer& operator=(
const PinnedBuffer&) =
delete;
657 bool ensureCapacity(
size_t bytes) {
658 if (capacity >= bytes)
return true;
659 if (ptr) { cudaFreeHost(ptr); ptr =
nullptr; capacity = 0; }
660 if (cudaHostAlloc(&ptr, bytes, cudaHostAllocDefault) != cudaSuccess)
return false;
667 struct DeviceBuffer {
671 ~DeviceBuffer() {
if (ptr) cudaFree(ptr); }
672 DeviceBuffer() =
default;
673 DeviceBuffer(
const DeviceBuffer&) =
delete;
674 DeviceBuffer& operator=(
const DeviceBuffer&) =
delete;
677 bool ensureCapacity(
size_t bytes) {
678 if (capacity >= bytes)
return true;
679 if (ptr) { cudaFree(ptr); ptr =
nullptr; capacity = 0; }
680 if (cudaMalloc(&ptr, bytes) != cudaSuccess)
return false;
688 Stage* addRawStage(Stage* stage);
690 struct OutputBuffer {
693 size_t allocated_size;
697 std::vector<OutputBuffer> getOutputBuffers()
const;
699 static void* loadCompressedData(
700 const std::string& filename,
701 const FZMFileHeader& header,
702 fz::stream_t stream = 0,
703 MemoryPool* pool =
nullptr
707 std::pair<std::vector<Stage*>, std::vector<Stage*>> identifyTopology();
708 void setupInputBuffers(
const std::vector<Stage*>& sources);
709 int autoDetectUnconnectedOutputs();
710 void detectMultiOutputScenario(
int pipeline_outputs);
711 void configureStreamsIfNeeded();
714 void typeCheckConnections();
715 void computeInputAlignment();
716 void notifyStagesFinalizeHooks();
717 void refinePoolSize();
718 void setupGraphModeInput();
719 void preallocatePadBuffer();
720 void preallocateConcatBuffers();
724 std::pair<const void*, size_t> prepareInputSource(
725 const void* d_input,
size_t input_size, fz::stream_t stream);
732 void propagateBufferSizes(
bool force_from_current_inputs =
false);
734 std::vector<Stage*> getSourceStages()
const;
735 std::vector<Stage*> getSinkStages()
const;
740 struct FwdStageDesc {
742 std::vector<int> output_buf_ids;
743 std::vector<int> input_buf_ids;
747 using PipelineOutputMap = std::unordered_map<int, std::pair<void*, size_t>>;
750 void buildOrReuseInvCache(
751 const PipelineOutputMap& po_map,
754 fz::stream_t stream);
772 size_t caller_capacity,
776 fz::stream_t stream);
785 void buildStaticBufferMetadata();
792 std::vector<size_t> readConcatSegmentSizes(
793 const void* d_blob,
size_t n, fz::stream_t stream)
const;
797 static FZMFileHeader parseHeaderFromMemory(
const void* data,
size_t size);
798 static size_t computeFilePoolSize(
const FZMFileHeader& fh,
size_t pool_override_bytes);
799 static std::pair<std::vector<std::unique_ptr<Stage>>, std::vector<FwdStageDesc>>
800 reconstructForwardTopology(
const FZMFileHeader& fh);
801 static std::unordered_map<Stage*, size_t> buildSourceSizesFromHeader(
802 const FZMFileHeader& fh,
const std::vector<FwdStageDesc>& fwd_topology);
809 static std::pair<std::unique_ptr<CompressionDAG>,
810 std::unordered_map<Stage*, int>>
812 const std::vector<FwdStageDesc>& fwd_stages,
813 const PipelineOutputMap& pipeline_outputs,
816 const std::unordered_map<Stage*, size_t>& source_sizes,
817 bool enable_profiling
822 struct OutputBufferInfo {
826 std::string stage_name;
827 std::string output_name;
830 std::vector<OutputBufferInfo> collectOutputBuffers()
const;
833 size_t calculateConcatSize(
const std::vector<OutputBufferInfo>& outputs)
const;
835 size_t writeConcatBuffer(
836 const std::vector<OutputBufferInfo>& outputs,
837 uint8_t* d_concat_bytes,
841 void concatOutputs(
void** d_output,
size_t* output_size, fz::stream_t stream);
845 std::unique_ptr<MemoryPool> mem_pool_;
846 std::unique_ptr<CompressionDAG> dag_;
849 std::vector<std::unique_ptr<Stage>> stages_;
850 std::unordered_map<Stage*, DAGNode*> stage_to_node_;
852 struct ConnectionInfo {
855 std::string output_name;
858 std::vector<ConnectionInfo> connections_;
862 bool warmup_on_finalize_;
863 bool pool_managed_decomp_;
868 bool was_compressed_;
870 bool profiling_enabled_;
873 bool coloring_enabled_ =
true;
874 PipelinePerfResult last_perf_result_;
876 std::vector<DAGNode*> input_nodes_;
877 std::vector<DAGNode*> output_nodes_;
878 std::vector<int> input_buffer_ids_;
879 std::vector<int> output_buffer_ids_;
881 PoolBuffer d_concat_buffer_;
886 std::vector<void*> d_decomp_outputs_;
889 PinnedBuffer h_concat_header_;
891 PinnedBuffer h_copy_descs_;
892 DeviceBuffer d_copy_descs_;
898 std::vector<size_t> source_input_sizes_;
902 size_t input_alignment_bytes_;
903 PoolBuffer d_pad_buf_;
907 size_t original_input_size_;
909 size_t input_size_hint_;
910 float pool_multiplier_;
914 std::array<size_t, 3> dims_;
925 std::unique_ptr<CompressionDAG> inv_dag;
926 std::unordered_map<Stage*, int> inv_result_map;
927 std::unordered_map<int, int> fwd_to_inv_ext_buf;
928 std::unordered_map<Stage*, size_t> source_sizes;
930 std::unique_ptr<InvDAGCache> inv_cache_;
932 struct BufferMetadata {
935 size_t allocated_size;
940 std::vector<BufferMetadata> buffer_metadata_;
942 bool graph_mode_enabled_;
943 bool graph_captured_;
947 PoolBuffer d_graph_input_;
948 size_t d_graph_input_size_;
950 fz::graph_t captured_graph_;
951 fz::graph_exec_t graph_exec_;