FZGPUModules 2.0
GPU-accelerated modular compression pipelines
Loading...
Searching...
No Matches
compressor.h
Go to the documentation of this file.
1
5#pragma once
6
7#include "pipeline/dag.h"
8#include "pipeline/perf.h"
9#include "pipeline/config.h"
10#include "stage/stage.h"
11#include "stage/stage_factory.h"
12#include "mem/mempool.h"
13#include "fzm_format.h"
14
15#include <array>
16#include <memory>
17#include <stdexcept>
18#include <unordered_map>
19
20namespace fz {
21
34class Pipeline {
35public:
41 explicit Pipeline(
42 size_t input_data_size = 0,
44 float pool_multiplier = 3.0f
45 );
46
54 explicit Pipeline(const std::string& config_path);
55
56 ~Pipeline();
57
58 // ── Configuration ─────────────────────────────────────────────────────────
59
62
64 void setNumStreams(int num_streams);
65
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_; }
75
76 // ── Builder API ───────────────────────────────────────────────────────────
77
82 template<typename StageT, typename... Args>
83 StageT* addStage(Args&&... args);
84
92 int connect(Stage* dependent, Stage* producer, const std::string& output_name = "output");
93
95 int connect(Stage* dependent, const std::vector<Stage*>& producers);
96
102 void finalize();
103
109 void warmup(cudaStream_t stream = 0);
110
112 void setWarmupOnFinalize(bool enable) { warmup_on_finalize_ = enable; }
113 bool isWarmupOnFinalizeEnabled() const { return warmup_on_finalize_; }
114
121 void setPoolManagedDecompOutput(bool enable) { pool_managed_decomp_ = enable; }
122 bool isPoolManagedDecompOutput() const { return pool_managed_decomp_; }
123
137 size_t getMaxCompressedSize(size_t input_bytes) const;
138
158 size_t getLastUncompressedSize() const {
159 return original_input_size_ > 0 ? original_input_size_ : input_size_;
160 }
161
162 // ── Execution ─────────────────────────────────────────────────────────────
163
176 const void* d_input,
177 size_t input_size,
178 void** d_output,
179 size_t* output_size,
180 cudaStream_t stream = 0
181 );
182
210 const void* d_input,
211 size_t input_size,
212 void* d_output_buf,
213 size_t output_buf_capacity,
214 size_t* actual_output_size,
215 cudaStream_t stream = 0
216 );
217
233 const void* d_input,
234 size_t input_size,
235 void** d_output,
236 size_t* output_size,
237 cudaStream_t stream = 0
238 );
239
261 const void* d_input,
262 size_t input_size,
263 void* d_output_buf,
264 size_t output_buf_capacity,
265 size_t* actual_output_size,
266 cudaStream_t stream = 0
267 );
268
307 const void* d_input,
308 size_t input_size,
309 void* d_output_buf,
310 size_t output_buf_capacity,
311 size_t* actual_output_size,
312 cudaStream_t stream = 0
313 );
314
339 void prepareInverse(size_t uncompressed_size);
340
342 void reset(cudaStream_t stream = 0);
343
344 // ── Profiling ─────────────────────────────────────────────────────────────
345
350 void enableProfiling(bool enable);
351 bool isProfilingEnabled() const { return profiling_enabled_; }
352
354 const PipelinePerfResult& getLastPerfResult() const { return last_perf_result_; }
355
357 CompressionDAG* getDAG() { return dag_.get(); }
358
360 size_t getPoolThreshold() const;
361
372
378 void enableBoundsCheck(bool enable) { dag_->enableBoundsCheck(enable); }
379 bool isBoundsCheckEnabled() const { return dag_->isBoundsCheckEnabled(); }
380
386 void setColoringEnabled(bool enable) { dag_->setColoringEnabled(enable); }
387 bool isColoringEnabled() const { return dag_->isColoringEnabled(); }
388 size_t getColorRegionCount() const { return dag_->getColorRegionCount(); }
389
390 // ── CUDA Graph Capture (compression-only) ─────────────────────────────────
391
399 void enableGraphMode(bool enable);
400 bool isGraphModeEnabled() const { return graph_mode_enabled_; }
401
412 void captureGraph(cudaStream_t stream = 0);
413 bool isGraphCaptured() const { return graph_captured_; }
414
415 size_t getPeakMemoryUsage() const;
416 size_t getCurrentMemoryUsage() const;
417 void printPipeline() const;
418
419 // ── File Serialization ────────────────────────────────────────────────────
420
423 FZMHeaderCore core;
424 std::vector<FZMStageInfo> stages;
425 std::vector<FZMBufferEntry> buffers;
426 };
427
429 void writeToFile(const std::string& filename, cudaStream_t stream = 0);
430
432 static FZMFileHeader readHeader(const std::string& filename);
433
436
437 // ── In-memory metadata header (decode without a prior compress) ────────────
438
453 std::vector<uint8_t> serializeHeaderToMemory() const;
454
474 void primeInverseFromHeader(const void* header_bytes, size_t header_size);
475
491 const std::string& filename,
492 void** d_output,
493 size_t* output_size,
494 cudaStream_t stream = 0,
495 PipelinePerfResult* perf_out = nullptr,
496 size_t pool_override_bytes = 0
497 );
498
517 const std::string& filename,
518 void** d_output,
519 size_t* output_size,
520 cudaStream_t stream = 0,
521 PipelinePerfResult* perf_out = nullptr
522 );
523
549 const void* header_bytes,
550 size_t header_size,
551 const void* d_blob,
552 size_t blob_size,
553 void** d_output,
554 size_t* output_size,
555 cudaStream_t stream = 0
556 );
557
558 // ── Config File ───────────────────────────────────────────────────────────
559
573 void loadConfig(const std::string& path);
574
583 void saveConfig(const std::string& path) const;
584
585private:
586 // ── RAII buffer wrappers (private implementation detail) ─────────────────
587
588 // Pool-allocated persistent device buffer.
589 struct PoolBuffer {
590 void* ptr = nullptr;
591 size_t capacity = 0;
592 MemoryPool* pool = nullptr;
593
594 ~PoolBuffer() { free(0); }
595 PoolBuffer() = default;
596 PoolBuffer(const PoolBuffer&) = delete;
597 PoolBuffer& operator=(const PoolBuffer&) = delete;
598
599 void free(cudaStream_t s) {
600 if (ptr && pool) { pool->free(ptr, s); ptr = nullptr; capacity = 0; }
601 }
602 bool allocate(MemoryPool* p, size_t bytes, cudaStream_t s,
603 const char* tag, bool persistent = false) {
604 free(s);
605 pool = p;
606 ptr = pool->allocate(bytes, s, tag, persistent);
607 if (ptr) capacity = bytes;
608 return ptr != nullptr;
609 }
610 };
611
612 // cudaHostAlloc pinned host buffer — grows on demand, never shrinks.
613 struct PinnedBuffer {
614 void* ptr = nullptr;
615 size_t capacity = 0;
616
617 ~PinnedBuffer() { if (ptr) cudaFreeHost(ptr); }
618 PinnedBuffer() = default;
619 PinnedBuffer(const PinnedBuffer&) = delete;
620 PinnedBuffer& operator=(const PinnedBuffer&) = delete;
621
622 // Returns false on CUDA allocation failure.
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;
627 capacity = bytes;
628 return true;
629 }
630 };
631
632 // cudaMalloc device buffer — grows on demand, never shrinks.
633 struct DeviceBuffer {
634 void* ptr = nullptr;
635 size_t capacity = 0;
636
637 ~DeviceBuffer() { if (ptr) cudaFree(ptr); }
638 DeviceBuffer() = default;
639 DeviceBuffer(const DeviceBuffer&) = delete;
640 DeviceBuffer& operator=(const DeviceBuffer&) = delete;
641
642 // Returns false on CUDA allocation failure.
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;
647 capacity = bytes;
648 return true;
649 }
650 };
651
652 // ── Internal helpers ──────────────────────────────────────────────────────
653
654 Stage* addRawStage(Stage* stage);
655
656 struct OutputBuffer {
657 void* d_ptr;
658 size_t actual_size;
659 size_t allocated_size;
660 std::string name;
661 int buffer_id;
662 };
663 std::vector<OutputBuffer> getOutputBuffers() const;
664
665 static void* loadCompressedData(
666 const std::string& filename,
667 const FZMFileHeader& header,
668 cudaStream_t stream = 0,
669 MemoryPool* pool = nullptr
670 );
671
672 void validate();
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();
678
679 // finalize() sub-steps
680 void typeCheckConnections();
681 void computeInputAlignment();
682 void notifyStagesFinalizeHooks();
683 void refinePoolSize();
684 void setupGraphModeInput();
685 void preallocatePadBuffer();
686 void preallocateConcatBuffers();
687
688 // compress() helper: handles graph-mode copy or alignment padding.
689 // Returns the effective source pointer and padded source size.
690 std::pair<const void*, size_t> prepareInputSource(
691 const void* d_input, size_t input_size, cudaStream_t stream);
692
698 void propagateBufferSizes(bool force_from_current_inputs = false);
699
700 std::vector<Stage*> getSourceStages() const;
701 std::vector<Stage*> getSinkStages() const;
702
703 // ── Inverse DAG helpers ───────────────────────────────────────────────────
704
706 struct FwdStageDesc {
707 Stage* stage;
708 std::vector<int> output_buf_ids;
709 std::vector<int> input_buf_ids;
710 };
711
713 using PipelineOutputMap = std::unordered_map<int, std::pair<void*, size_t>>;
714
715 // decompress() helper: builds or reuses the inverse DAG cache.
716 void buildOrReuseInvCache(
717 const PipelineOutputMap& po_map,
718 Stage* src_stage,
719 size_t src_sz,
720 cudaStream_t stream);
721
734 void decompressCore(
735 const void* d_input,
736 size_t input_size,
737 void* caller_output,
738 size_t caller_capacity,
739 bool synchronize,
740 void** d_output,
741 size_t* output_size,
742 cudaStream_t stream);
743
751 void buildStaticBufferMetadata();
752
758 std::vector<size_t> readConcatSegmentSizes(
759 const void* d_blob, size_t n, cudaStream_t stream) const;
760
761 // decompressFromFile() helpers.
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);
769
775 static std::pair<std::unique_ptr<CompressionDAG>,
776 std::unordered_map<Stage*, int>>
777 buildInverseDAG(
778 const std::vector<FwdStageDesc>& fwd_stages,
779 const PipelineOutputMap& pipeline_outputs,
780 MemoryPool* pool,
781 MemoryStrategy strategy,
782 const std::unordered_map<Stage*, size_t>& source_sizes,
783 bool enable_profiling
784 );
785
786 // ── Concat helpers ────────────────────────────────────────────────────────
787
788 struct OutputBufferInfo {
789 int buffer_id;
790 void* d_ptr;
791 size_t actual_size;
792 std::string stage_name;
793 std::string output_name;
794 };
795
796 std::vector<OutputBufferInfo> collectOutputBuffers() const;
797
799 size_t calculateConcatSize(const std::vector<OutputBufferInfo>& outputs) const;
800
801 size_t writeConcatBuffer(
802 const std::vector<OutputBufferInfo>& outputs,
803 uint8_t* d_concat_bytes,
804 cudaStream_t stream
805 ) const;
806
807 void concatOutputs(void** d_output, size_t* output_size, cudaStream_t stream);
808
809 // ── Member variables ──────────────────────────────────────────────────────
810
811 std::unique_ptr<MemoryPool> mem_pool_;
812 std::unique_ptr<CompressionDAG> dag_;
813 MemoryStrategy strategy_;
814
815 std::vector<std::unique_ptr<Stage>> stages_;
816 std::unordered_map<Stage*, DAGNode*> stage_to_node_;
817
818 struct ConnectionInfo {
819 Stage* dependent;
820 Stage* producer;
821 std::string output_name;
822 int output_index;
823 };
824 std::vector<ConnectionInfo> connections_;
825
826 int num_streams_;
827 bool is_finalized_;
828 bool warmup_on_finalize_;
829 bool pool_managed_decomp_;
830
831 // is_compressed_: true after the first successful compress() (gates writeToFile).
832 // was_compressed_: true between compress() and the next reset() (gates captureGraph).
833 bool is_compressed_;
834 bool was_compressed_;
835
836 bool profiling_enabled_;
837 PipelinePerfResult last_perf_result_;
838
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_;
843
844 PoolBuffer d_concat_buffer_;
845 bool needs_concat_;
846
847 // Pool-persistent decompress output buffers (one per source stage).
848 // Only used when pool_managed_decomp_ == true.
849 std::vector<void*> d_decomp_outputs_;
850
851 // Pinned host buffer for concat header (one H2D copy instead of N).
852 PinnedBuffer h_concat_header_;
853 // Persistent pinned host + device descriptor buffers for the gather kernel.
854 PinnedBuffer h_copy_descs_;
855 DeviceBuffer d_copy_descs_;
856
857 size_t input_size_;
858
859 // Per-source input sizes from the most recent compress(), ordered to match
860 // input_nodes_. Used by decompress() to size each inverse result buffer.
861 std::vector<size_t> source_input_sizes_;
862
863 // Input alignment in bytes — LCM of all stage getRequiredInputAlignment() values.
864 // compress() zero-pads to this boundary transparently.
865 size_t input_alignment_bytes_;
866 PoolBuffer d_pad_buf_;
867
868 // Original (pre-padding) input size. decompress() uses this to trim the
869 // reported output back to what the caller provided. 0 when no padding.
870 size_t original_input_size_;
871
872 size_t input_size_hint_;
873 float pool_multiplier_;
874
875 // Dataset dimensions (x=fast, y, z). Pushed to each stage on addStage() and
876 // again at finalize(). Default {0,1,1} = 1-D, infer x from input size.
877 std::array<size_t, 3> dims_;
878
887 struct InvDAGCache {
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;
892 };
893 std::unique_ptr<InvDAGCache> inv_cache_;
894
895 struct BufferMetadata {
896 int buffer_id;
897 size_t actual_size;
898 size_t allocated_size;
899 std::string name;
900 DAGNode* producer;
901 int output_index;
902 };
903 std::vector<BufferMetadata> buffer_metadata_;
904
905 bool graph_mode_enabled_;
906 bool graph_captured_;
907
908 // Fixed device input buffer whose address is baked into the captured graph.
909 // compress() copies user input here before cudaGraphLaunch().
910 PoolBuffer d_graph_input_;
911 size_t d_graph_input_size_;
912
913 cudaGraph_t captured_graph_;
914 cudaGraphExec_t graph_exec_;
915};
916
917// ── Template implementation ───────────────────────────────────────────────────
918
919template<typename StageT, typename... Args>
920StageT* Pipeline::addStage(Args&&... args) {
921 if (is_finalized_) {
922 throw std::runtime_error("Cannot add stages after finalization");
923 }
924
925 auto stage_ptr = std::make_unique<StageT>(std::forward<Args>(args)...);
926 StageT* stage = stage_ptr.get();
927
928 stage->setDims(dims_);
929
930 DAGNode* node = dag_->addStage(stage, stage->getName());
931 size_t num_outputs = stage->getNumOutputs();
932 auto output_names = stage->getOutputNames();
933
934 // Pre-allocate all output slots as unconnected (size=1 placeholder).
935 // connect() will promote any that get wired to downstream stages.
936 for (size_t i = 0; i < num_outputs; i++) {
937 std::string out_name = i < output_names.size() ? output_names[i] : std::to_string(i);
938 dag_->addUnconnectedOutput(node, 1, i, stage->getName() + "." + out_name + "_unconnected");
939 }
940
941 stage_to_node_[stage] = node;
942 stages_.push_back(std::move(stage_ptr));
943 return stage;
944}
945
946} // namespace fz
Definition dag.h:92
Definition mempool.h:82
void free(void *ptr, cudaStream_t stream)
void * allocate(size_t size, cudaStream_t stream, const std::string &tag="", bool persistent=false)
Definition compressor.h:34
void setDims(size_t x, size_t y=1, size_t z=1)
Definition compressor.h:72
void decompressFromFileInstance(const std::string &filename, void **d_output, size_t *output_size, cudaStream_t stream=0, PipelinePerfResult *perf_out=nullptr)
static void decompressFromFile(const std::string &filename, void **d_output, size_t *output_size, cudaStream_t stream=0, PipelinePerfResult *perf_out=nullptr, size_t pool_override_bytes=0)
int connect(Stage *dependent, const std::vector< Stage * > &producers)
void setPoolManagedDecompOutput(bool enable)
Definition compressor.h:121
std::vector< uint8_t > serializeHeaderToMemory() const
size_t getPoolThreshold() const
void warmup(cudaStream_t stream=0)
void enableBoundsCheck(bool enable)
Definition compressor.h:378
bool isMemPoolFallbackMode() const
void saveConfig(const std::string &path) const
size_t getLastUncompressedSize() const
Definition compressor.h:158
void loadConfig(const std::string &path)
void setWarmupOnFinalize(bool enable)
Definition compressor.h:112
void compress(const void *d_input, size_t input_size, void *d_output_buf, size_t output_buf_capacity, size_t *actual_output_size, cudaStream_t stream=0)
CompressionDAG * getDAG()
Definition compressor.h:357
StageT * addStage(Args &&... args)
Definition compressor.h:920
void decompress(const void *d_input, size_t input_size, void *d_output_buf, size_t output_buf_capacity, size_t *actual_output_size, cudaStream_t stream=0)
const PipelinePerfResult & getLastPerfResult() const
Definition compressor.h:354
Pipeline(const std::string &config_path)
void writeToFile(const std::string &filename, cudaStream_t stream=0)
int connect(Stage *dependent, Stage *producer, const std::string &output_name="output")
void decompressFromMemory(const void *header_bytes, size_t header_size, const void *d_blob, size_t blob_size, void **d_output, size_t *output_size, cudaStream_t stream=0)
Pipeline(size_t input_data_size=0, MemoryStrategy strategy=MemoryStrategy::MINIMAL, float pool_multiplier=3.0f)
void setMemoryStrategy(MemoryStrategy strategy)
void prepareInverse(size_t uncompressed_size)
static FZMFileHeader readHeader(const std::string &filename)
void setColoringEnabled(bool enable)
Definition compressor.h:386
size_t getMaxCompressedSize(size_t input_bytes) const
void finalize()
void primeInverseFromHeader(const void *header_bytes, size_t header_size)
void enableGraphMode(bool enable)
void enableProfiling(bool enable)
FZMFileHeader buildHeader() const
void setNumStreams(int num_streams)
void captureGraph(cudaStream_t stream=0)
void reset(cudaStream_t stream=0)
void decompressInto(const void *d_input, size_t input_size, void *d_output_buf, size_t output_buf_capacity, size_t *actual_output_size, cudaStream_t stream=0)
void decompress(const void *d_input, size_t input_size, void **d_output, size_t *output_size, cudaStream_t stream=0)
void compress(const void *d_input, size_t input_size, void **d_output, size_t *output_size, cudaStream_t stream=0)
Definition stage.h:30
TOML-based pipeline configuration file support.
Compression DAG wiring, execution, and memory strategy types.
FZM binary file format definitions — structs, enums, and helpers.
Stream-ordered CUDA memory pool for pipeline buffer management.
Definition fzm_format.h:25
MemoryStrategy
Definition dag.h:23
@ MINIMAL
Allocate on-demand, free at last consumer. Lowest peak memory.
Pipeline and per-stage profiling result types.
Base class interface for all compression stages.
Factory function for reconstructing pipeline stages from serialized FZM headers.
Definition dag.h:52
Fixed-size FZM file header core (80 bytes).
Definition fzm_format.h:224
Definition perf.h:78
Definition compressor.h:422