• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

llnl-asr / dftracer-utils / 36502607218

29 Sep 2026 12:20AM UTC coverage: 57.996% (+5.4%) from 52.58%
36502607218

push

github

rayandrew
Merge branch 'feat/view-as-lazyframe' into 'develop'

feat: view as lazyframe

See merge request dftracer/dftracer-utils!25

74077 of 162805 branches covered (45.5%)

Branch coverage included in aggregate %.

3652 of 3978 new or added lines in 49 files covered. (91.8%)

3464 existing lines in 85 files now uncovered.

65902 of 78556 relevant lines covered (83.89%)

179668.29 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

61.52
/src/dftracer/utils/python/indexer.cpp
1
#include <dftracer/utils/core/runtime.h>
2
#include <dftracer/utils/core/tasks/coro_scope.h>
3
#include <dftracer/utils/python/indexer.h>
4
#include <dftracer/utils/python/py_errors.h>
5
#include <dftracer/utils/python/py_method.h>
6
#include <dftracer/utils/python/py_runtime_mixin.h>
7
#include <dftracer/utils/python/py_str_helpers.h>
8
#include <dftracer/utils/python/py_type_helpers.h>
9
#include <dftracer/utils/python/runtime.h>
10
#include <dftracer/utils/trace/internal/utils.h>
11
#include <dftracer/utils/utilities/indexer/index_builder_utility.h>
12
#include <dftracer/utils/utilities/indexer/index_database.h>
13
#include <dftracer/utils/utilities/indexer/internal/helpers.h>
14
#include <structmember.h>
15

16
#include <cstring>
17
#include <memory>
18

19
static void CheckpointIndexer_dealloc(CheckpointIndexerObject *self) {
26 ✔
20
    if (self->handle) {
26 ✔
21
        // The Python wrapper owns only the native indexer handle. The
22
        // underlying RocksDB instance remains manager-owned and may continue to
23
        // live process-wide for the same .dftindex path.
24
        dftu_indexer_destroy(self->handle);
10 ✔
25
        self->handle = NULL;
10 ✔
26
    }
5 ✔
27
    Py_XDECREF(self->gz_path);
26 ✔
28
    Py_XDECREF(self->index_path);
26 ✔
29
    Py_XDECREF(self->runtime_obj);
26 ✔
30
    Py_TYPE(self)->tp_free((PyObject *)self);
26 ✔
31
}
26 ✔
32

33
static void CheckpointIndexer_release_handle(CheckpointIndexerObject *self) {
14 ✔
34
    if (self->handle) {
14 ✔
35
        // Releasing the handle drops this wrapper's native indexer state only.
36
        // Shared RocksDB lifetime is managed separately by RocksDBManager.
37
        dftu_indexer_destroy(self->handle);
14 ✔
38
        self->handle = NULL;
14 ✔
39
    }
7 ✔
40
}
14 ✔
41

42
static PyObject *CheckpointIndexer_new(PyTypeObject *type, PyObject *,
16 ✔
43
                                       PyObject *) {
44
    CheckpointIndexerObject *self;
45
    self = (CheckpointIndexerObject *)type->tp_alloc(type, 0);
16 ✔
46
    if (self != NULL) {
16 ✔
47
        self->handle = NULL;
16 ✔
48
        self->gz_path = NULL;
16 ✔
49
        self->index_path = NULL;
16 ✔
50
        self->checkpoint_size = 0;
16 ✔
51
        self->build_bloom = 0;
16 ✔
52
        self->runtime_obj = NULL;
16 ✔
53
    }
8 ✔
54
    return (PyObject *)self;
16 ✔
55
}
56

57
static int CheckpointIndexer_init(CheckpointIndexerObject *self, PyObject *args,
16 ✔
58
                                  PyObject *kwds) {
59
    static const char *kwlist[] = {"gz_path",
60
                                   "index_path",
61
                                   "checkpoint_size",
62
                                   "force_rebuild",
63
                                   "build_bloom",
64
                                   "runtime",
65
                                   NULL};
66
    const char *gz_path;
67
    const char *index_path = NULL;
16 ✔
68
    std::uint64_t checkpoint_size =
16 ✔
69
        dftracer::utils::constants::indexer::DEFAULT_CHECKPOINT_SIZE;
70
    int force_rebuild = 0;
16 ✔
71
    int build_bloom = 0;
16 ✔
72
    PyObject *runtime_arg = NULL;
16 ✔
73

74
    if (!PyArg_ParseTupleAndKeywords(
16 !
75
            args, kwds, "s|snppO", const_cast<char **>(kwlist), &gz_path,
8 ✔
76
            &index_path, &checkpoint_size, &force_rebuild, &build_bloom,
77
            &runtime_arg)) {
UNCOV
78
        return -1;
×
79
    }
80

81
    if (runtime_arg && runtime_arg != Py_None) {
16 !
UNCOV
82
        if (PyObject_TypeCheck(runtime_arg, &RuntimeType)) {
×
83
            Py_INCREF(runtime_arg);
×
84
            self->runtime_obj = runtime_arg;
×
85
        } else {
UNCOV
86
            PyObject *native = PyObject_GetAttrString(runtime_arg, "_native");
×
87
            if (native && PyObject_TypeCheck(native, &RuntimeType)) {
×
88
                self->runtime_obj = native;
×
89
            } else {
90
                Py_XDECREF(native);
×
UNCOV
91
                PyErr_SetString(PyExc_TypeError,
×
92
                                "runtime must be a Runtime instance or None");
UNCOV
93
                return -1;
×
94
            }
95
        }
96
    }
97

98
    self->gz_path = PyUnicode_FromString(gz_path);
16 !
99
    if (!self->gz_path) {
16 ✔
UNCOV
100
        return -1;
×
101
    }
102

103
    if (index_path) {
16 ✔
104
        self->index_path = PyUnicode_FromString(index_path);
14 !
105
    } else {
7 ✔
106
        const std::string resolved_index_path =
107
            dftracer::utils::trace::internal::determine_index_path(gz_path, "");
3 !
108
        self->index_path = PyUnicode_FromString(resolved_index_path.c_str());
2 !
109
    }
2 ✔
110

111
    if (!self->index_path) {
16 ✔
UNCOV
112
        Py_DECREF(self->gz_path);
×
UNCOV
113
        return -1;
×
114
    }
115

116
    self->checkpoint_size = checkpoint_size;
16 ✔
117
    self->build_bloom = build_bloom;
16 ✔
118

119
    const char *index_path_str = as_utf8(self->index_path);
16 !
120
    if (!index_path_str) {
16 ✔
UNCOV
121
        return -1;
×
122
    }
123

124
    self->handle = dftu_indexer_create(gz_path, index_path_str, checkpoint_size,
24 !
125
                                       force_rebuild);
8 ✔
126
    if (!self->handle) {
16 ✔
127
        PyErr_SetString(PyExc_RuntimeError, "Failed to create indexer");
2 !
128
        return -1;
2 ✔
129
    }
130

131
    return 0;
14 ✔
132
}
8 ✔
133

134
static dftracer::utils::Runtime *get_indexer_runtime(
2 ✔
135
    CheckpointIndexerObject *self) {
136
    if (self->runtime_obj) {
2 !
UNCOV
137
        return ((RuntimeObject *)self->runtime_obj)->runtime.get();
×
138
    }
139
    return dftracer::utils::python::get_default_runtime();
2 ✔
140
}
1 ✔
141

142
static PyObject *CheckpointIndexer_build(CheckpointIndexerObject *self,
10 ✔
143
                                         PyObject *Py_UNUSED(ignored)) {
144
    if (!self->handle) {
10 ✔
UNCOV
145
        PyErr_SetString(PyExc_RuntimeError, "Indexer not initialized");
×
UNCOV
146
        return NULL;
×
147
    }
148

149
    // Use IndexBatchBuilderUtility when bloom is requested.
150
    // Otherwise, use the simpler dftu_indexer_build which only creates
151
    // checkpoints.
152
    if (self->build_bloom) {
10 ✔
153
        using namespace dftracer::utils;
154
        using namespace dftracer::utils::utilities::indexer;
155

156
        const char *gz = as_utf8(self->gz_path);
2 !
157
        const char *idx = as_utf8(self->index_path);
2 !
158
        if (!gz || !idx) {
2 !
UNCOV
159
            return NULL;
×
160
        }
161

162
        auto batch_config = std::make_shared<IndexBuildBatchConfig>();
2 !
163
        batch_config->file_paths.emplace_back(gz);
2 !
164
        batch_config->checkpoint_size =
3 ✔
165
            static_cast<std::size_t>(self->checkpoint_size);
2 ✔
166
        batch_config->parallelism = 1;
2 ✔
167
        batch_config->rebuild_root_summaries = true;
2 ✔
168

169
        std::string idx_str(idx);
2 !
170
        auto pos = idx_str.find_last_of('/');
2 ✔
171
        if (pos != std::string::npos) {
2 !
172
            batch_config->index_dir = idx_str.substr(0, pos);
2 !
173
        }
1 ✔
174

175
        Runtime *rt = get_indexer_runtime(self);
2 !
176
        IndexBuildBatchResult batch_result;
2 ✔
177

178
        if (!run_blocking([&] {
3 !
179
                rt->submit(
7 !
180
                      run_coro_scope(
3 !
181
                          rt->executor(),
2 ✔
182
                          [](CoroScope &scope,
8 !
183
                             std::shared_ptr<IndexBuildBatchConfig> cfg,
184
                             IndexBuildBatchResult *out)
185
                              -> coro::CoroTask<void> {
1 !
186
                              *out = co_await IndexBatchBuilderUtility::process(
8 !
187
                                  &scope, std::move(cfg));
3 ✔
188
                          },
4 !
189
                          batch_config, &batch_result),
2 ✔
190
                      "indexer-build")
1 !
191
                    .get();
2 !
192
            })) {
2 ✔
193
            return NULL;
×
194
        }
195

196
        if (batch_result.failed > 0 && !batch_result.results.empty()) {
2 !
197
            const auto &result = batch_result.results[0];
×
198
            if (!result.success) {
×
199
                PyErr_SetString(PyExc_RuntimeError,
×
200
                                result.error_message.c_str());
201
                return NULL;
×
202
            }
203
        }
204
    } else {
2 !
205
        // Simple checkpoint-only build
206
        int result;
207
        Py_BEGIN_ALLOW_THREADS result = dftu_indexer_build(self->handle);
8 ✔
208
        Py_END_ALLOW_THREADS
8 ✔
209

210
            if (result < 0) {
8 ✔
211
            PyErr_SetString(PyExc_RuntimeError, "Failed to build index");
×
212
            return NULL;
×
213
        }
214
    }
215

216
    Py_RETURN_NONE;
10 ✔
217
}
5 ✔
218

219
static PyObject *CheckpointIndexer_need_rebuild(CheckpointIndexerObject *self,
12 ✔
220
                                                PyObject *Py_UNUSED(ignored)) {
221
    if (!self->handle) {
12 ✔
222
        PyErr_SetString(PyExc_RuntimeError, "Indexer not initialized");
×
223
        return NULL;
×
224
    }
225

226
    int result = dftu_indexer_need_rebuild(self->handle);
12 ✔
227
    return PyBool_FromLong(result);
12 ✔
228
}
6 ✔
229

230
static PyObject *CheckpointIndexer_exists(CheckpointIndexerObject *self,
×
231
                                          PyObject *Py_UNUSED(ignored)) {
232
    if (!self->handle) {
×
233
        PyErr_SetString(PyExc_RuntimeError, "Indexer not initialized");
×
234
        return NULL;
×
235
    }
236

237
    int result = dftu_indexer_exists(self->handle);
×
238
    return PyBool_FromLong(result);
×
239
}
240

241
static PyObject *CheckpointIndexer_get_max_bytes(CheckpointIndexerObject *self,
6 ✔
242
                                                 PyObject *Py_UNUSED(ignored)) {
243
    if (!self->handle) {
6 ✔
244
        PyErr_SetString(PyExc_RuntimeError, "Indexer not initialized");
×
245
        return NULL;
×
246
    }
247

248
    uint64_t result = dftu_indexer_get_max_bytes(self->handle);
6 ✔
249
    return PyLong_FromUnsignedLongLong(result);
6 ✔
250
}
3 ✔
251

252
static PyObject *CheckpointIndexer_get_num_lines(CheckpointIndexerObject *self,
10 ✔
253
                                                 PyObject *Py_UNUSED(ignored)) {
254
    if (!self->handle) {
10 ✔
255
        PyErr_SetString(PyExc_RuntimeError, "Indexer not initialized");
×
256
        return NULL;
×
257
    }
258

259
    uint64_t result = dftu_indexer_get_num_lines(self->handle);
10 ✔
260
    return PyLong_FromUnsignedLongLong(result);
10 ✔
261
}
5 ✔
262

263
static PyObject *CheckpointIndexer_has_bloom(CheckpointIndexerObject *self,
2 ✔
264
                                             void *) {
265
    const char *idx = as_utf8(self->index_path);
2 ✔
266
    const char *gz = as_utf8(self->gz_path);
2 ✔
267
    if (!idx || !gz) {
2 !
UNCOV
268
        Py_RETURN_FALSE;
×
269
    }
270
    try {
271
        using namespace dftracer::utils::utilities::indexer;
272
        using namespace dftracer::utils::utilities::indexer::internal;
273
        IndexDatabase db(
1 !
274
            idx, dftracer::utils::utilities::indexer::IndexOpenMode::ReadOnly);
3 !
275
        std::string logical = get_logical_path(gz);
2 !
276
        int fid = db.get_file_info_id(logical);
2 !
277
        if (fid >= 0 && db.has_bloom_data(fid)) {
2 !
278
            Py_RETURN_TRUE;
2 ✔
279
        }
280
    } catch (...) {
3 !
UNCOV
281
    }
×
UNCOV
282
    Py_RETURN_FALSE;
×
283
}
1 ✔
284

285
static PyObject *CheckpointIndexer_gz_path(CheckpointIndexerObject *self,
6 ✔
286
                                           void *) {
287
    Py_INCREF(self->gz_path);
6 !
288
    return self->gz_path;
6 ✔
289
}
290

291
static PyObject *CheckpointIndexer_index_path(CheckpointIndexerObject *self,
4 ✔
292
                                              void *) {
293
    Py_INCREF(self->index_path);
4 !
294
    return self->index_path;
4 ✔
295
}
296

297
static PyObject *CheckpointIndexer_checkpoint_size(
4 ✔
298
    CheckpointIndexerObject *self, void *) {
299
    return PyLong_FromUnsignedLongLong(self->checkpoint_size);
4 ✔
300
}
301

302
static PyObject *CheckpointIndexer_enter(CheckpointIndexerObject *self,
12 !
303
                                         PyObject *Py_UNUSED(ignored)) {
304
    Py_INCREF(self);
6 ✔
305
    return (PyObject *)self;
12 ✔
306
}
307

308
static PyObject *CheckpointIndexer_close(CheckpointIndexerObject *self,
2 ✔
309
                                         PyObject *Py_UNUSED(ignored)) {
310
    CheckpointIndexer_release_handle(self);
2 ✔
311
    Py_RETURN_NONE;
2 ✔
312
}
313

314
static PyObject *CheckpointIndexer_exit(CheckpointIndexerObject *self,
12 ✔
315
                                        PyObject *) {
316
    CheckpointIndexer_release_handle(self);
12 ✔
317
    Py_RETURN_NONE;
12 ✔
318
}
319

320
static PyMethodDef CheckpointIndexer_methods[] = {
321
    {"build", DFTU_PYCFUNCTION(CheckpointIndexer_build), METH_NOARGS,
322
     "build()\n"
323
     "--\n"
324
     "\n"
325
     "Build or rebuild the index.\n"},
326
    {"need_rebuild", DFTU_PYCFUNCTION(CheckpointIndexer_need_rebuild),
327
     METH_NOARGS, "Check if a rebuild is needed."},
328
    {"exists", DFTU_PYCFUNCTION(CheckpointIndexer_exists), METH_NOARGS,
329
     "Check if the .dftindex store exists."},
330
    {"get_max_bytes", DFTU_PYCFUNCTION(CheckpointIndexer_get_max_bytes),
331
     METH_NOARGS, "Get the maximum uncompressed bytes in the indexed file."},
332
    {"get_num_lines", DFTU_PYCFUNCTION(CheckpointIndexer_get_num_lines),
333
     METH_NOARGS, "Get the total number of lines in the indexed file."},
334
    {"close", DFTU_PYCFUNCTION(CheckpointIndexer_close), METH_NOARGS,
335
     "Release this Python wrapper's native indexer handle.\n"
336
     "\n"
337
     "The shared RocksDB instance for the same .dftindex path remains managed\n"
338
     "by the native RocksDBManager cache."},
339
    {"__enter__", DFTU_PYCFUNCTION(CheckpointIndexer_enter), METH_NOARGS,
340
     "Enter the runtime context for the with statement."},
341
    {"__exit__", DFTU_PYCFUNCTION(CheckpointIndexer_exit), METH_VARARGS,
342
     "Release this Python wrapper on context exit.\n"
343
     "\n"
344
     "This does not force-close the shared RocksDB instance for the same\n"
345
     ".dftindex path."},
346
    {NULL} /* Sentinel */
347
};
348

349
static PyGetSetDef CheckpointIndexer_getsetters[] = {
350
    {"gz_path", (getter)CheckpointIndexer_gz_path, NULL,
351
     "Path to the gzip file", NULL},
352
    {"index_path", (getter)CheckpointIndexer_index_path, NULL,
353
     "Path to the .dftindex store", NULL},
354
    {"checkpoint_size", (getter)CheckpointIndexer_checkpoint_size, NULL,
355
     "Checkpoint size in bytes", NULL},
356
    {"has_bloom", (getter)CheckpointIndexer_has_bloom, NULL,
357
     "Whether bloom data exists in index", NULL},
358
    {NULL} /* Sentinel */
359
};
360

361
PyTypeObject CheckpointIndexerType = {
362
    PyVarObject_HEAD_INIT(
363
        NULL, 0) "dftracer_utils_ext.CheckpointIndexer", /* tp_name */
364
    sizeof(CheckpointIndexerObject),                     /* tp_basicsize */
365
    0,                                                   /* tp_itemsize */
366
    (destructor)CheckpointIndexer_dealloc,               /* tp_dealloc */
367
    0,                                        /* tp_vectorcall_offset */
368
    0,                                        /* tp_getattr */
369
    0,                                        /* tp_setattr */
370
    0,                                        /* tp_as_async */
371
    0,                                        /* tp_repr */
372
    0,                                        /* tp_as_number */
373
    0,                                        /* tp_as_sequence */
374
    0,                                        /* tp_as_mapping */
375
    0,                                        /* tp_hash */
376
    0,                                        /* tp_call */
377
    0,                                        /* tp_str */
378
    0,                                        /* tp_getattro */
379
    0,                                        /* tp_setattro */
380
    0,                                        /* tp_as_buffer */
381
    Py_TPFLAGS_DEFAULT | Py_TPFLAGS_BASETYPE, /* tp_flags */
382
    "CheckpointIndexer(gz_path, index_path=None, checkpoint_size=33554432, "
383
    "force_rebuild=False, build_bloom=False, runtime=None)\n"
384
    "--\n"
385
    "\n"
386
    "Checkpoint indexer for single-file checkpoint-level operations on a "
387
    "gzip trace.\n"
388
    "\n"
389
    "Args:\n"
390
    "    gz_path (str): Path to the gzip trace file.\n"
391
    "    index_path (str or None): Path to the .dftindex store. If None,\n"
392
    "        uses the root-local \".dftindex\" next to gz_path.\n"
393
    "    checkpoint_size (int): Checkpoint size in bytes for index\n"
394
    "        building (default 1 MB).\n"
395
    "    force_rebuild (bool): If True, rebuild the index even if it\n"
396
    "        exists.\n"
397
    "    build_bloom (bool): If True, build bloom filter data in the\n"
398
    "        index.\n"
399
    "    runtime (Runtime or None): Runtime instance for thread pool\n"
400
    "        control. If None, uses the default global Runtime.\n", /* tp_doc */
401
    0,                                /* tp_traverse */
402
    0,                                /* tp_clear */
403
    0,                                /* tp_richcompare */
404
    0,                                /* tp_weaklistoffset */
405
    0,                                /* tp_iter */
406
    0,                                /* tp_iternext */
407
    CheckpointIndexer_methods,        /* tp_methods */
408
    0,                                /* tp_members */
409
    CheckpointIndexer_getsetters,     /* tp_getset */
410
    0,                                /* tp_base */
411
    0,                                /* tp_dict */
412
    0,                                /* tp_descr_get */
413
    0,                                /* tp_descr_set */
414
    0,                                /* tp_dictoffset */
415
    (initproc)CheckpointIndexer_init, /* tp_init */
416
    0,                                /* tp_alloc */
417
    CheckpointIndexer_new,            /* tp_new */
418
};
419

420
int dftracer::utils::python::init_checkpoint_indexer(PyObject *m) {
2 ✔
421
    if (register_type(m, &CheckpointIndexerType, "CheckpointIndexer") < 0)
2 ✔
UNCOV
422
        return -1;
×
423

424
    return 0;
2 ✔
425
}
1 ✔
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc