• 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

18.86
/src/dftracer/utils/python/index_database.cpp
1
#include <dftracer/utils/python/index_database.h>
2
#include <dftracer/utils/python/py_errors.h>
3
#include <dftracer/utils/python/py_list_helpers.h>
4
#include <dftracer/utils/python/py_method.h>
5
#include <dftracer/utils/python/py_runtime_mixin.h>
6
#include <dftracer/utils/python/py_str_helpers.h>
7
#include <dftracer/utils/python/py_type_helpers.h>
8
#include <dftracer/utils/python/sst_distribution.h>
9
#include <dftracer/utils/utilities/indexer/index_database.h>
10
#include <dftracer/utils/utilities/indexer/index_database_sst_writer_context.h>
11

12
#include <new>
13
#include <string>
14
#include <unordered_set>
15
#include <vector>
16

17
using dftracer::utils::utilities::indexer::IndexDatabase;
18
using dftracer::utils::utilities::indexer::SstArtifactRegistry;
19

20
static void IndexDatabase_dealloc(IndexDatabaseObject *self) {
8 ✔
21
    self->db.~shared_ptr<IndexDatabase>();
8 ✔
22
    Py_TYPE(self)->tp_free((PyObject *)self);
8 ✔
23
}
8 ✔
24

25
static PyObject *IndexDatabase_new(PyTypeObject *type, PyObject * /*args*/,
8 ✔
26
                                   PyObject * /*kwds*/) {
27
    auto *self = (IndexDatabaseObject *)type->tp_alloc(type, 0);
8 ✔
28
    if (!self) return NULL;
8 ✔
29
    new (&self->db) std::shared_ptr<IndexDatabase>();
8 ✔
30
    return (PyObject *)self;
8 ✔
31
}
4 ✔
32

33
static int IndexDatabase_init(IndexDatabaseObject *self, PyObject *args,
8 ✔
34
                              PyObject *kwds) {
35
    static const char *kwlist[] = {"index_path", NULL};
36
    const char *index_path;
37
    if (!PyArg_ParseTupleAndKeywords(
8 !
38
            args, kwds, "s", const_cast<char **>(kwlist), &index_path)) {
4 ✔
UNCOV
39
        return -1;
×
40
    }
41
    try {
42
        self->db = std::make_shared<IndexDatabase>(index_path);
8 !
43
    } catch (const std::exception &e) {
4 !
44
        dftracer::utils::python::set_typed_py_error(e);
×
UNCOV
45
        return -1;
×
UNCOV
46
    }
×
47
    return 0;
8 ✔
48
}
4 ✔
49

50
static PyObject *IndexDatabase_init_schema(IndexDatabaseObject *self,
×
51
                                           PyObject * /*ignored*/) {
52
    if (!self->db) {
×
UNCOV
53
        PyErr_SetString(PyExc_RuntimeError, "IndexDatabase not initialised");
×
54
        return NULL;
×
55
    }
UNCOV
56
    if (!run_blocking([&] { self->db->init_schema(); })) return NULL;
×
UNCOV
57
    Py_RETURN_NONE;
×
58
}
59

UNCOV
60
static PyObject *IndexDatabase_register_files(IndexDatabaseObject *self,
×
61
                                              PyObject *args, PyObject *kwds) {
62
    static const char *kwlist[] = {"paths", NULL};
63
    PyObject *paths_obj;
UNCOV
64
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "O",
×
65
                                     const_cast<char **>(kwlist), &paths_obj)) {
UNCOV
66
        return NULL;
×
67
    }
68
    std::vector<std::string> paths;
×
UNCOV
69
    if (!parse_str_list(paths_obj, "paths", paths)) return NULL;
×
70

71
    std::vector<int> ids;
×
72
    if (!run_blocking_r([&] { return self->db->register_files(paths); }, ids)) {
×
73
        return NULL;
×
74
    }
75

76
    PyObject *out = PyList_New(static_cast<Py_ssize_t>(ids.size()));
×
UNCOV
77
    if (!out) return NULL;
×
UNCOV
78
    for (Py_ssize_t i = 0; i < static_cast<Py_ssize_t>(ids.size()); ++i) {
×
79
        PyList_SET_ITEM(out, i, PyLong_FromLong(ids[i]));
×
80
    }
81
    return out;
×
82
}
×
83

84
static PyObject *IndexDatabase_reserve_file_id_range(IndexDatabaseObject *self,
×
85
                                                     PyObject *args) {
86
    Py_ssize_t count;
87
    if (!PyArg_ParseTuple(args, "n", &count)) return NULL;
×
UNCOV
88
    if (count < 0) {
×
UNCOV
89
        PyErr_SetString(PyExc_ValueError, "count must be >= 0");
×
90
        return NULL;
×
91
    }
92
    int first;
93
    if (!run_blocking_r(
×
UNCOV
94
            [&] {
×
UNCOV
95
                return self->db->reserve_file_id_range(
×
96
                    static_cast<std::size_t>(count));
×
97
            },
98
            first)) {
99
        return NULL;
×
100
    }
UNCOV
101
    return PyLong_FromLong(first);
×
102
}
103

104
static PyObject *IndexDatabase_bulk_ingest(IndexDatabaseObject *self,
×
105
                                           PyObject *args, PyObject *kwds) {
106
    static const char *kwlist[] = {"registry", "skip_cfs", NULL};
107
    PyObject *registry_obj;
UNCOV
108
    PyObject *skip_cfs_obj = NULL;
×
UNCOV
109
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "O|O",
×
110
                                     const_cast<char **>(kwlist), &registry_obj,
111
                                     &skip_cfs_obj)) {
112
        return NULL;
×
113
    }
114

115
    SstArtifactRegistry *registry =
UNCOV
116
        dftracer::utils::python::sst_artifact_registry_get(registry_obj);
×
117
    if (!registry) {
×
118
        PyErr_SetString(PyExc_TypeError,
×
119
                        "expected an SstArtifactRegistry instance");
UNCOV
120
        return NULL;
×
121
    }
122

UNCOV
123
    std::unordered_set<std::string> skip_cfs;
×
124
    if (skip_cfs_obj && skip_cfs_obj != Py_None) {
×
125
        PyObject *seq =
UNCOV
126
            PySequence_Fast(skip_cfs_obj, "skip_cfs must be an iterable");
×
127
        if (!seq) return NULL;
×
128
        Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
129
        for (Py_ssize_t i = 0; i < n; ++i) {
×
130
            PyObject *item = PySequence_Fast_GET_ITEM(seq, i);
×
131
            const char *s = as_utf8(item);
×
132
            if (!s) {
×
133
                Py_DECREF(seq);
×
UNCOV
134
                return NULL;
×
135
            }
UNCOV
136
            skip_cfs.emplace(s);
×
137
        }
138
        Py_DECREF(seq);
×
139
    }
140

UNCOV
141
    if (!run_blocking([&] { self->db->bulk_ingest(*registry, skip_cfs); }))
×
142
        return NULL;
×
143
    Py_RETURN_NONE;
×
144
}
×
145

UNCOV
146
static PyObject *IndexDatabase_write_agg_file_markers(IndexDatabaseObject *self,
×
147
                                                      PyObject *args) {
148
    PyObject *ids_obj;
UNCOV
149
    if (!PyArg_ParseTuple(args, "O", &ids_obj)) return NULL;
×
150

UNCOV
151
    PyObject *seq = PySequence_Fast(ids_obj, "file_ids must be an iterable");
×
152
    if (!seq) return NULL;
×
153
    Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
154
    std::vector<int> file_ids;
×
155
    file_ids.reserve(static_cast<std::size_t>(n));
×
156
    for (Py_ssize_t i = 0; i < n; ++i) {
×
157
        PyObject *item = PySequence_Fast_GET_ITEM(seq, i);
×
158
        long v = PyLong_AsLong(item);
×
159
        if (v == -1 && PyErr_Occurred()) {
×
160
            Py_DECREF(seq);
×
UNCOV
161
            return NULL;
×
162
        }
UNCOV
163
        file_ids.push_back(static_cast<int>(v));
×
164
    }
165
    Py_DECREF(seq);
×
166

UNCOV
167
    if (!run_blocking([&] { self->db->write_agg_file_markers(file_ids); }))
×
168
        return NULL;
×
169
    Py_RETURN_NONE;
×
170
}
×
171

UNCOV
172
static PyObject *IndexDatabase_write_agg_global_config(
×
173
    IndexDatabaseObject *self, PyObject *args, PyObject *kwds) {
174
    static const char *kwlist[] = {"time_interval_us", "config_hash",
175
                                   "group_by_file", NULL};
176
    unsigned long long time_interval_us = 0;
×
177
    unsigned int config_hash = 0;
×
178
    int group_by_file = 1;
×
UNCOV
179
    if (!PyArg_ParseTupleAndKeywords(
×
180
            args, kwds, "K|Ip", const_cast<char **>(kwlist), &time_interval_us,
181
            &config_hash, &group_by_file)) {
182
        return NULL;
×
183
    }
184
    if (!run_blocking([&] {
×
185
            self->db->write_agg_global_config(
×
186
                static_cast<std::uint64_t>(time_interval_us),
×
187
                static_cast<std::uint32_t>(config_hash), group_by_file != 0);
×
UNCOV
188
        })) {
×
189
        return NULL;
×
190
    }
UNCOV
191
    Py_RETURN_NONE;
×
192
}
193

UNCOV
194
static PyObject *IndexDatabase_write_aggregation_tracker(
×
195
    IndexDatabaseObject *self, PyObject *args) {
196
    PyObject *blobs_obj;
197
    if (!PyArg_ParseTuple(args, "O", &blobs_obj)) return NULL;
×
198
    PyObject *seq = PySequence_Fast(blobs_obj, "blobs must be an iterable");
×
199
    if (!seq) return NULL;
×
200
    Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
201
    std::vector<std::string> blobs;
×
202
    blobs.reserve(static_cast<std::size_t>(n));
×
203
    for (Py_ssize_t i = 0; i < n; ++i) {
×
204
        PyObject *item = PySequence_Fast_GET_ITEM(seq, i);
×
205
        if (item == Py_None) continue;
×
206
        char *buf = nullptr;
×
207
        Py_ssize_t len = 0;
×
UNCOV
208
        if (PyBytes_Check(item)) {
×
209
            if (PyBytes_AsStringAndSize(item, &buf, &len) < 0) {
×
210
                Py_DECREF(seq);
×
UNCOV
211
                return NULL;
×
212
            }
213
        } else {
214
            Py_DECREF(seq);
×
215
            PyErr_SetString(PyExc_TypeError,
×
216
                            "blobs entries must be bytes or None");
217
            return NULL;
×
218
        }
UNCOV
219
        if (len > 0) blobs.emplace_back(buf, static_cast<std::size_t>(len));
×
220
    }
221
    Py_DECREF(seq);
×
222
    if (!run_blocking([&] { self->db->write_aggregation_tracker(blobs); }))
×
223
        return NULL;
×
UNCOV
224
    Py_RETURN_NONE;
×
225
}
×
226

227
static PyObject *IndexDatabase_rebuild_root_summaries(IndexDatabaseObject *self,
×
228
                                                      PyObject * /*ignored*/) {
UNCOV
229
    if (!run_blocking([&] { self->db->rebuild_root_summaries(); })) return NULL;
×
UNCOV
230
    Py_RETURN_NONE;
×
231
}
232

233
static PyObject *IndexDatabase_find_stale_files(IndexDatabaseObject *self,
8 ✔
234
                                                PyObject *args,
235
                                                PyObject *kwds) {
236
    static const char *kwlist[] = {"paths", NULL};
237
    PyObject *paths_obj;
238
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "O",
8 !
239
                                     const_cast<char **>(kwlist), &paths_obj)) {
UNCOV
240
        return NULL;
×
241
    }
242
    std::vector<std::string> paths;
8 ✔
243
    if (!parse_str_list(paths_obj, "paths", paths)) return NULL;
8 !
244

245
    IndexDatabase::StaleCheckResult result;
8 ✔
246
    if (!run_blocking_r([&] { return self->db->find_stale_files(paths); },
16 !
247
                        result)) {
UNCOV
248
        return NULL;
×
249
    }
250

251
    PyObject *changed = str_list_from(result.changed);
8 !
252
    PyObject *added = str_list_from(result.added);
8 !
253
    PyObject *removed = str_list_from(result.removed);
8 !
254
    PyObject *d = PyDict_New();
8 !
255
    if (!changed || !added || !removed || !d) {
8 !
UNCOV
256
        Py_XDECREF(changed);
×
UNCOV
257
        Py_XDECREF(added);
×
UNCOV
258
        Py_XDECREF(removed);
×
UNCOV
259
        Py_XDECREF(d);
×
UNCOV
260
        return NULL;
×
261
    }
262
    PyObject *so = PyBool_FromLong(result.schema_outdated);
8 !
263
    PyObject *st = PyBool_FromLong(result.stale());
8 !
264
    PyDict_SetItemString(d, "changed", changed);
8 !
265
    PyDict_SetItemString(d, "added", added);
8 !
266
    PyDict_SetItemString(d, "removed", removed);
8 !
267
    PyDict_SetItemString(d, "schema_outdated", so);
8 !
268
    PyDict_SetItemString(d, "stale", st);
8 !
269
    Py_DECREF(changed);
4 !
270
    Py_DECREF(added);
4 !
271
    Py_DECREF(removed);
4 !
272
    Py_DECREF(so);
4 !
273
    Py_DECREF(st);
4 !
274
    return d;
8 ✔
275
}
8 ✔
276

277
static PyMethodDef IndexDatabase_methods[] = {
278
    {"init_schema", DFTU_PYCFUNCTION(IndexDatabase_init_schema), METH_NOARGS,
279
     "Idempotently initialise the schema version key."},
280
    {"register_files", DFTU_PYCFUNCTION(IndexDatabase_register_files),
281
     METH_VARARGS | METH_KEYWORDS,
282
     "register_files(paths) -> list[int]\n"
283
     "Register each path in the DEFAULT-CF file registry and return the "
284
     "assigned file_ids. Idempotent for files with matching hash."},
285
    {"find_stale_files", DFTU_PYCFUNCTION(IndexDatabase_find_stale_files),
286
     METH_VARARGS | METH_KEYWORDS,
287
     "find_stale_files(paths) -> dict\n"
288
     "Stat-only (mtime + size) staleness check of the given trace paths "
289
     "against the index. Returns {changed, added, removed, schema_outdated, "
290
     "stale}."},
291
    {"reserve_file_id_range",
292
     DFTU_PYCFUNCTION(IndexDatabase_reserve_file_id_range), METH_VARARGS,
293
     "reserve_file_id_range(count) -> int\n"
294
     "Atomically reserve `count` contiguous file_ids, return the first."},
295
    {"bulk_ingest", DFTU_PYCFUNCTION(IndexDatabase_bulk_ingest),
296
     METH_VARARGS | METH_KEYWORDS,
297
     "bulk_ingest(registry, skip_cfs=None) -> None\n"
298
     "Ingest all SSTs collected in the SstArtifactRegistry.\n"
299
     "skip_cfs is an optional iterable of CF names whose SSTs are left "
300
     "outside the unified DB (used by distributed builds to keep "
301
     "AGGREGATION/SYSTEM_METRICS SSTs addressable by manifest)."},
302
    {"rebuild_root_summaries",
303
     DFTU_PYCFUNCTION(IndexDatabase_rebuild_root_summaries), METH_NOARGS,
304
     "Recompute ROOT_* summary column families from per-file CFs."},
305
    {"write_agg_global_config",
306
     DFTU_PYCFUNCTION(IndexDatabase_write_agg_global_config),
307
     METH_VARARGS | METH_KEYWORDS,
308
     "write_agg_global_config(time_interval_us, config_hash=0, "
309
     "group_by_file=True) -> None\n"
310
     "Write the AGG_GLOBAL_CONFIG_KEY marker into the AGGREGATION CF. "
311
     "Required for `iter_arrow_dfanalyzer_all` on distributed builds "
312
     "(which never materialise the key via worker SSTs) or "
313
     "post-consolidate indices."},
314
    {"write_agg_file_markers",
315
     DFTU_PYCFUNCTION(IndexDatabase_write_agg_file_markers), METH_VARARGS,
316
     "write_agg_file_markers(file_ids) -> None\n"
317
     "Write per-file aggregation completion markers (\\xFF\\xFF + file_id) "
318
     "into the AGGREGATION CF. Required after distributed_index otherwise "
319
     "`ensure_indexed()` concludes aggregation is incomplete and re-runs "
320
     "the entire build."},
321
    {"write_aggregation_tracker",
322
     DFTU_PYCFUNCTION(IndexDatabase_write_aggregation_tracker), METH_VARARGS,
323
     "write_aggregation_tracker(blobs) -> None\n"
324
     "Merge a list of serialized AssociationTracker bytes and write the "
325
     "result to the AGGREGATION CF under the `__tracker__` key."},
326
    {NULL}};
327

328
PyTypeObject IndexDatabaseType = {
329
    PyVarObject_HEAD_INIT(NULL, 0) "dftracer_utils_ext.IndexDatabase",
330
    sizeof(IndexDatabaseObject),
331
    0,
332
    (destructor)IndexDatabase_dealloc,
333
    0,
334
    0,
335
    0,
336
    0,
337
    0,
338
    0,
339
    0,
340
    0,
341
    0,
342
    0,
343
    0,
344
    0,
345
    0,
346
    0,
347
    Py_TPFLAGS_DEFAULT,
348
    "Handle to a .dftindex RocksDB store.",
349
    0,
350
    0,
351
    0,
352
    0,
353
    0,
354
    0,
355
    IndexDatabase_methods,
356
    0,
357
    0,
358
    0,
359
    0,
360
    0,
361
    0,
362
    0,
363
    (initproc)IndexDatabase_init,
364
    0,
365
    IndexDatabase_new,
366
};
367

368
int dftracer::utils::python::init_index_database(PyObject *m) {
2 ✔
369
    if (register_type(m, &IndexDatabaseType, "IndexDatabase") < 0) return -1;
2 ✔
370
    return 0;
2 ✔
371
}
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