FZGPUModules 2.0
GPU-accelerated modular compression pipelines
Loading...
Searching...
No Matches
rre_stage.h
Go to the documentation of this file.
1#pragma once
2
28#include "stage/stage.h"
29#include "fzm_format.h"
30#include "backend/types.h"
31#include <cstdint>
32#include <cstring>
33#include <memory>
34#include <stdexcept>
35#include <string>
36#include <unordered_map>
37#include <vector>
38
39namespace fz {
40
55class RREStage : public Stage {
56public:
57 RREStage()
58 : is_inverse_(false)
59 , chunk_size_(16384)
60 , word_size_(1)
61 , actual_output_size_(0)
62 , cached_orig_bytes_(0)
63 , d_scratch_(nullptr)
64 , d_sizes_dev_(nullptr)
65 , d_clean_dev_(nullptr)
66 , d_dst_off_dev_(nullptr)
67 , scratch_capacity_(0)
68 {}
69
70 ~RREStage() override;
71
72 // ── Stage control ──────────────────────────────────────────────────────
73 void setInverse(bool inv) override { is_inverse_ = inv; }
74 bool isInverse() const override { return is_inverse_; }
75
79 bool isGraphCompatible() const override { return !is_inverse_; }
80
81 void setChunkSize(size_t bytes) { chunk_size_ = static_cast<uint32_t>(bytes); }
82 void setWordSize(size_t bytes) { word_size_ = static_cast<uint8_t>(bytes); }
83
84 size_t getChunkSize() const { return chunk_size_; }
85 size_t getRequiredInputAlignment() const override { return chunk_size_; }
86 int getWordSize() const { return static_cast<int>(word_size_); }
87
88 // Variable-length coder = the sink of a chunk-cooperative fused chain. Any
89 // byte-word chunk_size the fusion harness supports fuses (matches the fused
90 // RRECoder<ChunkBytes> op) — see chunk_geometry.h's kSupportedChunkBytes.
91 FusionSpec getFusionSpec() const override {
92 if (is_inverse_ || word_size_ != 1 ||
93 (chunk_size_ != 4096u && chunk_size_ != 8192u && chunk_size_ != 16384u)) return {};
94 return FusionSpec{FusionAccess::SegmentCodec, chunk_size_};
95 }
96
98 FusedOpDecl getFusedOp() const override {
99 if (!getFusionSpec().fusable()) return {};
100 return FusedOpDecl{FusionStrategy::ChunkCooperative, "RRECoder",
101 "fused/chunk_fusion/chunk_fusion.cuh", {}};
102 }
104 void setFusedArchiveResult(size_t archive_bytes, size_t orig_bytes) override {
105 setFusedResult(archive_bytes, orig_bytes);
106 }
110 void setFusedResult(size_t archive_bytes, size_t orig_bytes) {
111 actual_output_size_ = archive_bytes;
112 cached_orig_bytes_ = static_cast<uint32_t>(orig_bytes);
113 tail_readback_pending_ = false;
114 }
115 uint32_t getCachedOrigBytes() const { return cached_orig_bytes_; }
116
117 // ── Execution ──────────────────────────────────────────────────────────
119 fz::stream_t stream,
120 MemoryPool* pool,
121 const std::vector<void*>& inputs,
122 const std::vector<void*>& outputs,
123 const std::vector<size_t>& sizes
124 ) override;
125 void postStreamSync(fz::stream_t stream) override;
126
127 // ── Metadata ───────────────────────────────────────────────────────────
128 std::string getName() const override { return "RRE"; }
129 size_t getNumInputs() const override { return 1; }
130 size_t getNumOutputs() const override { return 1; }
131
132 std::vector<size_t> estimateOutputSizes(
133 const std::vector<size_t>& input_sizes
134 ) const override {
135 if (is_inverse_) {
136 if (cached_orig_bytes_ > 0)
137 return {static_cast<size_t>(cached_orig_bytes_)};
138 return {input_sizes.empty() ? 0 : input_sizes[0]};
139 }
140 // Forward: worst case = original data + stream header.
141 const size_t n_bytes = input_sizes.empty() ? 0 : input_sizes[0];
142 const size_t n_chunks = (n_bytes + chunk_size_ - 1) / chunk_size_;
143 const size_t hdr = 4 + 4 + 4 * n_chunks;
144 // postStreamSync()/getActualOutputSizesByName() always round the final
145 // size up to a 4-byte boundary and zero-fill the pad, even when the
146 // real total isn't already aligned (e.g. a partial final chunk stored
147 // raw at a byte count that isn't a multiple of 4) -- reserve that pad
148 // here too, or the caller's allocation is up to 3 bytes short and the
149 // memset in postStreamSync writes out of bounds.
150 const size_t worst = n_bytes + hdr;
151 return {(worst + 3) & ~size_t(3)};
152 }
153
154 std::unordered_map<std::string, size_t>
156 size_t getActualOutputSize(int index) const override;
157
167 const std::vector<size_t>& input_sizes
168 ) const override {
169 if (is_inverse_ || input_sizes.empty()) return 0;
170 const size_t in_bytes = input_sizes[0];
171 const size_t n_chunks = (in_bytes + chunk_size_ - 1) / chunk_size_;
172 return n_chunks * (static_cast<size_t>(chunk_size_) + 3 * sizeof(uint32_t));
173 }
174
175 uint16_t getStageTypeId() const override {
176 return static_cast<uint16_t>(StageType::RRE);
177 }
178
179 uint8_t getOutputDataType(size_t) const override {
180 return static_cast<uint8_t>(DataType::UINT8);
181 }
182
183 // ── Serialization ──────────────────────────────────────────────────────
185 size_t output_index, uint8_t* buf, size_t max_size
186 ) const override {
187 (void)output_index;
188 if (max_size < 9) return 0;
189 std::memcpy(buf, &chunk_size_, sizeof(uint32_t));
190 buf[4] = word_size_;
191 std::memcpy(buf + 5, &cached_orig_bytes_, sizeof(uint32_t));
192 return 9;
193 }
194
195 void deserializeHeader(const uint8_t* buf, size_t size) override {
196 if (size >= 4) std::memcpy(&chunk_size_, buf, sizeof(uint32_t));
197 if (size >= 5) word_size_ = buf[4];
198 if (size >= 9) std::memcpy(&cached_orig_bytes_, buf + 5, sizeof(uint32_t));
199 }
200
201 size_t getMaxHeaderSize(size_t) const override { return 9; }
202
203 void saveState() override {
204 saved_chunk_size_ = chunk_size_;
205 saved_word_size_ = word_size_;
206 saved_cached_orig_bytes_ = cached_orig_bytes_;
207 }
208
209 void restoreState() override {
210 chunk_size_ = saved_chunk_size_;
211 word_size_ = saved_word_size_;
212 cached_orig_bytes_ = saved_cached_orig_bytes_;
213 }
214
215private:
216 bool is_inverse_;
217 uint32_t chunk_size_;
218 uint32_t saved_chunk_size_ = 0;
219 uint8_t word_size_;
220 uint8_t saved_word_size_ = 0;
221 size_t actual_output_size_;
222 uint32_t cached_orig_bytes_ = 0;
223 uint32_t saved_cached_orig_bytes_ = 0;
224
225 // ── Persistent forward scratch buffers ───────────────────────────────────
226 uint8_t* d_scratch_;
227 uint32_t* d_sizes_dev_;
228 uint32_t* d_clean_dev_;
229 uint32_t* d_dst_off_dev_;
230 mutable bool tail_readback_pending_ = false;
231 mutable fz::stream_t tail_readback_stream_ = nullptr;
232 mutable uint32_t tail_last_index_ = 0;
233 mutable uint8_t* tail_output_ptr_ = nullptr;
234 size_t scratch_capacity_;
235 MemoryPool* scratch_pool_owner_ = nullptr;
236 bool scratch_from_pool_ = false;
238 std::weak_ptr<const void> scratch_alive_;
239};
240
241} // namespace fz
Definition mempool.h:82
Definition rre_stage.h:55
void setInverse(bool inv) override
Definition rre_stage.h:73
FusedOpDecl getFusedOp() const override
Chunk-cooperative coder op (the swappable variable-length sink). Stateless.
Definition rre_stage.h:98
void postStreamSync(fz::stream_t stream) override
size_t getMaxHeaderSize(size_t) const override
Definition rre_stage.h:201
void execute(fz::stream_t stream, MemoryPool *pool, const std::vector< void * > &inputs, const std::vector< void * > &outputs, const std::vector< size_t > &sizes) override
size_t estimateScratchBytes(const std::vector< size_t > &input_sizes) const override
Definition rre_stage.h:166
size_t getActualOutputSize(int index) const override
size_t getRequiredInputAlignment() const override
Definition rre_stage.h:85
uint8_t getOutputDataType(size_t) const override
Definition rre_stage.h:179
FusionSpec getFusionSpec() const override
Definition rre_stage.h:91
bool isGraphCompatible() const override
Definition rre_stage.h:79
std::vector< size_t > estimateOutputSizes(const std::vector< size_t > &input_sizes) const override
Definition rre_stage.h:132
void setFusedResult(size_t archive_bytes, size_t orig_bytes)
Definition rre_stage.h:110
void setFusedArchiveResult(size_t archive_bytes, size_t orig_bytes) override
Base-class tail hook → the existing coder result setter (archive, orig).
Definition rre_stage.h:104
void deserializeHeader(const uint8_t *buf, size_t size) override
Definition rre_stage.h:195
std::string getName() const override
Definition rre_stage.h:128
std::unordered_map< std::string, size_t > getActualOutputSizesByName() const override
size_t serializeHeader(size_t output_index, uint8_t *buf, size_t max_size) const override
Definition rre_stage.h:184
uint16_t getStageTypeId() const override
Definition rre_stage.h:175
void saveState() override
Definition rre_stage.h:203
Definition stage.h:31
FZM binary file format definitions — structs, enums, and helpers.
Definition dag.h:24
@ RRE
Repeated-word bitmap reducer with recursive bitmap compression (LC component)
Base class interface for all compression stages.
A stage's contribution to a generated fused kernel — the device-op it maps to, where its source lives...
Definition fusion.h:161
A stage's fusion contract. Stages that can participate in a fused kernel override Stage::getFusionSpe...
Definition fusion.h:51
Backend-neutral GPU type aliases.