• 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

56.5
/src/dftracer/utils/python/batch_indexer.cpp
1
#include <dftracer/utils/core/common/constants.h>
2
#include <dftracer/utils/core/common/filesystem.h>
3
#include <dftracer/utils/core/common/hash/hash_combine.h>
4
#include <dftracer/utils/core/common/string_intern.h>
5
#include <dftracer/utils/core/coro/task.h>
6
#include <dftracer/utils/core/coro/when_all.h>
7
#include <dftracer/utils/core/rocksdb/db_manager.h>
8
#include <dftracer/utils/core/runtime.h>
9
#include <dftracer/utils/core/tasks/coro_scope.h>
10
#include <dftracer/utils/python/batch_indexer.h>
11
#include <dftracer/utils/python/indexer.h>
12
#include <dftracer/utils/python/py_dict_helpers.h>
13
#include <dftracer/utils/python/py_list_helpers.h>
14
#include <dftracer/utils/python/py_method.h>
15
#include <dftracer/utils/python/py_runtime_mixin.h>
16
#include <dftracer/utils/python/py_seq_helpers.h>
17
#include <dftracer/utils/python/py_str_helpers.h>
18
#include <dftracer/utils/python/py_type_helpers.h>
19
#include <dftracer/utils/python/runtime.h>
20
#include <dftracer/utils/query/query.h>
21
#include <dftracer/utils/trace/aggregators/aggregation_config.h>
22
#include <dftracer/utils/trace/aggregators/aggregation_serialization.h>
23
#include <dftracer/utils/trace/aggregators/aggregator_types.h>
24
#include <dftracer/utils/trace/aggregators/event_aggregator.h>
25
#include <dftracer/utils/trace/aggregators/system_metrics.h>
26
#include <dftracer/utils/trace/aggregators/system_metrics_serialization.h>
27
#include <dftracer/utils/trace/indexing/index_resolver_utility.h>
28
#include <dftracer/utils/trace/indexing/resolve_and_build.h>
29
#include <dftracer/utils/trace/internal/utils.h>
30
#include <dftracer/utils/utilities/indexer/index_database.h>
31

32
#include <algorithm>
33
#include <chrono>
34
#include <cstdio>
35
#include <optional>
36
#include <string>
37
#include <unordered_map>
38
#include <unordered_set>
39
#include <vector>
40

41
using dftracer::utils::CoroScope;
42
using dftracer::utils::Runtime;
43
using dftracer::utils::coro::CoroTask;
44
using namespace dftracer::utils::trace::indexing;
45
using namespace dftracer::utils::trace::aggregators;
46

47
// ---------------------------------------------------------------------------
48
// BatchIndexer - directory-level indexer with resolve/build pattern
49
// ---------------------------------------------------------------------------
50

51
static void Indexer_dealloc(IndexerObject* self) {
330 ✔
52
    Py_XDECREF(self->runtime_obj);
330 ✔
53
    Py_XDECREF(self->directory);
330 ✔
54
    Py_XDECREF(self->files);
330 ✔
55
    Py_XDECREF(self->index_dir);
330 ✔
56
    Py_XDECREF(self->group_keys);
330 ✔
57
    Py_XDECREF(self->custom_metric_fields);
330 ✔
58
    Py_XDECREF(self->bloom_fields);
330 ✔
59
    Py_TYPE(self)->tp_free((PyObject*)self);
330 ✔
60
}
330 ✔
61

62
static PyObject* Indexer_new(PyTypeObject* type, PyObject*, PyObject*) {
330 ✔
63
    IndexerObject* self = (IndexerObject*)type->tp_alloc(type, 0);
330 ✔
64
    if (self) {
330 ✔
65
        self->runtime_obj = nullptr;
330 ✔
66
        self->directory = nullptr;
330 ✔
67
        self->files = nullptr;
330 ✔
68
        self->index_dir = nullptr;
330 ✔
69
        self->require_checkpoint = 1;
330 ✔
70
        self->require_bloom = 1;
330 ✔
71
        self->build_bloom = 1;
330 ✔
72
        self->require_aggregation = 0;
330 ✔
73
        self->bloom_fields = nullptr;
330 ✔
74
        self->false_positive_rate = ChunkIndexerConfig{}.false_positive_rate;
330 ✔
75
        self->expected_entries =
330 ✔
76
            ChunkIndexerConfig{}.expected_entries_per_chunk;
330 ✔
77
        self->auto_fields = ChunkIndexerConfig{}.auto_fields ? 1 : 0;
330 !
78
        self->auto_max_distinct = ChunkIndexerConfig{}.auto_max_distinct;
330 ✔
79
        self->time_interval_ms = 5000.0;
330 ✔
80
        self->group_keys = nullptr;
330 ✔
81
        self->custom_metric_fields = nullptr;
330 ✔
82
        self->compute_percentiles = 0;
330 ✔
83
        self->group_by_file = 1;
330 ✔
84
        self->checkpoint_size =
330 ✔
85
            dftracer::utils::constants::indexer::DEFAULT_CHECKPOINT_SIZE;
86
        self->parallelism = 0;
330 ✔
87
        self->force_rebuild = 0;
330 ✔
88
    }
165 ✔
89
    return (PyObject*)self;
330 ✔
90
}
91

92
static int Indexer_init(IndexerObject* self, PyObject* args, PyObject* kwds) {
330 ✔
93
    static const char* kwlist[] = {
94
        "directory",           "files",
95
        "index_dir",           "require_checkpoint",
96
        "require_bloom",       "build_bloom",
97
        "require_aggregation", "time_interval_ms",
98
        "group_keys",          "custom_metric_fields",
99
        "compute_percentiles", "group_by_file",
100
        "checkpoint_size",     "parallelism",
101
        "force_rebuild",       "runtime",
102
        "bloom_fields",        "false_positive_rate",
103
        "expected_entries",    "auto_fields",
104
        "auto_max_distinct",   nullptr};
105

106
    const char* directory = "";
330 ✔
107
    PyObject* files_obj = Py_None;
330 ✔
108
    const char* index_dir = "";
330 ✔
109
    int require_checkpoint = 1;
330 ✔
110
    int require_bloom = 1;
330 ✔
111
    int build_bloom = 1;
330 ✔
112
    int require_aggregation = 0;
330 ✔
113
    double time_interval_ms = 5000.0;
330 ✔
114
    PyObject* group_keys_obj = Py_None;
330 ✔
115
    PyObject* custom_metrics_obj = Py_None;
330 ✔
116
    int compute_percentiles = 0;
330 ✔
117
    int group_by_file = 1;
330 ✔
118
    Py_ssize_t checkpoint_size = static_cast<Py_ssize_t>(
330 ✔
119
        dftracer::utils::constants::indexer::DEFAULT_CHECKPOINT_SIZE);
120
    Py_ssize_t parallelism = 0;
330 ✔
121
    int force_rebuild = 0;
330 ✔
122
    PyObject* runtime_arg = nullptr;
330 ✔
123
    PyObject* bloom_fields_obj = Py_None;
330 ✔
124
    double false_positive_rate = self->false_positive_rate;
330 ✔
125
    Py_ssize_t expected_entries =
330 ✔
126
        static_cast<Py_ssize_t>(self->expected_entries);
330 ✔
127
    int auto_fields = self->auto_fields;
330 ✔
128
    Py_ssize_t auto_max_distinct =
330 ✔
129
        static_cast<Py_ssize_t>(self->auto_max_distinct);
330 ✔
130

131
    if (!PyArg_ParseTupleAndKeywords(
330 !
132
            args, kwds, "|sOsppppdOOppnnpOOdnpn", const_cast<char**>(kwlist),
165 ✔
133
            &directory, &files_obj, &index_dir, &require_checkpoint,
134
            &require_bloom, &build_bloom, &require_aggregation,
135
            &time_interval_ms, &group_keys_obj, &custom_metrics_obj,
136
            &compute_percentiles, &group_by_file, &checkpoint_size,
137
            &parallelism, &force_rebuild, &runtime_arg, &bloom_fields_obj,
138
            &false_positive_rate, &expected_entries, &auto_fields,
139
            &auto_max_distinct)) {
UNCOV
140
        return -1;
×
141
    }
142
    if (!(false_positive_rate > 0.0 && false_positive_rate < 1.0)) {
330 !
143
        PyErr_SetString(PyExc_ValueError,
2 !
144
                        "false_positive_rate must be in (0, 1)");
145
        return -1;
2 ✔
146
    }
147
    if (expected_entries <= 0) {
328 !
NEW
148
        PyErr_SetString(PyExc_ValueError, "expected_entries must be > 0");
×
NEW
149
        return -1;
×
150
    }
151
    if (auto_max_distinct <= 0) {
328 !
NEW
152
        PyErr_SetString(PyExc_ValueError, "auto_max_distinct must be > 0");
×
NEW
153
        return -1;
×
154
    }
155
    self->auto_fields = auto_fields;
328 ✔
156
    self->auto_max_distinct = static_cast<std::size_t>(auto_max_distinct);
328 ✔
157
    if (bloom_fields_obj != Py_None) {
328 ✔
158
        std::vector<std::string> check;
328 ✔
159
        if (!dftracer::utils::python::parse_string_seq(
328 !
160
                bloom_fields_obj, "bloom_fields must be a sequence of str",
164 ✔
161
                check))
NEW
162
            return -1;
×
163
        Py_INCREF(bloom_fields_obj);
328 !
164
        self->bloom_fields = bloom_fields_obj;
328 ✔
165
    }
328 !
166
    self->false_positive_rate = false_positive_rate;
328 ✔
167
    self->expected_entries = static_cast<std::size_t>(expected_entries);
328 ✔
168

169
    // Validate: at least one of directory or files must be provided
170
    bool has_directory = directory && directory[0] != '\0';
328 ✔
171
    bool has_files = files_obj && files_obj != Py_None &&
616 ✔
172
                     PyList_Check(files_obj) && PyList_Size(files_obj) > 0;
616 !
173

174
    if (!has_directory && !has_files) {
328 ✔
175
        PyErr_SetString(PyExc_ValueError,
2 !
176
                        "At least one of 'directory' or 'files' must be "
177
                        "provided");
178
        return -1;
2 ✔
179
    }
180

181
    // Store runtime
182
    if (runtime_arg && runtime_arg != Py_None) {
326 !
183
        if (PyObject_TypeCheck(runtime_arg, &RuntimeType)) {
×
184
            Py_INCREF(runtime_arg);
×
UNCOV
185
            self->runtime_obj = runtime_arg;
×
186
        } else {
187
            PyObject* native = PyObject_GetAttrString(runtime_arg, "_native");
×
188
            if (native && PyObject_TypeCheck(native, &RuntimeType)) {
×
UNCOV
189
                self->runtime_obj = native;
×
190
            } else {
191
                Py_XDECREF(native);
×
UNCOV
192
                PyErr_SetString(PyExc_TypeError,
×
193
                                "runtime must be a Runtime instance or None");
UNCOV
194
                return -1;
×
195
            }
196
        }
197
    }
198

199
    self->directory = PyUnicode_FromString(directory);
326 !
200
    self->index_dir = PyUnicode_FromString(index_dir);
326 !
201
    self->require_checkpoint = require_checkpoint;
326 ✔
202
    self->require_bloom = require_bloom;
326 ✔
203
    self->build_bloom = build_bloom;
326 ✔
204
    self->require_aggregation = require_aggregation;
326 ✔
205
    self->time_interval_ms = time_interval_ms;
326 ✔
206
    self->compute_percentiles = compute_percentiles;
326 ✔
207
    self->group_by_file = group_by_file;
326 ✔
208
    self->checkpoint_size = static_cast<std::size_t>(checkpoint_size);
326 ✔
209
    self->parallelism = static_cast<std::size_t>(parallelism);
326 ✔
210
    self->force_rebuild = force_rebuild;
326 ✔
211

212
    // Store files list
213
    if (has_files) {
326 ✔
214
        Py_INCREF(files_obj);
288 !
215
        self->files = files_obj;
288 ✔
216
    } else {
144 ✔
217
        self->files = nullptr;
38 ✔
218
    }
219

220
    // Store group_keys
221
    if (group_keys_obj && group_keys_obj != Py_None) {
326 !
UNCOV
222
        Py_INCREF(group_keys_obj);
×
UNCOV
223
        self->group_keys = group_keys_obj;
×
224
    } else {
225
        self->group_keys = nullptr;
326 ✔
226
    }
227

228
    // Store custom_metric_fields
229
    if (custom_metrics_obj && custom_metrics_obj != Py_None) {
326 !
UNCOV
230
        Py_INCREF(custom_metrics_obj);
×
UNCOV
231
        self->custom_metric_fields = custom_metrics_obj;
×
232
    } else {
233
        self->custom_metric_fields = nullptr;
326 ✔
234
    }
235

236
    return 0;
326 ✔
237
}
165 ✔
238

239
// The args fields the bloom tier covers, as the index names them. False with
240
// a Python error set when bloom_fields holds a non-str.
241
static bool requested_bloom_fields(IndexerObject* self,
960 ✔
242
                                   std::vector<std::string>& out) {
243
    if (!self->bloom_fields) {
960 ✔
NEW
244
        out.clear();
×
NEW
245
        return true;
×
246
    }
247
    out.clear();
960 ✔
248
    if (!dftracer::utils::python::parse_string_seq(
960 !
249
            self->bloom_fields, "bloom_fields must be a sequence of str", out))
480 ✔
NEW
250
        return false;
×
251
    for (std::string& f : out) f = extra_dimension_name(f);
986 ✔
252
    return true;
960 ✔
253
}
480 ✔
254

255
static Runtime* get_batch_indexer_runtime(IndexerObject* self) {
960 ✔
256
    if (self->runtime_obj) {
960 !
UNCOV
257
        return ((RuntimeObject*)self->runtime_obj)->runtime.get();
×
258
    }
259
    return dftracer::utils::python::get_default_runtime();
960 ✔
260
}
480 ✔
261

262
static std::optional<AggregationConfig> build_aggregation_config(
960 ✔
263
    IndexerObject* self) {
264
    if (!self->require_aggregation) {
960 ✔
265
        return std::nullopt;
750 ✔
266
    }
267

268
    AggregationConfig config;
210 !
269
    config.time_interval_us =
210 ✔
270
        static_cast<std::uint64_t>(self->time_interval_ms * 1000.0);
210 ✔
271

272
    if (self->group_keys && PyList_Check(self->group_keys)) {
210 !
UNCOV
273
        Py_ssize_t n = PyList_Size(self->group_keys);
×
274
        for (Py_ssize_t i = 0; i < n; i++) {
×
275
            const char* s = as_utf8(PyList_GetItem(self->group_keys, i));
×
UNCOV
276
            if (s) config.extra_group_keys.emplace_back(s);
×
277
        }
278
    }
279
    if (self->custom_metric_fields &&
210 !
280
        PyList_Check(self->custom_metric_fields)) {
×
281
        Py_ssize_t n = PyList_Size(self->custom_metric_fields);
×
UNCOV
282
        for (Py_ssize_t i = 0; i < n; i++) {
×
283
            const char* s =
284
                as_utf8(PyList_GetItem(self->custom_metric_fields, i));
×
UNCOV
285
            if (s) config.custom_metric_fields.emplace_back(s);
×
286
        }
287
    }
288

289
    config.compute_percentiles = self->compute_percentiles != 0;
210 ✔
290
    config.group_by_file = self->group_by_file != 0;
210 ✔
291
    return config;
210 !
292
}
585 ✔
293

294
// ---------------------------------------------------------------------------
295
// resolve() - check what exists vs needs building
296
// ---------------------------------------------------------------------------
297

298
static PyObject* Indexer_resolve(IndexerObject* self,
654 ✔
299
                                 PyObject* Py_UNUSED(ignored)) {
300
    const char* directory = as_utf8(self->directory);
654 !
301
    const char* index_dir = as_utf8(self->index_dir);
654 !
302

303
    ResolverInput input;
654 ✔
304
    input.directory = directory ? directory : "";
654 !
305
    input.index_dir = index_dir ? index_dir : "";
654 !
306
    input.require_checkpoints = self->require_checkpoint;
654 ✔
307
    input.require_bloom = self->require_bloom;
654 ✔
308
    input.require_aggregation = self->require_aggregation;
654 ✔
309
    input.checkpoint_size = self->checkpoint_size;
654 ✔
310
    input.aggregation_config = build_aggregation_config(self);
654 !
311
    if (self->build_bloom) {
654 !
312
        if (!requested_bloom_fields(self, input.bloom_fields)) return nullptr;
654 !
313
        if (self->auto_fields)
654 ✔
314
            input.bloom_fields.emplace_back(AUTO_FIELDS_MARKER);
632 !
315
    }
327 ✔
316

317
    // Add files if provided
318
    if (self->files && PyList_Check(self->files)) {
654 !
319
        Py_ssize_t n = PyList_Size(self->files);
594 !
320
        for (Py_ssize_t i = 0; i < n; i++) {
1,290 ✔
321
            const char* s = as_utf8(PyList_GetItem(self->files, i));
696 !
322
            if (s) input.files.emplace_back(s);
696 !
323
        }
348 ✔
324
    }
297 ✔
325

326
    ResolverResult result;
654 ✔
327

328
    if (!run_blocking([&] {
981 !
329
            Runtime* rt = get_batch_indexer_runtime(self);
654 ✔
330
            rt->submit(run_coro_scope(
2,616 !
331
                           rt->executor(),
654 ✔
332
                           [](CoroScope& scope, ResolverInput in,
2,616 !
333
                              ResolverResult* out) -> CoroTask<void> {
327 !
334
                               IndexResolverUtility resolver;
981 ✔
335
                               *out = co_await resolver(scope, std::move(in));
1,635 !
336
                           },
1,635 !
337
                           std::move(input), &result),
654 ✔
338
                       "batch-indexer-resolve")
327 !
339
                .get();
654 !
340
        })) {
654 ✔
UNCOV
341
        return nullptr;
×
342
    }
343

344
    // Build result dict
345
    PyObject* dict = PyDict_New();
654 !
346
    if (!dict) return nullptr;
654 ✔
347

348
    dict_set_steal(dict, "total_files",
654 !
349
                   PyLong_FromSize_t(result.all_files.size()));
327 !
350
    dict_set_steal(dict, "index_path",
654 !
351
                   PyUnicode_FromString(result.index_path.c_str()));
327 !
352
    dict_set_steal(dict, "aggregation_interval_us",
654 !
353
                   PyLong_FromUnsignedLongLong(result.stored_time_interval_us));
654 !
354
    dict_set_steal(dict, "needs_rebuild",
654 !
355
                   PyBool_FromLong(result.needs_augmentation));
654 !
356

357
    // Ready files
358
    PyObject* ready_list = PyList_New(result.cached.size());
654 !
359
    for (std::size_t i = 0; i < result.cached.size(); ++i) {
1,056 ✔
360
        PyList_SetItem(
402 !
361
            ready_list, i,
201 ✔
362
            PyUnicode_FromString(result.cached[i].file_path.c_str()));
402 !
363
    }
201 ✔
364
    PyDict_SetItemString(dict, "ready", ready_list);
654 !
365

366
    // Needs work files (union of all needs_* lists)
367
    std::vector<std::string> needs_work;
654 ✔
368
    for (const auto& item : result.needs_checkpoint) {
1,020 ✔
369
        needs_work.push_back(item.file_path);
366 !
370
    }
371
    for (const auto& item : result.needs_bloom) {
660 ✔
372
        bool found = false;
6 ✔
373
        for (const auto& existing : needs_work) {
6 ✔
374
            if (existing == item.file_path) {
×
375
                found = true;
×
376
                break;
×
377
            }
378
        }
379
        if (!found) needs_work.push_back(item.file_path);
6 !
380
    }
381
    for (const auto& item : result.needs_aggregation) {
654 ✔
UNCOV
382
        bool found = false;
×
383
        for (const auto& existing : needs_work) {
×
384
            if (existing == item.file_path) {
×
385
                found = true;
×
386
                break;
×
387
            }
388
        }
UNCOV
389
        if (!found) needs_work.push_back(item.file_path);
×
390
    }
391

392
    PyObject* needs_list = PyList_New(needs_work.size());
654 !
393
    for (std::size_t i = 0; i < needs_work.size(); ++i) {
1,026 ✔
394
        PyList_SetItem(needs_list, i,
372 !
395
                       PyUnicode_FromString(needs_work[i].c_str()));
372 !
396
    }
186 ✔
397
    PyDict_SetItemString(dict, "needs_work", needs_list);
654 !
398

399
    return dict;
654 ✔
400
}
654 ✔
401

402
// ---------------------------------------------------------------------------
403
// build() - build missing index tiers
404
// ---------------------------------------------------------------------------
405

406
static PyObject* Indexer_build(IndexerObject* self,
306 ✔
407
                               PyObject* Py_UNUSED(ignored)) {
408
    const char* directory = as_utf8(self->directory);
306 !
409
    const char* index_dir = as_utf8(self->index_dir);
306 !
410

411
    ResolveAndBuildInput input;
306 ✔
412
    input.directory = directory ? directory : "";
306 !
413
    input.index_dir = index_dir ? index_dir : "";
306 !
414
    input.require_checkpoints = self->require_checkpoint;
306 ✔
415
    input.require_bloom = self->require_bloom;
306 ✔
416
    input.build_bloom = self->build_bloom;
306 ✔
417
    input.require_aggregation = self->require_aggregation;
306 ✔
418
    input.aggregation_config = build_aggregation_config(self);
306 !
419
    if (!requested_bloom_fields(self, input.bloom_config.extra_dimensions))
306 !
NEW
420
        return nullptr;
×
421
    input.bloom_config.false_positive_rate = self->false_positive_rate;
306 ✔
422
    input.bloom_config.auto_fields = self->auto_fields != 0;
306 ✔
423
    input.bloom_config.auto_max_distinct = self->auto_max_distinct;
306 ✔
424
    input.bloom_config.expected_entries_per_chunk = self->expected_entries;
306 ✔
425
    input.checkpoint_size = self->checkpoint_size;
306 ✔
426
    input.parallelism = self->parallelism;
306 ✔
427
    input.force_rebuild = self->force_rebuild;
306 ✔
428

429
    // Add files if provided
430
    if (self->files && PyList_Check(self->files)) {
306 !
431
        Py_ssize_t n = PyList_Size(self->files);
280 !
432
        for (Py_ssize_t i = 0; i < n; i++) {
610 ✔
433
            const char* s = as_utf8(PyList_GetItem(self->files, i));
330 !
434
            if (s) input.files.emplace_back(s);
330 !
435
        }
165 ✔
436
    }
140 ✔
437

438
    if (!run_blocking([&] {
459 !
439
            Runtime* rt = get_batch_indexer_runtime(self);
306 ✔
440
            rt->submit(run_coro_scope(
1,224 !
441
                           rt->executor(),
306 ✔
442
                           [](CoroScope& scope,
1,224 !
443
                              ResolveAndBuildInput in) -> CoroTask<void> {
153 !
444
                               co_await resolve_and_build_index(&scope,
1,224 !
445
                                                                std::move(in));
459 ✔
446
                           },
612 !
447
                           std::move(input)),
306 ✔
448
                       "batch-indexer-build")
153 !
449
                .get();
306 !
450
        })) {
306 ✔
UNCOV
451
        return nullptr;
×
452
    }
453

454
    Py_RETURN_NONE;
306 ✔
455
}
306 ✔
456

457
// ---------------------------------------------------------------------------
458
// ensure_indexed() - resolve + build if needed
459
// ---------------------------------------------------------------------------
460

461
static PyObject* Indexer_ensure_indexed(IndexerObject* self,
306 ✔
462
                                        PyObject* Py_UNUSED(ignored)) {
463
    // First resolve
464
    PyObject* status = Indexer_resolve(self, nullptr);
306 ✔
465
    if (!status) return nullptr;
306 ✔
466

467
    // Build if files need work, or the aggregation tier must be rebuilt
468
    // (stored time interval differs from the requested one).
469
    PyObject* needs_work = PyDict_GetItemString(status, "needs_work");
306 ✔
470
    PyObject* needs_rebuild = PyDict_GetItemString(status, "needs_rebuild");
306 ✔
471
    bool work_pending = needs_work && PyList_Size(needs_work) > 0;
306 ✔
472
    bool rebuild_pending = needs_rebuild && PyObject_IsTrue(needs_rebuild);
306 ✔
473
    if (work_pending || rebuild_pending) {
306 ✔
474
        Py_DECREF(status);
152 ✔
475

476
        // Build
477
        PyObject* result = Indexer_build(self, nullptr);
304 ✔
478
        if (!result) return nullptr;
304 ✔
479
        Py_DECREF(result);
152 ✔
480

481
        // Re-resolve
482
        status = Indexer_resolve(self, nullptr);
304 ✔
483
    }
152 ✔
484

485
    return status;
306 ✔
486
}
153 ✔
487

488
// ---------------------------------------------------------------------------
489
// get_checkpoint_indexer() - get a single-file checkpoint indexer
490
// ---------------------------------------------------------------------------
491

492
static PyObject* Indexer_get_checkpoint_indexer(IndexerObject* self,
10 ✔
493
                                                PyObject* args) {
494
    const char* file_path = nullptr;
10 ✔
495
    if (!PyArg_ParseTuple(args, "s", &file_path)) {
10 !
UNCOV
496
        return nullptr;
×
497
    }
498

499
    // Determine index path using BatchIndexer's index_dir setting
500
    const char* index_dir = as_utf8(self->index_dir);
10 !
501
    std::string index_path =
502
        dftracer::utils::trace::internal::determine_index_path(
5 !
503
            file_path, index_dir ? index_dir : "");
15 !
504

505
    // Create IndexerObject
506
    CheckpointIndexerObject* indexer =
5 ✔
507
        (CheckpointIndexerObject*)CheckpointIndexerType.tp_alloc(
10 !
508
            &CheckpointIndexerType, 0);
509
    if (!indexer) {
10 ✔
UNCOV
510
        return nullptr;
×
511
    }
512

513
    indexer->handle = nullptr;
10 ✔
514
    indexer->gz_path = PyUnicode_FromString(file_path);
10 !
515
    indexer->index_path = PyUnicode_FromString(index_path.c_str());
10 !
516
    indexer->checkpoint_size = self->checkpoint_size;
10 ✔
517
    indexer->build_bloom = 0;
10 ✔
518

519
    // Share runtime reference
520
    if (self->runtime_obj) {
10 !
UNCOV
521
        Py_INCREF(self->runtime_obj);
×
UNCOV
522
        indexer->runtime_obj = self->runtime_obj;
×
523
    } else {
524
        indexer->runtime_obj = nullptr;
10 ✔
525
    }
526

527
    // Create the native handle
528
    indexer->handle = dftu_indexer_create(file_path, index_path.c_str(),
15 !
529
                                          self->checkpoint_size, 0);
5 ✔
530
    if (!indexer->handle) {
10 !
531
        Py_DECREF((PyObject*)indexer);
×
532
        PyErr_SetString(PyExc_RuntimeError,
×
533
                        "Failed to create checkpoint indexer");
UNCOV
534
        return nullptr;
×
535
    }
536

537
    return (PyObject*)indexer;
10 ✔
538
}
10 ✔
539

540
static std::optional<std::string> resolve_index_path(IndexerObject* self) {
20 ✔
541
    PyObject* status = Indexer_resolve(self, nullptr);
20 !
542
    if (!status) return std::nullopt;
20 ✔
543
    PyObject* obj = PyDict_GetItemString(status, "index_path");
20 !
544
    const char* path = obj ? as_utf8(obj) : nullptr;
20 !
545
    if (!path || path[0] == '\0') {
20 !
546
        Py_DECREF(status);
UNCOV
547
        PyErr_SetString(PyExc_RuntimeError, "No index path available");
×
UNCOV
548
        return std::nullopt;
×
549
    }
550
    std::string result(path);
30 !
551
    Py_DECREF(status);
10 !
552
    return result;
20 !
553
}
20 ✔
554

555
static PyObject* Indexer_get_hash_table(IndexerObject* self, PyObject* args) {
12 ✔
556
    const char* type_str = nullptr;
12 ✔
557
    if (!PyArg_ParseTuple(args, "s", &type_str)) {
12 !
558
        return nullptr;
×
559
    }
560

561
    using dftracer::utils::utilities::indexer::IndexDatabase;
562
    using HashType = IndexDatabase::HashType;
563

564
    HashType type;
565
    if (std::strcmp(type_str, "file") == 0) {
12 ✔
566
        type = HashType::FILE;
4 ✔
567
    } else if (std::strcmp(type_str, "host") == 0) {
10 ✔
568
        type = HashType::HOST;
4 ✔
569
    } else if (std::strcmp(type_str, "string") == 0) {
6 ✔
570
        type = HashType::STRING;
2 ✔
571
    } else if (std::strcmp(type_str, "proc") == 0) {
3 ✔
UNCOV
572
        type = HashType::PROC;
×
573
    } else {
574
        PyErr_SetString(PyExc_ValueError,
2 !
575
                        "type must be 'file', 'host', 'string', or 'proc'");
576
        return nullptr;
2 ✔
577
    }
578

579
    auto idx_opt = resolve_index_path(self);
10 !
580
    if (!idx_opt) return nullptr;
10 ✔
581
    std::string index_path = std::move(*idx_opt);
10 ✔
582

583
    std::unordered_map<std::string, std::string> hash_map;
10 ✔
584
    if (!run_blocking_r(
10 !
585
            [&] {
15 ✔
586
                IndexDatabase db(index_path,
10 ✔
587
                                 dftracer::utils::utilities::indexer::
588
                                     IndexOpenMode::ReadOnly);
5 !
589
                return db.query_hash_table(type);
15 !
590
            },
10 ✔
591
            hash_map)) {
UNCOV
592
        return nullptr;
×
593
    }
594

595
    PyObject* dict = PyDict_New();
10 !
596
    if (!dict) return nullptr;
10 ✔
597

598
    for (const auto& [hash, name] : hash_map) {
10 !
UNCOV
599
        PyObject* key = PyUnicode_FromStringAndSize(hash.data(), hash.size());
×
UNCOV
600
        PyObject* val = PyUnicode_FromStringAndSize(name.data(), name.size());
×
UNCOV
601
        PyDict_SetItem(dict, key, val);
×
602
        Py_DECREF(key);
×
603
        Py_DECREF(val);
×
604
    }
605

606
    return dict;
10 ✔
607
}
11 ✔
608

609
static PyObject* Indexer_query_file_pids(IndexerObject* self, PyObject* args) {
4 ✔
610
    int file_id;
611
    if (!PyArg_ParseTuple(args, "i", &file_id)) {
4 !
612
        return nullptr;
×
613
    }
614

615
    using dftracer::utils::utilities::indexer::IndexDatabase;
616

617
    auto idx_opt = resolve_index_path(self);
4 !
618
    if (!idx_opt) return nullptr;
4 ✔
619
    std::string index_path = std::move(*idx_opt);
4 ✔
620

621
    std::unordered_set<std::uint64_t> pids;
4 ✔
622
    if (!run_blocking_r(
4 !
623
            [&] {
6 ✔
624
                IndexDatabase db(index_path,
4 ✔
625
                                 dftracer::utils::utilities::indexer::
626
                                     IndexOpenMode::ReadOnly);
2 !
627
                return db.query_file_pids(file_id);
6 !
628
            },
4 ✔
629
            pids)) {
UNCOV
630
        return nullptr;
×
631
    }
632

633
    PyObject* set = PySet_New(nullptr);
4 !
634
    if (!set) return nullptr;
4 ✔
635

636
    for (auto pid : pids) {
6 !
637
        PyObject* val = PyLong_FromUnsignedLongLong(pid);
2 !
638
        PySet_Add(set, val);
2 !
639
        Py_DECREF(val);
1 !
640
    }
641

642
    return set;
4 ✔
643
}
4 ✔
644

645
static PyObject* Indexer_query_all_file_pids(IndexerObject* self,
6 ✔
646
                                             PyObject* Py_UNUSED(ignored)) {
647
    using dftracer::utils::utilities::indexer::IndexDatabase;
648

649
    auto idx_opt = resolve_index_path(self);
6 !
650
    if (!idx_opt) return nullptr;
6 !
651
    std::string index_path = std::move(*idx_opt);
6 ✔
652

653
    std::unordered_map<int, std::unordered_set<std::uint64_t>> all_pids;
6 ✔
654
    if (!run_blocking_r(
6 !
655
            [&] {
9 ✔
656
                IndexDatabase db(index_path,
6 ✔
657
                                 dftracer::utils::utilities::indexer::
658
                                     IndexOpenMode::ReadOnly);
3 !
659
                return db.query_all_file_pids();
9 !
660
            },
6 ✔
661
            all_pids)) {
UNCOV
662
        return nullptr;
×
663
    }
664

665
    PyObject* dict = PyDict_New();
6 !
666
    if (!dict) return nullptr;
6 ✔
667

668
    for (const auto& [file_id, pids] : all_pids) {
14 !
669
        PyObject* key = PyLong_FromLong(file_id);
8 !
670
        PyObject* set = PySet_New(nullptr);
8 !
671
        for (auto pid : pids) {
22 ✔
672
            PyObject* val = PyLong_FromUnsignedLongLong(pid);
14 !
673
            PySet_Add(set, val);
14 !
674
            Py_DECREF(val);
7 !
675
        }
676
        PyDict_SetItem(dict, key, set);
8 !
677
        Py_DECREF(key);
4 !
678
        Py_DECREF(set);
4 !
679
    }
680

681
    return dict;
6 ✔
682
}
6 ✔
683

UNCOV
684
static PyObject* Indexer_query_file_info(IndexerObject* self,
×
685
                                         PyObject* Py_UNUSED(ignored)) {
686
    using dftracer::utils::utilities::indexer::IndexDatabase;
687

UNCOV
688
    auto idx_opt = resolve_index_path(self);
×
UNCOV
689
    if (!idx_opt) return nullptr;
×
UNCOV
690
    std::string index_path = std::move(*idx_opt);
×
691

UNCOV
692
    std::unordered_map<std::string, int> file_ids;
×
UNCOV
693
    std::unordered_map<int, std::unordered_set<std::uint64_t>> all_pids;
×
694

695
    if (!run_blocking([&] {
×
696
            IndexDatabase db(
UNCOV
697
                index_path,
×
698
                dftracer::utils::utilities::indexer::IndexOpenMode::ReadOnly);
×
699
            file_ids = db.query_all_file_info_ids();
×
700
            all_pids = db.query_all_file_pids();
×
701
        })) {
×
UNCOV
702
        return nullptr;
×
703
    }
704

UNCOV
705
    auto data_dir = fs::weakly_canonical(fs::path(index_path)).parent_path();
×
706

UNCOV
707
    PyObject* id_to_path = PyDict_New();
×
708
    if (!id_to_path) return nullptr;
×
UNCOV
709
    for (const auto& [logical_name, fid] : file_ids) {
×
710
        auto resolved = (data_dir / logical_name).string();
×
711
        PyObject* key = PyLong_FromLong(fid);
×
712
        PyObject* val = PyUnicode_FromStringAndSize(
×
713
            resolved.data(), static_cast<Py_ssize_t>(resolved.size()));
×
UNCOV
714
        PyDict_SetItem(id_to_path, key, val);
×
715
        Py_DECREF(key);
×
716
        Py_DECREF(val);
×
UNCOV
717
    }
×
718

719
    PyObject* pid_dict = PyDict_New();
×
720
    if (!pid_dict) {
×
721
        Py_DECREF(id_to_path);
×
722
        return nullptr;
×
723
    }
724
    for (const auto& [file_id, pids] : all_pids) {
×
725
        PyObject* key = PyLong_FromLong(file_id);
×
UNCOV
726
        PyObject* set = PySet_New(nullptr);
×
UNCOV
727
        for (auto pid : pids) {
×
728
            PyObject* val = PyLong_FromUnsignedLongLong(pid);
×
UNCOV
729
            PySet_Add(set, val);
×
730
            Py_DECREF(val);
×
731
        }
UNCOV
732
        PyDict_SetItem(pid_dict, key, set);
×
733
        Py_DECREF(key);
×
734
        Py_DECREF(set);
×
735
    }
736

737
    PyObject* result = PyTuple_Pack(2, id_to_path, pid_dict);
×
738
    Py_DECREF(id_to_path);
×
739
    Py_DECREF(pid_dict);
×
740
    return result;
×
UNCOV
741
}
×
742

743
#ifdef DFTRACER_UTILS_ENABLE_ARROW
UNCOV
744
static PyObject* count_hash_entries_fn(PyObject* /*self*/, PyObject* args) {
×
UNCOV
745
    const char* index_path = nullptr;
×
UNCOV
746
    const char* type_str = nullptr;
×
UNCOV
747
    if (!PyArg_ParseTuple(args, "ss", &index_path, &type_str)) return nullptr;
×
748

749
    using dftracer::utils::utilities::indexer::IndexDatabase;
750
    using HashType = IndexDatabase::HashType;
751

752
    HashType type;
UNCOV
753
    if (std::strcmp(type_str, "file") == 0) {
×
UNCOV
754
        type = HashType::FILE;
×
UNCOV
755
    } else if (std::strcmp(type_str, "host") == 0) {
×
UNCOV
756
        type = HashType::HOST;
×
UNCOV
757
    } else if (std::strcmp(type_str, "string") == 0) {
×
UNCOV
758
        type = HashType::STRING;
×
UNCOV
759
    } else if (std::strcmp(type_str, "proc") == 0) {
×
UNCOV
760
        type = HashType::PROC;
×
761
    } else {
UNCOV
762
        PyErr_SetString(PyExc_ValueError,
×
763
                        "type must be 'file', 'host', 'string', or 'proc'");
UNCOV
764
        return nullptr;
×
765
    }
766

UNCOV
767
    std::uint64_t count = 0;
×
UNCOV
768
    if (!run_blocking_r(
×
UNCOV
769
            [&] {
×
UNCOV
770
                IndexDatabase db(index_path,
×
771
                                 dftracer::utils::utilities::indexer::
772
                                     IndexOpenMode::ReadOnly);
×
UNCOV
773
                return db.count_hash_entries(type);
×
UNCOV
774
            },
×
775
            count)) {
UNCOV
776
        return nullptr;
×
777
    }
UNCOV
778
    return PyLong_FromUnsignedLongLong(count);
×
779
}
780

781
static PyMethodDef BatchIndexerModuleMethods[] = {
782
    {"count_hash_entries", DFTU_PYCFUNCTION(count_hash_entries_fn),
783
     METH_VARARGS,
784
     "count_hash_entries(index_path, type)\n"
785
     "--\n\n"
786
     "Number of hashes of `type` ('file', 'host', 'string', 'proc') in the\n"
787
     "index at `index_path`. Counted by iteration, so a table holding tens\n"
788
     "of millions of entries is not materialised to take its length.\n"},
789
    {nullptr, nullptr, 0, nullptr}};
790
#endif
791

792
static PyMethodDef Indexer_methods[] = {
793
    {"get_checkpoint_indexer", DFTU_PYCFUNCTION(Indexer_get_checkpoint_indexer),
794
     METH_VARARGS,
795
     "get_checkpoint_indexer(file_path)\n"
796
     "--\n\n"
797
     "Get a checkpoint indexer for a specific file.\n\n"
798
     "Args:\n"
799
     "    file_path: Path to the trace file (.pfw/.pfw.gz)\n\n"
800
     "Returns:\n"
801
     "    Indexer instance for checkpoint-level operations.\n"},
802
    {"resolve", DFTU_PYCFUNCTION(Indexer_resolve), METH_NOARGS,
803
     "resolve()\n"
804
     "--\n\n"
805
     "Check what files exist vs need indexing.\n\n"
806
     "Returns:\n"
807
     "    dict with 'total_files', 'ready', 'needs_work', 'index_path'\n"},
808
    {"build", DFTU_PYCFUNCTION(Indexer_build), METH_NOARGS,
809
     "build()\n"
810
     "--\n\n"
811
     "Build all missing index tiers based on require_* flags.\n"},
812
    {"ensure_indexed", DFTU_PYCFUNCTION(Indexer_ensure_indexed), METH_NOARGS,
813
     "ensure_indexed()\n"
814
     "--\n\n"
815
     "Resolve and build if needed.\n\n"
816
     "Returns:\n"
817
     "    dict with index status after building.\n"},
818
    {"get_hash_table", DFTU_PYCFUNCTION(Indexer_get_hash_table), METH_VARARGS,
819
     "get_hash_table(type)\n"
820
     "--\n\n"
821
     "Query hash table mappings.\n\n"
822
     "Args:\n"
823
     "    type: 'file', 'host', 'string', or 'proc'\n\n"
824
     "Returns:\n"
825
     "    dict mapping hash values to resolved names.\n"},
826
    {"query_file_pids", DFTU_PYCFUNCTION(Indexer_query_file_pids), METH_VARARGS,
827
     "query_file_pids(file_id)\n"
828
     "--\n\n"
829
     "Query PIDs observed in a specific file.\n\n"
830
     "Args:\n"
831
     "    file_id: Integer file ID from index.\n\n"
832
     "Returns:\n"
833
     "    set of PIDs.\n"},
834
    {"query_all_file_pids", DFTU_PYCFUNCTION(Indexer_query_all_file_pids),
835
     METH_NOARGS,
836
     "query_all_file_pids()\n"
837
     "--\n\n"
838
     "Query PIDs for all indexed files.\n\n"
839
     "Returns:\n"
840
     "    dict mapping file_id to set of PIDs.\n"},
841
    {"query_file_info", DFTU_PYCFUNCTION(Indexer_query_file_info), METH_NOARGS,
842
     "query_file_info()\n"
843
     "--\n\n"
844
     "Query file ID to path mapping and per-file PIDs in one call.\n\n"
845
     "Returns:\n"
846
     "    tuple of (dict[int, str], dict[int, set[int]]).\n"},
847
    {nullptr}};
848

849
static PyGetSetDef Indexer_getsetters[] = {{nullptr}};
850

851
PyTypeObject IndexerType = {
852
    PyVarObject_HEAD_INIT(nullptr, 0) "dftracer_utils_ext.Indexer",
853
    sizeof(IndexerObject),
854
    0,
855
    (destructor)Indexer_dealloc,
856
    0,
857
    0,
858
    0,
859
    0,
860
    0,
861
    0,
862
    0,
863
    0,
864
    0,
865
    0,
866
    0,
867
    0,
868
    0,
869
    0,
870
    Py_TPFLAGS_DEFAULT | Py_TPFLAGS_BASETYPE,
871
    "BatchIndexer(directory='', files=None, index_dir='',\n"
872
    "             require_checkpoint=True, require_bloom=True,\n"
873
    "             time_interval_ms=5000.0, group_keys=None,\n"
874
    "             custom_metric_fields=None, compute_percentiles=False,\n"
875
    "             parallelism=0, force_rebuild=False, runtime=None)\n"
876
    "--\n\n"
877
    "Indexer with tiered index building.\n\n"
878
    "At least one of 'directory' or 'files' must be provided.\n"
879
    "- directory: scan for .pfw/.pfw.gz files\n"
880
    "- files: list of specific file paths\n\n"
881
    "Supports:\n"
882
    "- Tier 1: Checkpoints (require_checkpoint)\n"
883
    "- Tier 2: Bloom filters (require_bloom)\n"
884
    "- Tier 3: Aggregation (require_aggregation + config params)\n",
885
    0,
886
    0,
887
    0,
888
    0,
889
    0,
890
    0,
891
    Indexer_methods,
892
    0,
893
    Indexer_getsetters,
894
    0,
895
    0,
896
    0,
897
    0,
898
    0,
899
    (initproc)Indexer_init,
900
    0,
901
    Indexer_new,
902
};
903

904
int dftracer::utils::python::init_indexer(PyObject* m) {
2 ✔
905
    if (register_type(m, &IndexerType, "Indexer") < 0) return -1;
2 ✔
906

907
#ifdef DFTRACER_UTILS_ENABLE_ARROW
908
    if (PyModule_AddFunctions(m, BatchIndexerModuleMethods) < 0) return -1;
2 ✔
909
#endif
910

911
    return 0;
2 ✔
912
}
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