• 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

0.49
/src/dftracer/utils/python/sst_distribution.cpp
1
#include <dftracer/utils/core/common/constants.h>
2
#include <dftracer/utils/core/runtime.h>
3
#include <dftracer/utils/core/tasks/coro_scope.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_seq_helpers.h>
8
#include <dftracer/utils/python/py_str_helpers.h>
9
#include <dftracer/utils/python/py_type_helpers.h>
10
#include <dftracer/utils/python/runtime.h>
11
#include <dftracer/utils/python/sst_distribution.h>
12
#include <dftracer/utils/trace/aggregators/aggregation_config.h>
13
#include <dftracer/utils/trace/aggregators/aggregation_key.h>
14
#include <dftracer/utils/trace/aggregators/association_tracker.h>
15
#include <dftracer/utils/trace/views/aggregation_fold.h>
16
#include <dftracer/utils/utilities/filesystem/pattern_directory_scanner_utility.h>
17
#include <dftracer/utils/utilities/indexer/file_partition.h>
18
#include <dftracer/utils/utilities/indexer/index_batch_sink.h>
19
#include <dftracer/utils/utilities/indexer/index_builder_utility.h>
20
#include <dftracer/utils/utilities/indexer/index_database_sst_writer_context.h>
21
#include <dftracer/utils/utilities/indexer/internal/common/gzip_member_scanner.h>
22
#include <fcntl.h>
23
#include <sys/stat.h>
24
#include <unistd.h>
25

26
#include <algorithm>
27
#include <atomic>
28
#include <cstdint>
29
#include <memory>
30
#include <mutex>
31
#include <new>
32
#include <optional>
33
#include <string>
34
#include <unordered_set>
35
#include <vector>
36

37
using dftracer::utils::Runtime;
38
using dftracer::utils::python::parse_string_seq;
39
using dftracer::utils::utilities::filesystem::FileEntry;
40
using dftracer::utils::utilities::filesystem::PatternDirectoryScannerUtility;
41
using dftracer::utils::utilities::filesystem::
42
    PatternDirectoryScannerUtilityInput;
43
using dftracer::utils::utilities::indexer::IndexBatchBuilderUtility;
44
using dftracer::utils::utilities::indexer::IndexBatchSink;
45
using dftracer::utils::utilities::indexer::IndexBuildBatchConfig;
46
using dftracer::utils::utilities::indexer::IndexBuildBatchResult;
47
using dftracer::utils::utilities::indexer::IndexDatabaseSstWriterContext;
48
using dftracer::utils::utilities::indexer::plan_lpt_partition;
49
using dftracer::utils::utilities::indexer::SstArtifactRegistry;
50
using dftracer::utils::utilities::indexer::internal::
51
    enumerate_gzip_member_candidates;
52
using dftracer::utils::utilities::indexer::internal::GzipMember;
53

54
// ---------------------------------------------------------------------------
55
// SstArtifactRegistry type
56
// ---------------------------------------------------------------------------
57

58
typedef struct {
59
    PyObject_HEAD std::shared_ptr<SstArtifactRegistry> registry;
60
} SstArtifactRegistryObject;
61

UNCOV
62
static void SstArtifactRegistry_dealloc(SstArtifactRegistryObject *self) {
×
63
    self->registry.~shared_ptr<SstArtifactRegistry>();
×
UNCOV
64
    Py_TYPE(self)->tp_free((PyObject *)self);
×
UNCOV
65
}
×
66

67
static PyObject *SstArtifactRegistry_new(PyTypeObject *type,
×
68
                                         PyObject * /*args*/,
69
                                         PyObject * /*kwds*/) {
70
    auto *self = (SstArtifactRegistryObject *)type->tp_alloc(type, 0);
×
UNCOV
71
    if (!self) return NULL;
×
UNCOV
72
    new (&self->registry) std::shared_ptr<SstArtifactRegistry>(
×
UNCOV
73
        std::make_shared<SstArtifactRegistry>());
×
UNCOV
74
    return (PyObject *)self;
×
75
}
76

77
namespace {
78

79
// Field names in the Artifacts dict returned by build_sst_batch and
80
// consumed by SstArtifactRegistry.append. Must match the field names on
81
// IndexDatabaseSstWriterContext::Artifacts.
82
constexpr const char *ARTIFACT_FIELDS[] = {
83
    "metadata_sst",
84
    "members_sst",
85
    "manifest_sst",
86
    "chunk_bloom_sst",
87
    "file_bloom_sst",
88
    "chunk_stats_sst",
89
    "chunk_dim_stats_sst",
90
    "dimensions_sst",
91
    "file_scalar_stats_sst",
92
    "file_cat_counts_sst",
93
    "file_pid_tid_counts_sst",
94
    "file_name_counts_sst",
95
    "name_dictionary_sst",
96
    "name_file_postings_sst",
97
    "name_chunk_postings_sst",
98
    "hash_tables_sst",
99
    "aggregation_sst",
100
    "system_metrics_sst",
101
};
102

103
/// Map a slot name to the matching Artifacts member. Kept in one place so
104
/// that adding a new CF requires updating only `ARTIFACT_FIELDS` plus
105
/// `dispatch_*` below.
106
std::optional<std::string> *artifacts_slot(
×
107
    IndexDatabaseSstWriterContext::Artifacts &a, std::string_view name) {
108
    if (name == "metadata_sst") return &a.metadata_sst;
×
109
    if (name == "members_sst") return &a.members_sst;
×
110
    if (name == "manifest_sst") return &a.manifest_sst;
×
UNCOV
111
    if (name == "chunk_bloom_sst") return &a.chunk_bloom_sst;
×
UNCOV
112
    if (name == "file_bloom_sst") return &a.file_bloom_sst;
×
UNCOV
113
    if (name == "chunk_stats_sst") return &a.chunk_stats_sst;
×
UNCOV
114
    if (name == "chunk_dim_stats_sst") return &a.chunk_dim_stats_sst;
×
UNCOV
115
    if (name == "dimensions_sst") return &a.dimensions_sst;
×
116
    if (name == "file_scalar_stats_sst") return &a.file_scalar_stats_sst;
×
UNCOV
117
    if (name == "file_cat_counts_sst") return &a.file_cat_counts_sst;
×
118
    if (name == "file_pid_tid_counts_sst") return &a.file_pid_tid_counts_sst;
×
119
    if (name == "file_name_counts_sst") return &a.file_name_counts_sst;
×
120
    if (name == "name_dictionary_sst") return &a.name_dictionary_sst;
×
UNCOV
121
    if (name == "name_file_postings_sst") return &a.name_file_postings_sst;
×
122
    if (name == "name_chunk_postings_sst") return &a.name_chunk_postings_sst;
×
123
    if (name == "hash_tables_sst") return &a.hash_tables_sst;
×
124
    if (name == "aggregation_sst") return &a.aggregation_sst;
×
125
    if (name == "system_metrics_sst") return &a.system_metrics_sst;
×
126
    return nullptr;
×
127
}
128

129
/// Convert a Python artifacts dict to the C++ Artifacts struct. Missing,
130
/// None, or empty-string entries become nullopt. Returns false on type
131
/// errors (exception set).
132
bool artifacts_from_dict(PyObject *dict,
×
133
                         IndexDatabaseSstWriterContext::Artifacts *out) {
134
    if (!PyDict_Check(dict)) {
×
UNCOV
135
        PyErr_SetString(PyExc_TypeError, "artifacts must be a dict");
×
136
        return false;
×
137
    }
UNCOV
138
    for (const char *field : ARTIFACT_FIELDS) {
×
139
        PyObject *val = PyDict_GetItemString(dict, field);  // borrowed
×
140
        if (!val || val == Py_None) continue;
×
141
        if (!PyUnicode_Check(val)) {
×
142
            PyErr_Format(PyExc_TypeError, "artifacts['%s'] must be str or None",
×
143
                         field);
144
            return false;
×
145
        }
146
        const char *s = as_utf8(val);
×
147
        if (!s) return false;
×
UNCOV
148
        if (s[0] == '\0') continue;
×
149
        auto *slot = artifacts_slot(*out, field);
×
150
        if (slot) *slot = std::string(s);
×
151
    }
152
    return true;
×
153
}
154

155
PyObject *artifacts_to_dict(const IndexDatabaseSstWriterContext::Artifacts &a) {
×
156
    PyObject *dict = PyDict_New();
×
157
    if (!dict) return NULL;
×
158
    auto set_field = [&](const char *name,
×
159
                         const std::optional<std::string> &slot) -> bool {
160
        PyObject *v = slot.has_value() ? PyUnicode_FromString(slot->c_str())
×
161
                                       : (Py_INCREF(Py_None), Py_None);
×
162
        if (!v) return false;
×
163
        int rc = PyDict_SetItemString(dict, name, v);
×
164
        Py_DECREF(v);
165
        return rc == 0;
×
166
    };
×
167
    if (!set_field("metadata_sst", a.metadata_sst) ||
×
168
        !set_field("members_sst", a.members_sst) ||
×
169
        !set_field("manifest_sst", a.manifest_sst) ||
×
170
        !set_field("chunk_bloom_sst", a.chunk_bloom_sst) ||
×
UNCOV
171
        !set_field("file_bloom_sst", a.file_bloom_sst) ||
×
172
        !set_field("chunk_stats_sst", a.chunk_stats_sst) ||
×
UNCOV
173
        !set_field("chunk_dim_stats_sst", a.chunk_dim_stats_sst) ||
×
UNCOV
174
        !set_field("dimensions_sst", a.dimensions_sst) ||
×
UNCOV
175
        !set_field("file_scalar_stats_sst", a.file_scalar_stats_sst) ||
×
UNCOV
176
        !set_field("file_cat_counts_sst", a.file_cat_counts_sst) ||
×
177
        !set_field("file_pid_tid_counts_sst", a.file_pid_tid_counts_sst) ||
×
UNCOV
178
        !set_field("file_name_counts_sst", a.file_name_counts_sst) ||
×
UNCOV
179
        !set_field("name_dictionary_sst", a.name_dictionary_sst) ||
×
180
        !set_field("name_file_postings_sst", a.name_file_postings_sst) ||
×
181
        !set_field("name_chunk_postings_sst", a.name_chunk_postings_sst) ||
×
182
        !set_field("hash_tables_sst", a.hash_tables_sst) ||
×
183
        !set_field("aggregation_sst", a.aggregation_sst) ||
×
184
        !set_field("system_metrics_sst", a.system_metrics_sst)) {
×
185
        Py_DECREF(dict);
×
UNCOV
186
        return NULL;
×
187
    }
UNCOV
188
    return dict;
×
189
}
190

191
}  // namespace
192

UNCOV
193
static PyObject *SstArtifactRegistry_append(SstArtifactRegistryObject *self,
×
194
                                            PyObject *args) {
195
    PyObject *dict;
UNCOV
196
    if (!PyArg_ParseTuple(args, "O", &dict)) return NULL;
×
UNCOV
197
    IndexDatabaseSstWriterContext::Artifacts a;
×
UNCOV
198
    if (!artifacts_from_dict(dict, &a)) return NULL;
×
UNCOV
199
    self->registry->append(std::move(a));
×
UNCOV
200
    Py_RETURN_NONE;
×
UNCOV
201
}
×
202

203
static PyMethodDef SstArtifactRegistry_methods[] = {
204
    {"append", DFTU_PYCFUNCTION(SstArtifactRegistry_append), METH_VARARGS,
205
     "append(artifacts_dict) -> None\n"
206
     "Add a per-batch Artifacts dict (as returned by build_sst_batch or "
207
     "IndexDatabaseSstWriterContext.commit) to the registry."},
208
    {NULL}};
209

210
static PyTypeObject SstArtifactRegistryType = {
211
    PyVarObject_HEAD_INIT(NULL, 0) "dftracer_utils_ext.SstArtifactRegistry",
212
    sizeof(SstArtifactRegistryObject),
213
    0,
214
    (destructor)SstArtifactRegistry_dealloc,
215
    0,
216
    0,
217
    0,
218
    0,
219
    0,
220
    0,
221
    0,
222
    0,
223
    0,
224
    0,
225
    0,
226
    0,
227
    0,
228
    0,
229
    Py_TPFLAGS_DEFAULT,
230
    "Thread-safe collector for SST artifact paths.",
231
    0,
232
    0,
233
    0,
234
    0,
235
    0,
236
    0,
237
    SstArtifactRegistry_methods,
238
    0,
239
    0,
240
    0,
241
    0,
242
    0,
243
    0,
244
    0,
245
    0,
246
    0,
247
    SstArtifactRegistry_new,
248
};
249

250
SstArtifactRegistry *dftracer::utils::python::sst_artifact_registry_get(
×
251
    PyObject *obj) {
UNCOV
252
    if (!PyObject_TypeCheck(obj, &SstArtifactRegistryType)) return nullptr;
×
UNCOV
253
    return ((SstArtifactRegistryObject *)obj)->registry.get();
×
254
}
255

256
// ---------------------------------------------------------------------------
257
// scan_files: parallel directory scan with size info
258
// ---------------------------------------------------------------------------
259

260
static PyObject *scan_files_fn(PyObject * /*self*/, PyObject *args,
×
261
                               PyObject *kwds) {
262
    static const char *kwlist[] = {"directory", "patterns", "recursive",
263
                                   "runtime", NULL};
264
    const char *directory;
265
    PyObject *patterns_obj = NULL;
×
266
    int recursive = 0;
×
UNCOV
267
    PyObject *runtime_arg = NULL;
×
268
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "s|OpO",
×
269
                                     const_cast<char **>(kwlist), &directory,
270
                                     &patterns_obj, &recursive, &runtime_arg)) {
UNCOV
271
        return NULL;
×
272
    }
273

UNCOV
274
    std::vector<std::string> patterns;
×
275
    if (patterns_obj && patterns_obj != Py_None) {
×
276
        if (!parse_string_seq(patterns_obj, "patterns must be a sequence",
×
277
                              patterns))
278
            return NULL;
×
279
    }
280

281
    Runtime *rt = nullptr;
×
UNCOV
282
    if (runtime_arg && runtime_arg != Py_None) {
×
283
        if (!PyObject_TypeCheck(runtime_arg, &RuntimeType)) {
×
UNCOV
284
            PyObject *native = PyObject_GetAttrString(runtime_arg, "_native");
×
285
            if (!native || !PyObject_TypeCheck(native, &RuntimeType)) {
×
UNCOV
286
                Py_XDECREF(native);
×
UNCOV
287
                PyErr_SetString(PyExc_TypeError,
×
288
                                "runtime must be a Runtime instance or None");
UNCOV
289
                return NULL;
×
290
            }
291
            rt = ((RuntimeObject *)native)->runtime.get();
×
292
            Py_DECREF(native);
×
293
        } else {
UNCOV
294
            rt = ((RuntimeObject *)runtime_arg)->runtime.get();
×
295
        }
296
    } else {
×
297
        rt = dftracer::utils::python::get_default_runtime();
×
298
    }
299

300
    PatternDirectoryScannerUtilityInput input(directory, patterns,
×
UNCOV
301
                                              recursive != 0, true);
×
UNCOV
302
    std::vector<FileEntry> entries;
×
UNCOV
303
    if (!run_blocking([&] {
×
UNCOV
304
            rt->submit(dftracer::utils::run_coro_scope(
×
UNCOV
305
                           rt->executor(),
×
306
                           [](dftracer::utils::CoroScope &scope,
×
307
                              PatternDirectoryScannerUtilityInput in,
308
                              std::vector<FileEntry> *out)
309
                               -> dftracer::utils::coro::CoroTask<void> {
×
310
                               PatternDirectoryScannerUtility scanner;
311
                               *out = co_await scanner(scope, in);
×
UNCOV
312
                           },
×
UNCOV
313
                           std::move(input), &entries),
×
314
                       "scan-files")
×
315
                .get();
×
316
        })) {
×
317
        return NULL;
×
318
    }
319

UNCOV
320
    PyObject *out = PyList_New(static_cast<Py_ssize_t>(entries.size()));
×
321
    if (!out) return NULL;
×
UNCOV
322
    for (std::size_t i = 0; i < entries.size(); ++i) {
×
323
        PyObject *t = Py_BuildValue("(sn)", entries[i].path.c_str(),
×
UNCOV
324
                                    (Py_ssize_t)entries[i].size);
×
325
        if (!t) {
×
326
            Py_DECREF(out);
×
UNCOV
327
            return NULL;
×
328
        }
UNCOV
329
        PyList_SET_ITEM(out, i, t);
×
330
    }
UNCOV
331
    return out;
×
332
}
×
333

334
// ---------------------------------------------------------------------------
335
// plan_lpt_partition: LPT bin-packing of (path, size) pairs
336
// ---------------------------------------------------------------------------
337

338
static PyObject *plan_lpt_partition_fn(PyObject * /*self*/, PyObject *args) {
×
339
    PyObject *entries_obj;
340
    Py_ssize_t num_workers;
UNCOV
341
    if (!PyArg_ParseTuple(args, "On", &entries_obj, &num_workers)) return NULL;
×
342
    if (num_workers <= 0) num_workers = 1;
×
343

344
    std::vector<FileEntry> entries;
×
345
    PyObject *seq = PySequence_Fast(entries_obj,
×
346
                                    "entries must be a sequence of "
347
                                    "(path, size) tuples");
348
    if (!seq) return NULL;
×
349
    Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
UNCOV
350
    entries.reserve(n);
×
351
    for (Py_ssize_t i = 0; i < n; ++i) {
×
UNCOV
352
        PyObject *item = PySequence_Fast_GET_ITEM(seq, i);
×
353
        const char *path = nullptr;
×
354
        Py_ssize_t size = 0;
×
355
        if (!PyArg_ParseTuple(item, "sn", &path, &size)) {
×
356
            Py_DECREF(seq);
×
357
            return NULL;
×
358
        }
UNCOV
359
        FileEntry fe;
×
UNCOV
360
        fe.path = path;
×
361
        fe.size = static_cast<std::size_t>(size);
×
362
        fe.is_regular_file = true;
×
UNCOV
363
        entries.push_back(std::move(fe));
×
364
    }
×
365
    Py_DECREF(seq);
×
366

367
    auto buckets = plan_lpt_partition(std::move(entries),
×
368
                                      static_cast<std::size_t>(num_workers));
×
369

370
    PyObject *out = PyList_New(static_cast<Py_ssize_t>(buckets.size()));
×
UNCOV
371
    if (!out) return NULL;
×
372
    for (std::size_t i = 0; i < buckets.size(); ++i) {
×
373
        PyObject *lst = PyList_New(static_cast<Py_ssize_t>(buckets[i].size()));
×
374
        if (!lst) {
×
375
            Py_DECREF(out);
×
UNCOV
376
            return NULL;
×
377
        }
378
        for (std::size_t j = 0; j < buckets[i].size(); ++j) {
×
UNCOV
379
            PyObject *t = Py_BuildValue("(sn)", buckets[i][j].path.c_str(),
×
380
                                        (Py_ssize_t)buckets[i][j].size);
×
UNCOV
381
            if (!t) {
×
382
                Py_DECREF(lst);
×
383
                Py_DECREF(out);
×
384
                return NULL;
×
385
            }
UNCOV
386
            PyList_SET_ITEM(lst, j, t);
×
387
        }
UNCOV
388
        PyList_SET_ITEM(out, i, lst);
×
389
    }
UNCOV
390
    return out;
×
UNCOV
391
}
×
392

393
// ---------------------------------------------------------------------------
394
// build_sst_batch: run the indexer pipeline with an SST sink and return
395
// the merged Artifacts dict.
396
// ---------------------------------------------------------------------------
397

UNCOV
398
static PyObject *build_sst_batch_fn(PyObject * /*self*/, PyObject *args,
×
399
                                    PyObject *kwds) {
400
    static const char *kwlist[] = {"files",
401
                                   "file_ids",
402
                                   "staging_dir",
403
                                   "batch_id",
404
                                   "index_dir",
405
                                   "checkpoint_size",
406
                                   "force_rebuild",
407
                                   "build_bloom",
408
                                   "bloom_dimensions",
409
                                   "parallelism",
410
                                   "flush_every_files",
411
                                   "runtime",
412
                                   "aggregation_config",
413
                                   "file_slices",
414
                                   "progress",
415
                                   NULL};
416
    PyObject *files_obj;
417
    PyObject *file_ids_obj;
418
    const char *staging_dir;
419
    const char *batch_id;
420
    const char *index_dir = "";
×
421
    Py_ssize_t checkpoint_size = static_cast<Py_ssize_t>(
×
422
        dftracer::utils::constants::indexer::DEFAULT_CHECKPOINT_SIZE);
423
    int force_rebuild = 0;
×
UNCOV
424
    int build_bloom = 1;
×
425
    PyObject *bloom_dims_obj = NULL;
×
UNCOV
426
    Py_ssize_t parallelism = 0;
×
UNCOV
427
    Py_ssize_t flush_every_files = 0;
×
UNCOV
428
    PyObject *runtime_arg = NULL;
×
UNCOV
429
    PyObject *aggregation_config_obj = NULL;
×
UNCOV
430
    PyObject *file_slices_obj = NULL;
×
431
    PyObject *progress_obj = NULL;
×
432

UNCOV
433
    if (!PyArg_ParseTupleAndKeywords(
×
434
            args, kwds, "OOss|snppOnnOOOO", const_cast<char **>(kwlist),
435
            &files_obj, &file_ids_obj, &staging_dir, &batch_id, &index_dir,
436
            &checkpoint_size, &force_rebuild, &build_bloom, &bloom_dims_obj,
437
            &parallelism, &flush_every_files, &runtime_arg,
438
            &aggregation_config_obj, &file_slices_obj, &progress_obj)) {
439
        return NULL;
×
440
    }
441

442
    // Unpack files.
443
    std::vector<std::string> files;
×
444
    {
445
        if (!parse_string_seq(files_obj, "files must be a sequence", files))
×
UNCOV
446
            return NULL;
×
447
    }
UNCOV
448
    if (files.empty()) {
×
UNCOV
449
        return PyDict_New();
×
450
    }
451

452
    // Unpack file_ids, parallel to files.
UNCOV
453
    std::vector<int> file_ids;
×
454
    {
455
        PyObject *seq =
456
            PySequence_Fast(file_ids_obj, "file_ids must be a sequence");
×
UNCOV
457
        if (!seq) return NULL;
×
UNCOV
458
        Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
459
        if (static_cast<std::size_t>(n) != files.size()) {
×
460
            Py_DECREF(seq);
×
461
            PyErr_SetString(PyExc_ValueError,
×
462
                            "file_ids must have the same length as files");
UNCOV
463
            return NULL;
×
464
        }
UNCOV
465
        file_ids.reserve(n);
×
466
        for (Py_ssize_t i = 0; i < n; ++i) {
×
UNCOV
467
            long v = PyLong_AsLong(PySequence_Fast_GET_ITEM(seq, i));
×
468
            if (v == -1 && PyErr_Occurred()) {
×
469
                Py_DECREF(seq);
×
470
                return NULL;
×
471
            }
UNCOV
472
            file_ids.push_back(static_cast<int>(v));
×
473
        }
474
        Py_DECREF(seq);
×
475
    }
476

477
    // Optional bloom dimensions override.
UNCOV
478
    std::vector<std::string> bloom_dims;
×
UNCOV
479
    if (bloom_dims_obj && bloom_dims_obj != Py_None) {
×
UNCOV
480
        if (!parse_string_seq(bloom_dims_obj,
×
481
                              "bloom_dimensions must be a sequence",
482
                              bloom_dims))
483
            return NULL;
×
484
    }
485

486
    // Resolve Runtime (matching CheckpointIndexer pattern).
487
    Runtime *rt = nullptr;
×
488
    if (runtime_arg && runtime_arg != Py_None) {
×
489
        if (PyObject_TypeCheck(runtime_arg, &RuntimeType)) {
×
490
            rt = ((RuntimeObject *)runtime_arg)->runtime.get();
×
491
        } else {
492
            PyObject *native = PyObject_GetAttrString(runtime_arg, "_native");
×
UNCOV
493
            if (!native || !PyObject_TypeCheck(native, &RuntimeType)) {
×
494
                Py_XDECREF(native);
×
UNCOV
495
                PyErr_SetString(PyExc_TypeError,
×
496
                                "runtime must be a Runtime instance or None");
UNCOV
497
                return NULL;
×
498
            }
UNCOV
499
            rt = ((RuntimeObject *)native)->runtime.get();
×
500
            Py_DECREF(native);
×
501
        }
502
    } else {
×
503
        rt = dftracer::utils::python::get_default_runtime();
×
504
    }
505

506
    // Build config + sink factory shared state.
507
    struct SharedArtifacts {
508
        std::mutex mu;
509
        std::vector<IndexDatabaseSstWriterContext::Artifacts> list;
510
    };
UNCOV
511
    auto artifacts = std::make_shared<SharedArtifacts>();
×
512
    auto staging = std::string(staging_dir);
×
UNCOV
513
    auto batch = std::string(batch_id);
×
514

515
    // Optional aggregation config, extracted from the Python dataclass.
516
    std::shared_ptr<dftracer::utils::trace::aggregators::AggregationConfig>
UNCOV
517
        agg_config_ptr;
×
UNCOV
518
    if (aggregation_config_obj && aggregation_config_obj != Py_None) {
×
519
        using dftracer::utils::trace::aggregators::AggregationConfig;
UNCOV
520
        auto cfg = std::make_shared<AggregationConfig>();
×
UNCOV
521
        auto pull_double = [&](const char *name, double fallback) -> double {
×
UNCOV
522
            PyObject *v = PyObject_GetAttrString(aggregation_config_obj, name);
×
UNCOV
523
            if (!v || v == Py_None) {
×
524
                Py_XDECREF(v);
×
525
                PyErr_Clear();
×
526
                return fallback;
×
527
            }
UNCOV
528
            double out = PyFloat_AsDouble(v);
×
529
            Py_DECREF(v);
UNCOV
530
            if (out == -1.0 && PyErr_Occurred()) return fallback;
×
531
            return out;
×
532
        };
×
UNCOV
533
        auto pull_bool = [&](const char *name, bool fallback) -> bool {
×
UNCOV
534
            PyObject *v = PyObject_GetAttrString(aggregation_config_obj, name);
×
535
            if (!v || v == Py_None) {
×
536
                Py_XDECREF(v);
×
537
                PyErr_Clear();
×
538
                return fallback;
×
539
            }
540
            int out = PyObject_IsTrue(v);
×
541
            Py_DECREF(v);
UNCOV
542
            return out > 0 ? true : fallback;
×
543
        };
×
544
        auto pull_string_list =
545
            [&](const char *name) -> std::vector<std::string> {
×
546
            std::vector<std::string> out;
×
547
            PyObject *v = PyObject_GetAttrString(aggregation_config_obj, name);
×
548
            if (!v || v == Py_None) {
×
549
                Py_XDECREF(v);
×
550
                PyErr_Clear();
×
551
                return out;
×
552
            }
553
            PyObject *seq = PySequence_Fast(v, "expected list of str");
×
554
            Py_DECREF(v);
×
555
            if (!seq) {
×
UNCOV
556
                PyErr_Clear();
×
557
                return out;
×
558
            }
UNCOV
559
            Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
560
            out.reserve(n);
×
561
            for (Py_ssize_t i = 0; i < n; ++i) {
×
562
                const char *s = as_utf8(PySequence_Fast_GET_ITEM(seq, i));
×
563
                if (s) out.emplace_back(s);
×
564
            }
565
            Py_DECREF(seq);
×
566
            return out;
×
UNCOV
567
        };
×
568
        double time_interval_ms = pull_double("time_interval_ms", 5000.0);
×
UNCOV
569
        cfg->time_interval_us =
×
570
            static_cast<std::uint64_t>(time_interval_ms * 1000.0);
×
571
        cfg->compute_percentiles = pull_bool("compute_percentiles", false);
×
572
        cfg->extra_group_keys = pull_string_list("group_keys");
×
UNCOV
573
        cfg->custom_metric_fields = pull_string_list("custom_metric_fields");
×
574
        agg_config_ptr = std::move(cfg);
×
575
    }
×
576

577
    // owned_member_maps must outlive rt->submit: FileSlice::members is raw.
578
    std::vector<std::vector<GzipMember>> owned_member_maps;
×
579
    std::vector<IndexBuildBatchConfig::FileSlice> parsed_slices;
×
UNCOV
580
    if (file_slices_obj && file_slices_obj != Py_None) {
×
581
        PyObject *seq =
582
            PySequence_Fast(file_slices_obj, "file_slices must be a sequence");
×
583
        if (!seq) return NULL;
×
584
        Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
585
        if (static_cast<std::size_t>(n) != files.size()) {
×
586
            Py_DECREF(seq);
×
587
            PyErr_SetString(PyExc_ValueError,
×
588
                            "file_slices must match files length");
589
            return NULL;
×
590
        }
591
        owned_member_maps.resize(n);
×
UNCOV
592
        parsed_slices.resize(n);
×
UNCOV
593
        for (Py_ssize_t i = 0; i < n; ++i) {
×
594
            PyObject *entry = PySequence_Fast_GET_ITEM(seq, i);
×
595
            if (entry == Py_None) {
×
596
                continue;  // leave slice default-constructed (members=null)
×
597
            }
598
            Py_ssize_t mb = 0, me = 0;
×
599
            int skip_scoped = 0;
×
600
            PyObject *members_obj = nullptr;
×
601
            if (!PyArg_ParseTuple(entry, "nnpO", &mb, &me, &skip_scoped,
×
602
                                  &members_obj)) {
603
                Py_DECREF(seq);
×
UNCOV
604
                return NULL;
×
605
            }
UNCOV
606
            PyObject *mseq = PySequence_Fast(
×
607
                members_obj, "file_slices[i].members must be a sequence");
608
            if (!mseq) {
×
609
                Py_DECREF(seq);
×
610
                return NULL;
×
611
            }
612
            Py_ssize_t mn = PySequence_Fast_GET_SIZE(mseq);
×
UNCOV
613
            auto &mv = owned_member_maps[i];
×
614
            mv.resize(mn);
×
615
            for (Py_ssize_t j = 0; j < mn; ++j) {
×
616
                PyObject *m = PySequence_Fast_GET_ITEM(mseq, j);
×
617
                unsigned long long c_offset = 0, c_size = 0;
×
UNCOV
618
                if (!PyArg_ParseTuple(m, "KK", &c_offset, &c_size)) {
×
619
                    Py_DECREF(mseq);
×
620
                    Py_DECREF(seq);
×
UNCOV
621
                    return NULL;
×
622
                }
UNCOV
623
                mv[j].c_offset = static_cast<std::uint64_t>(c_offset);
×
624
                mv[j].c_size = static_cast<std::uint64_t>(c_size);
×
625
            }
626
            Py_DECREF(mseq);
×
UNCOV
627
            parsed_slices[i].members = &mv;
×
628
            parsed_slices[i].member_begin = static_cast<std::size_t>(mb);
×
629
            parsed_slices[i].member_end = static_cast<std::size_t>(me);
×
630
            parsed_slices[i].skip_file_scoped_writes = skip_scoped != 0;
×
631
        }
632
        Py_DECREF(seq);
×
633
    }
634

UNCOV
635
    auto batch_config = std::make_shared<IndexBuildBatchConfig>();
×
UNCOV
636
    batch_config->file_paths = std::move(files);
×
637
    batch_config->preassigned_file_ids = std::move(file_ids);
×
UNCOV
638
    if (!parsed_slices.empty()) {
×
639
        batch_config->file_slices = parsed_slices;
×
640
    }
UNCOV
641
    batch_config->index_dir = index_dir;
×
UNCOV
642
    batch_config->build_bloom = build_bloom != 0;
×
643
    batch_config->checkpoint_size = static_cast<std::size_t>(checkpoint_size);
×
644
    batch_config->force_rebuild = force_rebuild != 0;
×
645
    batch_config->bloom_dimensions = std::move(bloom_dims);
×
646
    batch_config->parallelism =
×
647
        parallelism > 0 ? static_cast<std::size_t>(parallelism)
×
648
                        : (rt ? std::max<std::size_t>(rt->threads(), 1) : 1);
×
UNCOV
649
    batch_config->flush_every_files =
×
UNCOV
650
        static_cast<std::size_t>(flush_every_files);
×
UNCOV
651
    batch_config->rebuild_root_summaries = false;
×
652

653
    // The build runs with the GIL released; re-acquire it per call. The
654
    // GIL-holding deleter drops the ref safely after the build.
655
    if (progress_obj && progress_obj != Py_None) {
×
656
        Py_INCREF(progress_obj);
×
657
        std::shared_ptr<PyObject> cb(progress_obj, [](PyObject *p) {
×
UNCOV
658
            PyGILState_STATE g = PyGILState_Ensure();
×
659
            Py_DECREF(p);
660
            PyGILState_Release(g);
×
661
        });
×
662
        batch_config->progress = [cb](std::size_t done, std::size_t total) {
×
663
            PyGILState_STATE g = PyGILState_Ensure();
×
664
            PyObject *r = PyObject_CallFunction(cb.get(), "nn",
×
665
                                                static_cast<Py_ssize_t>(done),
666
                                                static_cast<Py_ssize_t>(total));
667
            if (r) {
×
668
                Py_DECREF(r);
669
            } else {
670
                // A failing progress callback must not abort the build.
671
                PyErr_Clear();
×
672
            }
673
            PyGILState_Release(g);
×
UNCOV
674
        };
×
UNCOV
675
    }
×
676

677
    if (agg_config_ptr) {
×
678
        auto agg_intern =
UNCOV
679
            dftracer::utils::trace::aggregators::intern_for_index(index_dir);
×
680
        // The fold writes aggregation SSTs through the batch build's own SST
681
        // sink (routed to aggregation.sst / system_metrics.sst), so they land
682
        // in `artifacts->list` with bloom/dict - no separate per-file sink.
UNCOV
683
        batch_config->agg_fold_factory =
×
UNCOV
684
            [agg_config_ptr,
×
685
             agg_intern](dftracer::utils::StringIntern &build_intern)
686
            -> std::unique_ptr<
687
                dftracer::utils::trace::views::detail::AggregationFold> {
688
            return std::make_unique<
689
                dftracer::utils::trace::views::detail::AggregationFold>(
690
                build_intern, agg_intern, *agg_config_ptr, /*config_hash=*/0);
×
691
        };
×
692
    }
×
693

694
    // Atomic: write phase calls sink_factory from N coroutines concurrently.
UNCOV
695
    auto batch_counter = std::make_shared<std::atomic<std::size_t>>(0);
×
UNCOV
696
    batch_config->sink_factory =
×
697
        [staging, batch, batch_counter]() -> std::unique_ptr<IndexBatchSink> {
×
698
        const std::size_t idx =
699
            batch_counter->fetch_add(1, std::memory_order_relaxed);
×
UNCOV
700
        std::string sub_batch = batch + "_" + std::to_string(idx);
×
701
        return std::make_unique<IndexDatabaseSstWriterContext>(staging,
×
702
                                                               sub_batch);
703
    };
×
UNCOV
704
    batch_config->sink_commit = [artifacts](IndexBatchSink &sink) {
×
705
        auto &sst = static_cast<IndexDatabaseSstWriterContext &>(sink);
×
706
        auto batch_artifacts = sst.commit();
×
707
        std::lock_guard<std::mutex> lock(artifacts->mu);
×
708
        if (!batch_artifacts.empty()) {
×
709
            artifacts->list.push_back(std::move(batch_artifacts));
×
710
        }
711
    };
×
712

713
    IndexBuildBatchResult result;
×
UNCOV
714
    if (!run_blocking([&] {
×
715
            rt->submit(dftracer::utils::run_coro_scope(
×
716
                           rt->executor(),
×
717
                           [](dftracer::utils::CoroScope &scope,
×
718
                              std::shared_ptr<IndexBuildBatchConfig> cfg,
719
                              IndexBuildBatchResult *out)
720
                               -> dftracer::utils::coro::CoroTask<void> {
×
721
                               *out =
×
722
                                   co_await IndexBatchBuilderUtility::process(
×
723
                                       &scope, std::move(cfg));
UNCOV
724
                           },
×
UNCOV
725
                           batch_config, &result),
×
726
                       "build-sst-batch")
×
727
                .get();
×
UNCOV
728
        })) {
×
729
        return NULL;
×
730
    }
731

732
    // If any file failed, surface the first error.
UNCOV
733
    if (result.failed > 0) {
×
UNCOV
734
        for (const auto &r : result.results) {
×
735
            if (!r.success) {
×
736
                PyErr_SetString(PyExc_RuntimeError, r.error_message.c_str());
×
737
                return NULL;
×
738
            }
739
        }
740
    }
741

742
    // One dict per committed sink + per-file aggregation below.
UNCOV
743
    PyObject *out_list = PyList_New(0);
×
UNCOV
744
    if (!out_list) return NULL;
×
745
    {
746
        std::lock_guard<std::mutex> lock(artifacts->mu);
×
UNCOV
747
        for (const auto &a : artifacts->list) {
×
748
            PyObject *main_dict = artifacts_to_dict(a);
×
749
            if (!main_dict || PyList_Append(out_list, main_dict) < 0) {
×
750
                Py_XDECREF(main_dict);
×
751
                Py_DECREF(out_list);
×
752
                return NULL;
×
753
            }
754
            Py_DECREF(main_dict);
×
755
        }
UNCOV
756
    }
×
757
    // The fold wrote its aggregation SSTs through the batch build's sink, so
758
    // they are already in `artifacts->list` above (no separate per-visitor
759
    // harvest). Combine the per-file trackers the folds produced out-of-band.
760
    using dftracer::utils::trace::aggregators::AssociationTracker;
UNCOV
761
    AssociationTracker combined;
×
UNCOV
762
    bool any_tracker = false;
×
UNCOV
763
    for (auto &ao : result.agg_outputs) {
×
UNCOV
764
        if (ao.tracker) {
×
UNCOV
765
            ao.tracker->finalize();
×
UNCOV
766
            combined.merge(*ao.tracker);
×
UNCOV
767
            any_tracker = true;
×
768
        }
769
    }
UNCOV
770
    PyObject *tracker_bytes = nullptr;
×
UNCOV
771
    if (any_tracker) {
×
UNCOV
772
        combined.finalize();
×
UNCOV
773
        std::string blob = combined.serialize();
×
774
        tracker_bytes = PyBytes_FromStringAndSize(
×
775
            blob.data(), static_cast<Py_ssize_t>(blob.size()));
×
776
    } else {
×
777
        tracker_bytes = PyBytes_FromStringAndSize(nullptr, 0);
×
778
    }
779
    if (!tracker_bytes) {
×
780
        Py_DECREF(out_list);
×
781
        return NULL;
×
782
    }
783
    PyObject *ret = PyTuple_Pack(2, out_list, tracker_bytes);
×
784
    Py_DECREF(out_list);
×
785
    Py_DECREF(tracker_bytes);
×
786
    return ret;
×
UNCOV
787
}
×
788

UNCOV
789
static PyObject *enable_aggregation_deterministic_ids_fn(PyObject * /*self*/,
×
790
                                                         PyObject * /*args*/) {
UNCOV
791
    dftracer::utils::trace::aggregators::enable_deterministic_intern_ids();
×
UNCOV
792
    Py_RETURN_NONE;
×
793
}
794

795
static PyObject *move_artifacts_fn(PyObject * /*self*/, PyObject *args,
×
796
                                   PyObject *kwds) {
797
    static const char *kwlist[] = {"artifacts", "dest_dir", NULL};
798
    PyObject *dict = NULL;
×
799
    const char *dest_dir = NULL;
×
UNCOV
800
    if (!PyArg_ParseTupleAndKeywords(
×
801
            args, kwds, "Os", const_cast<char **>(kwlist), &dict, &dest_dir)) {
802
        return NULL;
×
803
    }
804
    IndexDatabaseSstWriterContext::Artifacts a;
×
805
    if (!artifacts_from_dict(dict, &a)) return NULL;
×
806
    IndexDatabaseSstWriterContext::Artifacts moved;
×
807
    if (!run_blocking_r([&] { return std::move(a).move_to(dest_dir); }, moved))
×
808
        return NULL;
×
809
    return artifacts_to_dict(moved);
×
UNCOV
810
}
×
811

812
namespace {
813

UNCOV
814
dftracer::utils::coro::CoroTask<void> scan_one_gzip_file(
×
815
    std::string path, std::vector<GzipMember> *out) {
×
816
    out->clear();
817
    int fd = ::open(path.c_str(), O_RDONLY);
×
818
    if (fd < 0) co_return;
×
819
    struct stat st;
820
    if (::fstat(fd, &st) == 0 && st.st_size >= 18) {
×
821
        co_await enumerate_gzip_member_candidates(
×
822
            fd, static_cast<std::uint64_t>(st.st_size), *out);
823
    }
824
    ::close(fd);
×
825
}
×
826

827
}  // namespace
828

829
static PyObject *enumerate_gzip_members_fn(PyObject * /*self*/, PyObject *args,
×
830
                                           PyObject *kwds) {
831
    static const char *kwlist[] = {"files", "runtime", NULL};
832
    PyObject *files_obj = NULL;
×
833
    PyObject *runtime_arg = NULL;
×
834
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "O|O",
×
835
                                     const_cast<char **>(kwlist), &files_obj,
836
                                     &runtime_arg)) {
UNCOV
837
        return NULL;
×
838
    }
839

840
    std::vector<std::string> files;
×
841
    {
842
        if (!parse_string_seq(files_obj, "files must be a sequence", files))
×
843
            return NULL;
×
844
    }
845

UNCOV
846
    Runtime *rt = nullptr;
×
UNCOV
847
    if (runtime_arg && runtime_arg != Py_None) {
×
848
        if (PyObject_TypeCheck(runtime_arg, &RuntimeType)) {
×
UNCOV
849
            rt = ((RuntimeObject *)runtime_arg)->runtime.get();
×
850
        } else {
UNCOV
851
            PyObject *native = PyObject_GetAttrString(runtime_arg, "_native");
×
UNCOV
852
            if (!native || !PyObject_TypeCheck(native, &RuntimeType)) {
×
UNCOV
853
                Py_XDECREF(native);
×
UNCOV
854
                PyErr_SetString(PyExc_TypeError,
×
855
                                "runtime must be a Runtime instance or None");
UNCOV
856
                return NULL;
×
857
            }
UNCOV
858
            rt = ((RuntimeObject *)native)->runtime.get();
×
859
            Py_DECREF(native);
×
860
        }
UNCOV
861
    } else {
×
UNCOV
862
        rt = dftracer::utils::python::get_default_runtime();
×
863
    }
864

UNCOV
865
    std::vector<std::vector<GzipMember>> results(files.size());
×
866
    if (!run_blocking([&] {
×
867
            rt->submit(
×
868
                  dftracer::utils::run_coro_scope(
×
UNCOV
869
                      rt->executor(),
×
870
                      [](dftracer::utils::CoroScope &scope,
×
871
                         const std::vector<std::string> *paths,
872
                         std::vector<std::vector<GzipMember>> *out)
873
                          -> dftracer::utils::coro::CoroTask<void> {
×
874
                          co_await scope.scope(
×
875
                              [paths, out](dftracer::utils::CoroScope &child)
×
876
                                  -> dftracer::utils::coro::CoroTask<void> {
×
877
                                  for (std::size_t i = 0; i < paths->size();
×
878
                                       ++i) {
879
                                      const std::string &path = (*paths)[i];
880
                                      auto *slot = &(*out)[i];
881
                                      child.spawn(
×
UNCOV
882
                                          [path,
×
883
                                           slot](dftracer::utils::CoroScope &)
884
                                              -> dftracer::utils::coro::CoroTask<
885
                                                  void> {
×
886
                                              co_await scan_one_gzip_file(path,
×
887
                                                                          slot);
UNCOV
888
                                          });
×
889
                                  }
890
                                  co_return;
891
                              });
×
892
                          co_return;
893
                      },
×
UNCOV
894
                      &files, &results),
×
895
                  "enumerate-gzip-members")
×
896
                .get();
×
897
        })) {
×
898
        return NULL;
×
899
    }
900

UNCOV
901
    PyObject *out_list = PyList_New(static_cast<Py_ssize_t>(results.size()));
×
902
    if (!out_list) return NULL;
×
UNCOV
903
    for (std::size_t i = 0; i < results.size(); ++i) {
×
UNCOV
904
        const auto &mv = results[i];
×
905
        PyObject *inner = PyList_New(static_cast<Py_ssize_t>(mv.size()));
×
906
        if (!inner) {
×
907
            Py_DECREF(out_list);
×
UNCOV
908
            return NULL;
×
909
        }
910
        for (std::size_t j = 0; j < mv.size(); ++j) {
×
911
            PyObject *t =
912
                Py_BuildValue("(KK)", (unsigned long long)mv[j].c_offset,
×
UNCOV
913
                              (unsigned long long)mv[j].c_size);
×
914
            if (!t) {
×
915
                Py_DECREF(inner);
×
916
                Py_DECREF(out_list);
×
UNCOV
917
                return NULL;
×
918
            }
919
            PyList_SET_ITEM(inner, j, t);
×
920
        }
UNCOV
921
        PyList_SET_ITEM(out_list, i, inner);
×
922
    }
UNCOV
923
    return out_list;
×
UNCOV
924
}
×
925

926
// LPT (longest-processing-time) work assignment for balanced SST distribution.
927
// so the Dask backend produces identical work distribution to MPI.
UNCOV
928
static PyObject *plan_work_units_fn(PyObject * /*self*/, PyObject *args,
×
929
                                    PyObject *kwds) {
930
    static const char *kwlist[] = {"member_map", "num_workers", "target_c_size",
931
                                   NULL};
932
    PyObject *map_obj = NULL;
×
UNCOV
933
    Py_ssize_t num_workers = 0;
×
UNCOV
934
    unsigned long long target_c_size = 0;
×
935
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "On|K",
×
936
                                     const_cast<char **>(kwlist), &map_obj,
937
                                     &num_workers, &target_c_size)) {
938
        return NULL;
×
939
    }
940
    if (num_workers <= 0) num_workers = 1;
×
941

942
    std::vector<std::vector<GzipMember>> member_map;
×
943
    {
944
        PyObject *seq =
945
            PySequence_Fast(map_obj, "member_map must be a sequence");
×
946
        if (!seq) return NULL;
×
947
        Py_ssize_t n = PySequence_Fast_GET_SIZE(seq);
×
948
        member_map.resize(n);
×
949
        for (Py_ssize_t i = 0; i < n; ++i) {
×
950
            PyObject *inner = PySequence_Fast_GET_ITEM(seq, i);
×
951
            PyObject *iseq =
952
                PySequence_Fast(inner, "member_map[i] must be a sequence");
×
UNCOV
953
            if (!iseq) {
×
954
                Py_DECREF(seq);
×
UNCOV
955
                return NULL;
×
956
            }
957
            Py_ssize_t ni = PySequence_Fast_GET_SIZE(iseq);
×
958
            member_map[i].resize(ni);
×
UNCOV
959
            for (Py_ssize_t j = 0; j < ni; ++j) {
×
UNCOV
960
                PyObject *t = PySequence_Fast_GET_ITEM(iseq, j);
×
961
                unsigned long long c_offset = 0, c_size = 0;
×
UNCOV
962
                if (!PyTuple_Check(t)) {
×
963
                    PyErr_SetString(
×
964
                        PyExc_TypeError,
965
                        "each member must be an (offset, size) tuple");
966
                    Py_DECREF(iseq);
×
967
                    Py_DECREF(seq);
×
968
                    return NULL;
×
969
                }
UNCOV
970
                if (!PyArg_ParseTuple(t, "KK", &c_offset, &c_size)) {
×
971
                    Py_DECREF(iseq);
×
972
                    Py_DECREF(seq);
×
UNCOV
973
                    return NULL;
×
974
                }
UNCOV
975
                member_map[i][j].c_offset =
×
976
                    static_cast<std::uint64_t>(c_offset);
977
                member_map[i][j].c_size = static_cast<std::uint64_t>(c_size);
×
978
            }
979
            Py_DECREF(iseq);
×
980
        }
981
        Py_DECREF(seq);
×
982
    }
983

984
    // Fallback: treat empty/non-gzip files as a single whole-file member.
985
    std::uint64_t total_c = 0;
×
UNCOV
986
    for (auto &mv : member_map) {
×
UNCOV
987
        if (mv.empty()) mv.push_back({0, 0});
×
988
        for (const auto &m : mv) total_c += m.c_size;
×
989
    }
990

991
    if (target_c_size == 0) {
×
992
        target_c_size =
×
993
            (total_c + static_cast<std::uint64_t>(num_workers) - 1) /
×
UNCOV
994
            std::max<std::uint64_t>(static_cast<std::uint64_t>(num_workers), 1);
×
995
    }
996

997
    struct Unit {
998
        std::size_t file_idx;
999
        std::size_t member_begin;
1000
        std::size_t member_end;
1001
        std::uint64_t c_size;
1002
    };
1003
    std::vector<Unit> units;
×
1004
    for (std::size_t fi = 0; fi < member_map.size(); ++fi) {
×
1005
        const auto &members = member_map[fi];
×
UNCOV
1006
        if (members.empty()) continue;
×
UNCOV
1007
        std::size_t begin = 0;
×
1008
        std::uint64_t accum = 0;
×
UNCOV
1009
        for (std::size_t i = 0; i < members.size(); ++i) {
×
1010
            accum += members[i].c_size;
×
UNCOV
1011
            const bool is_last = (i + 1 == members.size());
×
1012
            if ((target_c_size > 0 && accum >= target_c_size) || is_last) {
×
UNCOV
1013
                units.push_back({fi, begin, i + 1, accum});
×
UNCOV
1014
                begin = i + 1;
×
UNCOV
1015
                accum = 0;
×
1016
            }
1017
        }
1018
    }
1019

1020
    std::vector<std::size_t> order(units.size());
×
1021
    for (std::size_t i = 0; i < order.size(); ++i) order[i] = i;
×
1022
    std::sort(order.begin(), order.end(), [&](std::size_t a, std::size_t b) {
×
1023
        if (units[a].c_size != units[b].c_size)
×
UNCOV
1024
            return units[a].c_size > units[b].c_size;
×
UNCOV
1025
        if (units[a].file_idx != units[b].file_idx)
×
1026
            return units[a].file_idx < units[b].file_idx;
×
1027
        return units[a].member_begin < units[b].member_begin;
×
1028
    });
1029
    const std::size_t nw = static_cast<std::size_t>(num_workers);
×
UNCOV
1030
    std::vector<std::uint64_t> loads(nw, 0);
×
UNCOV
1031
    std::vector<std::vector<std::size_t>> per_worker(nw);
×
UNCOV
1032
    for (std::size_t ord : order) {
×
UNCOV
1033
        std::size_t best = 0;
×
UNCOV
1034
        for (std::size_t r = 1; r < nw; ++r)
×
UNCOV
1035
            if (loads[r] < loads[best]) best = r;
×
UNCOV
1036
        per_worker[best].push_back(ord);
×
UNCOV
1037
        loads[best] += std::max<std::uint64_t>(units[ord].c_size, 1);
×
1038
    }
1039

1040
    PyObject *out = PyList_New(static_cast<Py_ssize_t>(nw));
×
1041
    if (!out) return NULL;
×
1042
    for (std::size_t w = 0; w < nw; ++w) {
×
1043
        // Keep per-worker slices sorted by (file_idx, member_begin) for
1044
        // deterministic, file-group-friendly iteration downstream.
1045
        auto &lst = per_worker[w];
×
1046
        std::sort(lst.begin(), lst.end(), [&](std::size_t a, std::size_t b) {
×
1047
            if (units[a].file_idx != units[b].file_idx)
×
1048
                return units[a].file_idx < units[b].file_idx;
×
1049
            return units[a].member_begin < units[b].member_begin;
×
1050
        });
UNCOV
1051
        PyObject *inner = PyList_New(static_cast<Py_ssize_t>(lst.size()));
×
UNCOV
1052
        if (!inner) {
×
1053
            Py_DECREF(out);
×
UNCOV
1054
            return NULL;
×
1055
        }
1056
        for (std::size_t k = 0; k < lst.size(); ++k) {
×
1057
            const auto &u = units[lst[k]];
×
1058
            PyObject *t = Py_BuildValue(
×
1059
                "(nnnK)", (Py_ssize_t)u.file_idx, (Py_ssize_t)u.member_begin,
×
1060
                (Py_ssize_t)u.member_end, (unsigned long long)u.c_size);
×
1061
            if (!t) {
×
1062
                Py_DECREF(inner);
×
1063
                Py_DECREF(out);
×
1064
                return NULL;
×
1065
            }
1066
            PyList_SET_ITEM(inner, k, t);
×
1067
        }
1068
        PyList_SET_ITEM(out, w, inner);
×
1069
    }
1070
    return out;
×
1071
}
×
1072

1073
// ---------------------------------------------------------------------------
1074
// Module registration
1075
// ---------------------------------------------------------------------------
1076

1077
static PyMethodDef SstDistributionMethods[] = {
1078
    {"build_sst_batch", DFTU_PYCFUNCTION(build_sst_batch_fn),
1079
     METH_VARARGS | METH_KEYWORDS,
1080
     "build_sst_batch(files, file_ids, staging_dir, batch_id, ...) "
1081
     "-> (list[dict], bytes)\n"
1082
     "Run the indexer pipeline with an SST sink and return "
1083
     "(artifact_dicts, tracker_blob). The tracker blob is the serialized "
1084
     "merged AssociationTracker from this batch's aggregation visitors "
1085
     "(empty bytes when no aggregation_config was passed)."},
1086
    {"plan_lpt_partition", DFTU_PYCFUNCTION(plan_lpt_partition_fn),
1087
     METH_VARARGS,
1088
     "plan_lpt_partition(entries, num_workers) -> list[list[(path, size)]]\n"
1089
     "Greedy Longest-Processing-Time-first bin-packing of (path, size) "
1090
     "tuples across num_workers buckets. Minimises the maximum per-worker "
1091
     "total size."},
1092
    {"scan_files", DFTU_PYCFUNCTION(scan_files_fn),
1093
     METH_VARARGS | METH_KEYWORDS,
1094
     "scan_files(directory, patterns=None, recursive=False, runtime=None) "
1095
     "-> list[(path, size)]\n"
1096
     "Parallel directory scan returning (path, size) tuples for regular "
1097
     "files matching the patterns."},
1098
    {"enable_aggregation_deterministic_ids",
1099
     DFTU_PYCFUNCTION(enable_aggregation_deterministic_ids_fn), METH_NOARGS,
1100
     "enable_aggregation_deterministic_ids() -> None\n"
1101
     "Flip the global aggregation StringIntern into deterministic-id mode "
1102
     "so the same string maps to the same 32-bit id in every worker "
1103
     "process. Call once at worker startup BEFORE any aggregation work."},
1104
    {"move_artifacts", DFTU_PYCFUNCTION(move_artifacts_fn),
1105
     METH_VARARGS | METH_KEYWORDS,
1106
     "move_artifacts(artifacts, dest_dir) -> dict\n"
1107
     "Move every populated SST in `artifacts` (as returned by "
1108
     "`build_sst_batch`) into `dest_dir` via the C++ rename/copy helper, "
1109
     "returning a fresh dict with the new paths. Single GIL release, no "
1110
     "per-file Python shutil.move overhead."},
1111
    {"enumerate_gzip_members", DFTU_PYCFUNCTION(enumerate_gzip_members_fn),
1112
     METH_VARARGS | METH_KEYWORDS,
1113
     "enumerate_gzip_members(files, runtime=None) -> list[list[(c_offset, "
1114
     "c_size)]]\n"
1115
     "Cooperative async scan of gzip member offsets across `files`. "
1116
     "Returns lists of (c_offset, c_size) parallel to `files`; empty for "
1117
     "non-gzip / unreadable files."},
1118
    {"plan_work_units", DFTU_PYCFUNCTION(plan_work_units_fn),
1119
     METH_VARARGS | METH_KEYWORDS,
1120
     "plan_work_units(member_map, num_workers, target_c_size=0) "
1121
     "-> list[list[(file_idx, member_begin, member_end, c_size)]]\n"
1122
     "Deterministic LPT assignment of intra-file gzip-member slices "
1123
     "across workers. Each worker's list contains (file_idx, "
1124
     "member_begin, member_end, c_size) tuples; a file sliced across "
1125
     "multiple workers appears in each owner's list with disjoint "
1126
     "[member_begin, member_end) ranges."},
1127
    {NULL, NULL, 0, NULL}};
1128

1129
int dftracer::utils::python::init_sst_distribution(PyObject *m) {
2 ✔
1130
    if (register_type(m, &SstArtifactRegistryType, "SstArtifactRegistry") < 0)
2 ✔
UNCOV
1131
        return -1;
×
1132
    if (PyModule_AddFunctions(m, SstDistributionMethods) < 0) return -1;
2 ✔
1133
    return 0;
2 ✔
1134
}
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