A source-code guide
How Rumi Reads
This guide begins with no assumed knowledge of Rumi, Karu, OpenZL, GeoZL, or C++. It follows one real rumi.read() call until the selected samples become a framework tensor.
Four responsibilities make one read
Rumi is a raster format and reader. It understands array dimensions, tiles, bands, time positions, and output layout. It delegates byte transport and frame decompression to smaller libraries with narrower jobs.
The orchestrator. It maps a tensor selection to compressed frames and maps decoded samples to an output array.
A read-only transport library. Given a locator, offset, and length, it retrieves those bytes from memory, a local file, HTTP, or supported object storage.
OpenZL is the graph-based compression runtime. GeoZL extends that runtime with raster-aware codec nodes registered into the same decoder context.
Six nouns used throughout the guide
The bytes of one .rumi file, supplied as memory, a path, or a URI.
The small external index that tells Rumi the array shape, frame layout, and compressed byte counts.
The requested time positions, bands, spatial window, and output axis pattern.
One independently compressed unit. Rumi fetches and decodes whole frames, then keeps only requested samples.
Rumi’s in-memory work record: where a frame lives, how large it is, and where its samples must go.
Karu’s report that one requested range finished, including its status, byte count, buffer, and caller tag.
read() is a function. TransportSource::read() is a method owned by the C++ type TransportSource. A name such as rumi_source* is an opaque native handle from Python’s point of view: Python can retain or release it but does not inspect the C++ object behind it.
One read, seen from end to end
Suppose a model asks for one spatial window, one time position, and two bands. The public call is small because the implementation carries those choices through every layer.
chip = rumi.read(
"scene.rumi",
header,
time=[2],
bands=[1, 3],
window=(512, 256, 256, 256),
framework="torch",
)
The external header tells Rumi which frames contain the request and where their bytes begin. It does not perform I/O or decoding. The Spec defines that arithmetic; this page explains the functions that execute it.
Turn a user request into native values
rumi.read() is synchronous: it returns a finished array or raises an exception. Its arguments describe three different things—the bytes to read, the samples to keep, and the array to return.
| Argument | Meaning | What the code does with it |
|---|---|---|
source | A bytes-like object, local path, or remote URI | Wraps it as MemorySource or TransportSource |
header | The external index returned by write() or rebuilt by info() | Parses shape, dtype, tile layout, frame axes, and compressed sizes |
time, bands | Zero-based lists or half-open (start, stop) ranges; None means all | Validates them and converts them to the C API’s one-based indices |
window | (row, column, height, width); None means the whole image | Checks bounds and converts it to offsets plus sizes |
pattern | The order of output axes such as band, time, row, and column | Compiles shape and byte strides before any frame runs |
framework | numpy, torch, jax, tensorflow, or dlpack | Chooses an exact consumer; dlpack leaves a RumiArray |
The data file and header are separate on purpose. Rumi can plan exact byte ranges from the small header before touching the potentially large data source. A caller that already has the header can therefore begin a remote read without first downloading metadata from the data object.
read() in bindings/python/rumi/_read.py parses the external header, validates the selection and resolves the framework consumer before opening the source. The native handle wrappers live in bindings/python/rumi/_native.py; the CFFI declarations and library loading remain isolated in _ffi.py.
spec = _Spec(header)
consumer = resolve_framework(framework, spec.fields.dtype)
times = _resolve_axis(time, "time", spec.fields.time_count)
bands = _resolve_axis(bands, "bands", spec.fields.samples_per_pixel)
window = _resolve_window(window, spec.fields.image_length, spec.fields.image_width)
pattern = _resolve_pattern(pattern)
arr = _read_one(_Source(source), spec, times, bands, window, pattern)
return consumer.convert(arr)
_Source opens anything. Only memory sources bypass Karu. A local path and a remote URI both enter through TransportSource; Karu resolves which transport they need._resolve_axis() expands a half-open range when necessary, validates Python’s zero-based band and time indices, and converts them to the C API’s one-based convention. A null pointer plus count zero means “all.” _resolve_window() turns (row, column, height, width) into the row and column offsets and sizes expected by C.
_Source owns the native source handle. For in-memory input, the C++ MemorySource only borrows the Python buffer, so _Source._keep keeps that buffer alive until the native source is released. _Spec calls rumi_spec_parse(), arranges for the native object to be freed with ffi.gc(), and exposes the parsed fields needed for Python-side bounds checks.
Python rejects mistakes with familiar TypeError and ValueError messages. The C API validates again because it is a public boundary that can also be called without Python. The core then checks format-dependent invariants near the code that uses them.
Share a stable C boundary, then allocate exactly once
CFFI, Python’s C foreign-function bridge, calls the public C API in core/src/capi.cpp. That API deliberately exposes plain handles, integers, pointers, and status codes—not C++ classes or exceptions.
SelectionSetup carries the selected one-based axes, compiled LayoutPlan, and exact output byte count. The core receives a destination that is already valid and large enough.prepare_selection() centralizes pointer/count checks, selected-axis resolution, output-layout compilation, and overflow-safe size calculation. Both rumi_read(), where the caller supplies storage, and rumi_read_dlpack(), where Rumi allocates storage, use this same preparation. The two entry points cannot silently disagree about shape or required bytes.
The read path allocates exactly need bytes, invokes read_window(), and only then wraps the allocation with finish_dlpack(). Every failure path frees the allocation. The dtype registry supplies the decoded byte width and exact DLPack code independently from the logical width stored in the file.
capi_call() catches C++ exceptions and records a Rumi status plus error text. Back in Python, _check() maps that status to an appropriate Python exception. No C++ exception or partially initialized tensor crosses the C ABI.
Compile the selection into frame tasks
With shape and storage settled, the core turns the logical selection into one FrameTask for every compressed frame that can contribute samples.
FrameTask is the unit shared by planning, transport, decoding, and placement. Its tag lets a Karu completion return to the exact task that requested it.The header’s frame_unit records a choice made when the file was written: a frame may hold several bands or times, or those axes may select separate frames. append_read_plan() walks an axis only when that axis changes the frame index. Selected planes that live together remain members of one task, avoiding a separate fetch for every plane inside that frame.
The function walks only the tile rows and columns intersecting the window. For each relevant band and time combination it calls Header::frame_index(), then obtains the compressed range with frame_offset() and frame_byte_count(). It also precomputes where every selected plane belongs in the output.
idx = h.frame_index(tile_row, tile_col, band, time)
task.offset = h.frame_offset(idx)
task.compressed_size = h.frame_byte_count(idx)
When a complete decoded frame maps byte-for-byte to one contiguous output region, the task stores a direct destination. Otherwise it stores crop coordinates and strides for a later copy_rect(). No transport decision has been made yet.
The tile grid, frame-unit rules, frame-index formula, and prefix-sum offset reconstruction are documented with their diagrams in the Rumi Spec. The implementation point on this page is that those calculations become concrete FrameTask records before any frame bytes are requested.
Ask for byte ranges, not raster concepts
Karu is a read-only transport library. Rumi gives it a locator plus offset and length; Karu returns exactly that byte range. It does not know about bands, tiles, frames, or output arrays.
A locator is Karu’s transport-independent description of where bytes live. karu_resolve() turns a local path, HTTP URL, or supported cloud URI into that description once when TransportSource::open() runs. Reads are positional rather than cursor-based, so independent frame tasks can request different parts of the same object concurrently.
karu_resolve(). At execution time, remote_locator() decides whether the task joins the remote batch.execute_plan() asks each source for remote_locator(). A memory source returns no locator and later performs a bounded memcpy. A local TransportSource also joins local_tasks; its Rumi worker calls the blocking convenience function karu_client_fetch(). Only a genuinely remote locator enters the asynchronous batch route.
For remote work, Rumi creates one karu_req per frame and stores &task in its caller-controlled tag. It submits the complete remote workload before decoding local tasks. Seeing every range at once lets Karu schedule and coalesce compatible requests while Rumi’s CPU workers remain busy.
karu_batch_next() returns completions as they become ready. Rumi checks the Karu status and exact byte count, recovers the FrameTask from the tag, attaches the compressed buffer, and sends the ready tasks to the decoder in waves. Karu’s I/O threads continue filling the next wave while Rumi decodes the current one.
The locator, byte offset, byte count, destination or returned buffer, status, and caller tag.
Why that range is a frame, how it must decode, and where its selected samples belong.
read_items() creates a TransportSession for the operation. The session lazily obtains a Karu client only when transport is needed; a thread-local cache can reuse that client on later reads when its configuration still matches. Pure memory reads never start the transport runtime.
The API is synchronous; the engine is concurrent
There is no await rumi.read(...). The Python call blocks until every requested sample is ready. The asynchronous work lives underneath that call, where Karu can progress remote I/O while Rumi workers read local frames and decode completed ones.
decode_tasks(local_tasks) uses the Rumi pool. Ready remote frames are then decoded in waves while Karu keeps filling the next wave.Which thread owns what?
Builds the plan, submits the remote batch, waits for Karu completions, groups currently ready frames, and returns the final result.
Own active remote transfers, connection reuse, retries, and completion delivery, four per client by default (KARU_IO_THREADS). They keep running while Rumi decodes.
When enabled, runs execute_task(): positional local or memory reads when needed, OpenZL decompression, validation, and output placement.
Keep blocking local-file and cloud-credential preparation away from the remote event loops.
karu_batch_next(batch, -1) blocks the coordinator until the first completion is available. Rumi then calls karu_batch_next(batch, 0) repeatedly to drain everything already ready without waiting. That group becomes one decoder wave. The loop repeats until Karu reports the batch ended.
The Rumi thread pool is process-wide and is created only when a plan has more than one frame task. Its default size is one unless RUMI_NUM_THREADS or rumi.set_num_threads() changes it; without a pool, Executor::run() performs the same tasks serially on the calling thread. The first genuinely parallel read fixes the pool size, so configuration must happen before that read. Each call submits its own ThreadPool::Batch and waits only for that batch.
append_read_plan() fixes every destination before execution. Workers take one output-plane tile row at a time, keeping frames that touch the same pages together. Remote completions may arrive in any order; their write-group IDs restore the grouping before decode.
Use one OpenZL runtime extended by GeoZL
OpenZL is the compression runtime that understands the graph stored inside a compressed frame. GeoZL is an extension library that adds raster-oriented transform nodes to that runtime. They are one decoding path, not two competing decoders.
An OpenZL frame carries the graph and parameters needed to reconstruct its values. Here, a graph is a recipe of connected transform nodes: the output of one operation becomes the input of another. Built-in nodes implement general operations; GeoZL contributes nodes designed for geospatial raster data. Rumi does not choose or rebuild that graph during a read—it supplies the compressed frame to OpenZL and validates the result against the Rumi header.
struct WorkerState {
ZL_DCtx* dctx = ZL_DCtx_create();
WorkerState() {
geozl_register_decoders(dctx);
}
};
geozl_register_decoders(dctx) is the integration point. After registration, OpenZL can execute graphs that contain both built-in nodes and GeoZL custom raster nodes.WorkerState in core/src/plan.cpp is declared thread_local: every worker thread receives its own instance the first time it decodes. The worker can reuse its OpenZL context plus compressed and decoded scratch buffers for later tasks, but no two threads share that mutable decoder state. Frames remain independent; reuse avoids rebuilding the runtime and reallocating scratch for every frame.
execute_task() calls ZL_DCtx_decompressTyped(), then checks that OpenZL produced numeric samples with the expected fixed width and exact decoded byte size. Boolean frames receive one additional check: every decoded byte must be zero or one. Only a validated frame can reach output placement.
If OpenZL reports an unknown custom transform, Rumi asks geozl_owns_ctid() whether that identifier belongs to GeoZL. The error can then distinguish “update GeoZL” from an unknown third-party OpenZL codec.
Decode directly or copy the requested rectangle
The output buffer and all of its strides were fixed before the plan ran. A completed frame therefore already knows whether it can decode directly into its final location or needs a temporary decoded frame.
copy_rect() copies complete rows when both layouts are contiguous and falls back to sample-wise copies when the requested output pattern uses a different pixel stride.The output pattern determines shape and strides. For example, "t y x b" places time first and band last; another order changes where a task must write without changing which source frames it needs. This is why compile_layout() runs before append_read_plan().
rumi_read_dlpack() allocates the exact output size before calling the core. finish_dlpack() wraps that allocation in a versioned DLPack tensor. DLPack is a small ownership-and-metadata protocol shared by array libraries: it describes the address, shape, strides, dtype, and device without copying the decoded array into a second framework-owned buffer.
Python wraps the native tensor in RumiArray. Calling numpy(), torch(), jax(), or tensorflow() transfers that allocation through DLPack. RumiArray.__dlpack__() permits the transfer exactly once; if no consumer takes it, the wrapper frees the native tensor itself. Framework compatibility is checked before any source is opened.
Output placement does not depend on completion order. Every task already carries its destination offsets, so a fast remote frame and a local frame can finish in either order without changing the returned array.
Read several images in one plan
read_many() does not loop over read(). It pairs each source with its header and, when supplied, its window, then sends the whole batch through one read_items() call and one execution plan. Without windows, it reads every image in full.
batch = rumi.read_many(
["north.rumi", "south.rumi"],
[north_header, south_header],
windows=[(0, 0, 256, 256), (512, 0, 256, 256)],
bands=[1, 3],
framework="torch",
)
k * n_stride. The final tensor preserves input order even though its frames may complete in another order.compatible_headers() requires the same tile size, band count, dtype, time count, and frame-axis behavior. A batch without windows also requires equal image dimensions. With explicit windows, the source images may differ in size, but every window must have the same height and width so the results fit one rectangular tensor.
The combined plan is the performance feature. Remote frame requests from every item are visible to one Karu submission, local and memory work share the same Rumi pool, and every ready frame can decode without waiting for earlier items to finish.
The read path by file
These are the implementation files to open in order when changing or debugging a read.
| File | What it owns | Main symbols |
|---|---|---|
bindings/python/rumi/_read.py | Public argument validation and read orchestration | read, read_many, _read_one |
bindings/python/rumi/_framework.py | Framework compatibility and exact DLPack transfer | resolve_framework, _FRAMEWORKS |
bindings/python/rumi/_dlpack.py | DLPack capsule and native tensor ownership | RumiArray, __dlpack__ |
bindings/python/rumi/_native.py | Native source and parsed-header handles | _Source, _Spec |
bindings/python/rumi/_ffi.py | C declarations, library loading, ABI checks, and status mapping | ffi, lib, _check |
core/src/capi.cpp | C validation, output layout, allocation, and DLPack wrapping | prepare_selection, rumi_read_dlpack, prepare_many |
core/src/read.cpp | Frame planning, source checks, Karu batching, and execution order | append_read_plan, read_items, execute_plan |
core/src/source.cpp | Memory sources, Karu locator resolution, and transport sessions | MemorySource, TransportSource, TransportSession |
core/src/plan.cpp | Worker-local OpenZL context, GeoZL decoder registration, frame validation, and placement | WorkerState, execute_task, copy_rect |
core/include/rumi/thread_pool.hpp | Process-wide workers and operation-scoped batches | ThreadPool, ThreadPool::Batch |
karu/src/runtime/engine.cpp | Remote event loops, local-file workers, credential workers, and completion queues | Engine::submit, io_loop, file_loop, credential_loop |
geozl/core/src/decoder_registry.c | Registration of GeoZL’s raster codecs into an OpenZL decoder context | geozl_register_decoders, geozl_owns_ctid |
openzl/decompress/decompress2.c | OpenZL graph decompression reached through its public decoder API | ZL_DCtx_decompressTyped |
If the wrong frames are requested, begin in append_read_plan(). If ranges fail, begin in execute_plan() and source.cpp. If compressed bytes arrive but decoding or placement fails, begin in execute_task(). If the native result is correct but the Python object fails, begin in RumiArray.
Use the Rumi Spec for the binary layout and frame-index formulas, or continue to Examples for public API usage.