Skip to content

Commit 38d2171

Browse files
committed
updates API surface with clearer ownership rules and removes bloated compress/decompress versions. also adds cli tool and updates tests.
1 parent a7eda3d commit 38d2171

24 files changed

Lines changed: 1218 additions & 1492 deletions

‎include/pipeline/compressor.h‎

Lines changed: 111 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -110,31 +110,41 @@ class Pipeline {
110110
bool isWarmupOnFinalizeEnabled() const { return warmup_on_finalize_; }
111111

112112
/**
113-
* When true, decompress() returns a pool-owned pointer (do NOT cudaFree).
113+
* When true (default), decompress() returns a pool-owned pointer (do NOT cudaFree).
114114
* Valid until the next decompress() call or Pipeline destruction.
115-
* When false (default), decompress() returns a freshly cudaMalloc'd pointer
115+
* When false, decompress() returns a freshly cudaMalloc'd pointer
116116
* that the caller must cudaFree().
117117
*/
118118
void setPoolManagedDecompOutput(bool enable) { pool_managed_decomp_ = enable; }
119119
bool isPoolManagedDecompOutput() const { return pool_managed_decomp_; }
120120

121-
// ── Execution ─────────────────────────────────────────────────────────────
122-
123121
/**
124-
* Input specification for multi-source compression.
125-
* @param source A DAG root stage already added to this pipeline.
126-
* @param d_data Device pointer to raw input for this source.
127-
* @param size Size of d_data in bytes.
122+
* Return the worst-case compressed output size in bytes for the given input.
123+
*
124+
* Must be called after finalize(). Use this to pre-allocate a caller-owned
125+
* output buffer before passing it to the user-owned compress() overload.
126+
*
127+
* The returned value is a tight upper bound derived from each stage's
128+
* estimateOutputSizes() chain — it should rarely exceed ~110% of the actual
129+
* compressed size for typical data.
130+
*
131+
* @param input_bytes Number of bytes you intend to compress.
132+
* @throws std::runtime_error if the pipeline is not yet finalized.
128133
*/
129-
struct InputSpec {
130-
Stage* source;
131-
const void* d_data;
132-
size_t size;
133-
};
134+
size_t getMaxCompressedSize(size_t input_bytes) const;
135+
136+
// ── Execution ─────────────────────────────────────────────────────────────
134137

135138
/**
136-
* Compress (single-source). Throws if the pipeline has more than one source.
137-
* *d_output is pool-owned — do NOT cudaFree.
139+
* Compress (pool-owned output). The pool retains the output buffer.
140+
*
141+
* @param d_input Device pointer to raw input data.
142+
* @param input_size Size of d_input in bytes.
143+
* @param d_output Receives a pool-owned pointer to the compressed output.
144+
* Do NOT call cudaFree() — valid until the next compress(),
145+
* reset(), or Pipeline destruction.
146+
* @param output_size Receives the exact compressed size in bytes.
147+
* @param stream CUDA stream for all GPU operations.
138148
*/
139149
void compress(
140150
const void* d_input,
@@ -145,39 +155,53 @@ class Pipeline {
145155
);
146156

147157
/**
148-
* Compress (multi-source). One InputSpec per source stage; order does not matter.
149-
* *d_output is pool-owned — do NOT cudaFree.
158+
* Compress (user-owned output). The compressed data is written into the
159+
* caller-provided device buffer.
160+
*
161+
* The buffer just needs to be large enough for the actual compressed output
162+
* of this specific call — which depends on the data. If the actual output
163+
* exceeds `output_buf_capacity` a `std::runtime_error` is thrown with the
164+
* actual and capacity sizes so the caller can retry with a larger buffer.
165+
*
166+
* Use `getMaxCompressedSize(input_bytes)` for a guaranteed safe upper bound.
167+
* Alternatively, if you know empirically that your data compresses to at most
168+
* X bytes for your workload, you can pass X directly and accept the small
169+
* risk of a runtime error on unusually incompressible inputs.
170+
*
171+
* Incompatible with CUDA Graph mode (the output address cannot be baked into
172+
* a captured graph). Throws if enableGraphMode(true) was set.
173+
*
174+
* @param d_input Device pointer to raw input data.
175+
* @param input_size Size of d_input in bytes.
176+
* @param d_output_buf Caller-allocated device buffer to write compressed
177+
* data into.
178+
* @param output_buf_capacity Capacity of d_output_buf in bytes. Must fit the
179+
* actual compressed output for this call.
180+
* @param actual_output_size Receives the exact compressed bytes written.
181+
* @param stream CUDA stream for all GPU operations.
150182
*/
151183
void compress(
152-
const std::vector<InputSpec>& inputs,
153-
void** d_output,
154-
size_t* output_size,
184+
const void* d_input,
185+
size_t input_size,
186+
void* d_output_buf,
187+
size_t output_buf_capacity,
188+
size_t* actual_output_size,
155189
cudaStream_t stream = 0
156190
);
157191

158192
/**
159-
* Per-source input size hint for multi-source pipelines. Call after addStage()
160-
* but before finalize() to seed propagateBufferSizes() for accurate estimates.
161-
*/
162-
void setInputSizeHint(Stage* source, size_t size) {
163-
per_source_hints_[source] = size;
164-
}
165-
166-
/**
167-
* Decompress (single-source). Inverse of compress().
193+
* Decompress. Inverse of compress().
168194
*
169195
* @param d_input nullptr to read from the forward DAG's live buffers
170196
* (simplest path, valid immediately after compress()).
171197
* Non-null for an external compressed buffer.
172-
* @param input_size Byte size of `d_input` (ignored when `d_input` is nullptr).
198+
* @param input_size Byte size of d_input (ignored when d_input is nullptr).
173199
* @param d_output Receives the decompressed device pointer.
174-
* Ownership depends on setPoolManagedDecompOutput().
200+
* Ownership depends on setPoolManagedDecompOutput():
201+
* false → caller-owned, must cudaFree.
202+
* true (default) → pool-owned, do NOT cudaFree.
175203
* @param output_size Receives the exact decompressed size in bytes.
176204
* @param stream CUDA stream for all GPU operations.
177-
*
178-
* For multi-source pipelines, `*d_output` is a concat buffer:
179-
* `[num_bufs:u32][size1:u64][data1][size2:u64][data2]...`
180-
* Use decompressMulti() to receive individual per-source buffers.
181205
*/
182206
void decompress(
183207
const void* d_input,
@@ -188,14 +212,32 @@ class Pipeline {
188212
);
189213

190214
/**
191-
* Decompress (multi-source). Returns one {device_ptr, size} pair per source,
192-
* in the same order as forward source discovery. Ownership follows
193-
* setPoolManagedDecompOutput().
215+
* Decompress into a caller-provided device buffer (user-owned output).
216+
*
217+
* The decompressed data is written directly into d_output_buf. No cudaMalloc
218+
* or pool allocation is performed — the caller owns the buffer entirely.
219+
*
220+
* The buffer just needs to be large enough for the actual decompressed output
221+
* of this call. If it is too small a `std::runtime_error` is thrown with the
222+
* actual size so the caller can retry. Typically the uncompressed size is
223+
* known from the file header (`FZMHeaderCore::uncompressed_size`) or from
224+
* the original compress() call.
225+
*
226+
* @param d_input See decompress() above.
227+
* @param input_size See decompress() above.
228+
* @param d_output_buf Caller-allocated device buffer to receive
229+
* decompressed data.
230+
* @param output_buf_capacity Capacity of d_output_buf in bytes.
231+
* @param actual_output_size Receives the exact bytes written.
232+
* @param stream CUDA stream for all GPU operations.
194233
*/
195-
std::vector<std::pair<void*, size_t>> decompressMulti(
196-
const void* d_input = nullptr,
197-
size_t input_size = 0,
198-
cudaStream_t stream = 0
234+
void decompress(
235+
const void* d_input,
236+
size_t input_size,
237+
void* d_output_buf,
238+
size_t output_buf_capacity,
239+
size_t* actual_output_size,
240+
cudaStream_t stream = 0
199241
);
200242

201243
/** Free non-persistent buffers and reset execution state for re-use. */
@@ -287,6 +329,8 @@ class Pipeline {
287329
* One-shot decompress from an FZM file. Reconstructs the pipeline from the
288330
* file header, allocates a pool, and runs decompression.
289331
*
332+
* Output is always caller-owned (caller must cudaFree *d_output).
333+
*
290334
* @param filename Path to the `.fzm` file.
291335
* @param d_output Receives the decompressed device pointer (caller must `cudaFree`).
292336
* @param output_size Receives the decompressed size in bytes.
@@ -304,6 +348,31 @@ class Pipeline {
304348
size_t pool_override_bytes = 0
305349
);
306350

351+
/**
352+
* One-shot decompress from an FZM file (instance overload).
353+
*
354+
* Behaves identically to the static `decompressFromFile()` overload but
355+
* respects the setPoolManagedDecompOutput() flag on this instance:
356+
* false → caller must `cudaFree(*d_output)`.
357+
* true (default) → *d_output is pool-owned; do NOT `cudaFree`.
358+
*
359+
* The distinct name avoids overload-resolution ambiguity at call sites
360+
* that are not member functions.
361+
*
362+
* @param filename Path to the `.fzm` file.
363+
* @param d_output Receives the decompressed device pointer.
364+
* @param output_size Receives the decompressed size in bytes.
365+
* @param stream CUDA stream for all GPU operations.
366+
* @param perf_out Optional timing result.
367+
*/
368+
void decompressFromFileInstance(
369+
const std::string& filename,
370+
void** d_output,
371+
size_t* output_size,
372+
cudaStream_t stream = 0,
373+
PipelinePerfResult* perf_out = nullptr
374+
);
375+
307376
// ── Config File ───────────────────────────────────────────────────────────
308377

309378
/**
@@ -493,10 +562,6 @@ class Pipeline {
493562
size_t input_size_hint_;
494563
float pool_multiplier_;
495564

496-
// Per-source size hints (set via setInputSizeHint()). Override input_size_hint_
497-
// for the matching source during propagateBufferSizes().
498-
std::unordered_map<Stage*, size_t> per_source_hints_;
499-
500565
// Dataset dimensions (x=fast, y, z). Used by convenience.h addLorenzo() to
501566
// select 1-D/2-D/3-D automatically. Default {0,1,1} = 1-D, infer x from input.
502567
std::array<size_t, 3> dims_;

‎src/pipeline/compressor.cpp‎

Lines changed: 7 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ Pipeline::Pipeline(size_t input_data_size, MemoryStrategy strategy, float pool_m
1919
num_streams_(1),
2020
is_finalized_(false),
2121
warmup_on_finalize_(false),
22-
pool_managed_decomp_(false),
22+
pool_managed_decomp_(true),
2323
is_compressed_(false),
2424
was_compressed_(false),
2525
profiling_enabled_(false),
@@ -298,7 +298,7 @@ void Pipeline::finalize() {
298298
// (CUB temp storage, RZE per-chunk scratch, etc.).
299299
// Skip when there is no input size hint — buffer sizes will be 1-byte
300300
// placeholders and the topology-derived value would be meaningless.
301-
if (input_size_hint_ > 0 || !per_source_hints_.empty()) {
301+
if (input_size_hint_ > 0) {
302302
// computeTopoPoolSize() already accounts for all intermediate buffers
303303
// and persistent stage scratch, so it is the accurate peak requirement.
304304
// Apply a small safety margin (10%) for transient CUB allocations that
@@ -344,7 +344,7 @@ void Pipeline::finalize() {
344344
"Graph mode requires PREALLOCATE memory strategy. "
345345
"Call setMemoryStrategy(MemoryStrategy::PREALLOCATE) before finalize().");
346346
}
347-
if (input_size_hint_ == 0 && per_source_hints_.empty()) {
347+
if (input_size_hint_ == 0) {
348348
throw std::runtime_error(
349349
"Graph mode requires a non-zero input size hint. "
350350
"Pass input_data_size to the Pipeline constructor.");
@@ -718,21 +718,17 @@ void Pipeline::configureStreamsIfNeeded() {
718718
}
719719

720720
void Pipeline::propagateBufferSizes(bool force_from_current_inputs) {
721-
bool has_hint = (input_size_hint_ > 0) || !per_source_hints_.empty();
721+
bool has_hint = (input_size_hint_ > 0);
722722
if (!has_hint && !force_from_current_inputs) {
723723
FZ_LOG(DEBUG, "No input size hint, using placeholder buffer sizes");
724724
return;
725725
}
726726

727727
if (!force_from_current_inputs) {
728-
// Seed each source's external input buffer with its per-source hint,
729-
// falling back to the global constructor hint when none is set.
728+
// Seed the source's external input buffer with the constructor hint.
730729
for (size_t i = 0; i < input_nodes_.size(); i++) {
731-
Stage* src_stage = input_nodes_[i]->stage;
732-
auto it = per_source_hints_.find(src_stage);
733-
size_t hint = (it != per_source_hints_.end()) ? it->second : input_size_hint_;
734-
if (hint > 0) {
735-
dag_->updateBufferSize(input_buffer_ids_[i], hint);
730+
if (input_size_hint_ > 0) {
731+
dag_->updateBufferSize(input_buffer_ids_[i], input_size_hint_);
736732
}
737733
}
738734
}

0 commit comments

Comments
 (0)