• 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

2.2
/src/dftracer/utils/python/trace_reader_iterator.cpp
1
#define PY_SSIZE_T_CLEAN
2
#include <Python.h>
3
#include <dftracer/utils/core/common/config.h>
4
#include <dftracer/utils/python/batch_byte_size.h>
5
#include <dftracer/utils/python/json.h>
6
#include <dftracer/utils/python/py_errors.h>
7
#include <dftracer/utils/python/py_method.h>
8
#include <dftracer/utils/python/py_type_helpers.h>
9
#include <dftracer/utils/python/trace_reader_iterator.h>
10
#ifdef DFTRACER_UTILS_ENABLE_ARROW
11
#include <nanoarrow/nanoarrow.h>
12

13
using ArrowExportResult =
14
    dftracer::utils::utilities::common::arrow::ArrowExportResult;
15

UNCOV
16
static void release_arrow_schema(PyObject *capsule) {
×
17
    auto *schema = static_cast<ArrowSchema *>(
UNCOV
18
        PyCapsule_GetPointer(capsule, "arrow_schema"));
×
19
    if (schema && schema->release) {
×
UNCOV
20
        schema->release(schema);
×
21
    }
UNCOV
22
    delete schema;
×
UNCOV
23
}
×
24

UNCOV
25
static void release_arrow_array(PyObject *capsule) {
×
26
    auto *array =
UNCOV
27
        static_cast<ArrowArray *>(PyCapsule_GetPointer(capsule, "arrow_array"));
×
28
    if (array && array->release) {
×
UNCOV
29
        array->release(array);
×
30
    }
UNCOV
31
    delete array;
×
UNCOV
32
}
×
33

UNCOV
34
static PyObject *ArrowBatchCapsule_arrow_c_array(ArrowBatchCapsuleObject *self,
×
35
                                                 PyObject *args) {
UNCOV
36
    PyObject *requested_schema = Py_None;
×
UNCOV
37
    if (!PyArg_ParseTuple(args, "|O", &requested_schema)) return NULL;
×
38

39
    if (!self->result || !self->result->valid()) {
×
UNCOV
40
        PyErr_SetString(PyExc_RuntimeError,
×
41
                        "Arrow data already exported via __arrow_c_array__. "
42
                        "Each batch can only be exported once.");
UNCOV
43
        return NULL;
×
44
    }
45

UNCOV
46
    auto *schema = new ArrowSchema;
×
UNCOV
47
    auto *array = new ArrowArray;
×
48

UNCOV
49
    self->result->release_schema().move(schema);
×
UNCOV
50
    self->result->release_array().move(array);
×
51

52
    PyObject *schema_capsule =
UNCOV
53
        PyCapsule_New(schema, "arrow_schema", release_arrow_schema);
×
54
    if (!schema_capsule) {
×
55
        if (schema->release) schema->release(schema);
×
56
        delete schema;
×
57
        if (array->release) array->release(array);
×
58
        delete array;
×
UNCOV
59
        return NULL;
×
60
    }
61

62
    PyObject *array_capsule =
UNCOV
63
        PyCapsule_New(array, "arrow_array", release_arrow_array);
×
UNCOV
64
    if (!array_capsule) {
×
65
        Py_DECREF(schema_capsule);
66
        if (array->release) array->release(array);
×
67
        delete array;
×
UNCOV
68
        return NULL;
×
69
    }
70

UNCOV
71
    PyObject *tuple = PyTuple_Pack(2, schema_capsule, array_capsule);
×
72
    Py_DECREF(schema_capsule);
73
    Py_DECREF(array_capsule);
UNCOV
74
    return tuple;
×
75
}
76

UNCOV
77
static PyObject *ArrowBatchCapsule_get_num_rows(ArrowBatchCapsuleObject *self,
×
78
                                                void *) {
UNCOV
79
    if (!self->result || !self->result->valid()) return PyLong_FromLong(0);
×
UNCOV
80
    return PyLong_FromLongLong(self->result->num_rows());
×
81
}
82

UNCOV
83
static PyObject *ArrowBatchCapsule_get_num_columns(
×
84
    ArrowBatchCapsuleObject *self, void *) {
UNCOV
85
    if (!self->result || !self->result->valid()) return PyLong_FromLong(0);
×
UNCOV
86
    return PyLong_FromLongLong(self->result->num_columns());
×
87
}
88

UNCOV
89
static void ArrowBatchCapsule_dealloc(ArrowBatchCapsuleObject *self) {
×
UNCOV
90
    delete self->result;
×
UNCOV
91
    Py_TYPE(self)->tp_free((PyObject *)self);
×
UNCOV
92
}
×
93

94
static PyMethodDef ArrowBatchCapsule_methods[] = {
95
    {"__arrow_c_array__", DFTU_PYCFUNCTION(ArrowBatchCapsule_arrow_c_array),
96
     METH_VARARGS,
97
     "Export as Arrow C Data Interface PyCapsule pair (schema, array)"},
98
    {NULL}};
99

100
static PyGetSetDef ArrowBatchCapsule_getsetters[] = {
101
    {"num_rows", (getter)ArrowBatchCapsule_get_num_rows, NULL, "Number of rows",
102
     NULL},
103
    {"num_columns", (getter)ArrowBatchCapsule_get_num_columns, NULL,
104
     "Number of columns", NULL},
105
    {NULL}};
106

107
PyTypeObject ArrowBatchCapsuleType = {
108
    PyVarObject_HEAD_INIT(NULL, 0) "dftracer_utils_ext._ArrowBatchCapsule",
109
    sizeof(ArrowBatchCapsuleObject),       /* tp_basicsize */
110
    0,                                     /* tp_itemsize */
111
    (destructor)ArrowBatchCapsule_dealloc, /* tp_dealloc */
112
    0,                                     /* tp_vectorcall_offset */
113
    0,                                     /* tp_getattr */
114
    0,                                     /* tp_setattr */
115
    0,                                     /* tp_as_async */
116
    0,                                     /* tp_repr */
117
    0,                                     /* tp_as_number */
118
    0,                                     /* tp_as_sequence */
119
    0,                                     /* tp_as_mapping */
120
    0,                                     /* tp_hash */
121
    0,                                     /* tp_call */
122
    0,                                     /* tp_str */
123
    0,                                     /* tp_getattro */
124
    0,                                     /* tp_setattro */
125
    0,                                     /* tp_as_buffer */
126
    Py_TPFLAGS_DEFAULT,                    /* tp_flags */
127
    "Internal Arrow batch wrapper implementing __arrow_c_array__ protocol",
128
    0,                                     /* tp_traverse */
129
    0,                                     /* tp_clear */
130
    0,                                     /* tp_richcompare */
131
    0,                                     /* tp_weaklistoffset */
132
    0,                                     /* tp_iter */
133
    0,                                     /* tp_iternext */
134
    ArrowBatchCapsule_methods,             /* tp_methods */
135
    0,                                     /* tp_members */
136
    ArrowBatchCapsule_getsetters,          /* tp_getset */
137
};
138

139
#endif                                     // DFTRACER_UTILS_ENABLE_ARROW
140

UNCOV
141
static void cancel_and_wait_batch_state(MemoryViewBatchIteratorState *bs) {
×
UNCOV
142
    bs->cancelled.store(true, std::memory_order_release);
×
UNCOV
143
    if (bs->channel) bs->channel->close();
×
UNCOV
144
    if (bs->task_future.valid()) bs->task_future.wait();
×
UNCOV
145
}
×
146

UNCOV
147
static void cancel_and_wait_json_dict_state(JsonDictIteratorState *js) {
×
UNCOV
148
    js->cancelled.store(true, std::memory_order_release);
×
UNCOV
149
    if (js->channel) js->channel->close();
×
UNCOV
150
    if (js->task_future.valid()) js->task_future.wait();
×
UNCOV
151
}
×
152

UNCOV
153
static void TraceReaderIterator_dealloc(TraceReaderIteratorObject *self) {
×
154
#ifdef DFTRACER_UTILS_ENABLE_ARROW
UNCOV
155
    if (self->arrow_state) {
×
UNCOV
156
        self->arrow_state->cancelled.store(true, std::memory_order_release);
×
UNCOV
157
        if (self->arrow_state->channel) self->arrow_state->channel->close();
×
UNCOV
158
        Py_BEGIN_ALLOW_THREADS if (self->arrow_state->task_future.valid()) {
×
UNCOV
159
            self->arrow_state->task_future.wait();
×
160
        }
UNCOV
161
        Py_END_ALLOW_THREADS self->arrow_state.reset();
×
162
    }
163
#endif
UNCOV
164
    if (self->json_dict_state) {
×
UNCOV
165
        Py_BEGIN_ALLOW_THREADS cancel_and_wait_json_dict_state(
×
166
            self->json_dict_state.get());
UNCOV
167
        Py_END_ALLOW_THREADS self->json_dict_state.reset();
×
168
    }
UNCOV
169
    if (self->batch_state) {
×
UNCOV
170
        Py_BEGIN_ALLOW_THREADS cancel_and_wait_batch_state(
×
171
            self->batch_state.get());
UNCOV
172
        Py_END_ALLOW_THREADS self->batch_state.reset();
×
173
    }
UNCOV
174
    Py_XDECREF(self->current_batch);
×
UNCOV
175
    self->current_batch = NULL;
×
UNCOV
176
    Py_TYPE(self)->tp_free((PyObject *)self);
×
UNCOV
177
}
×
178

UNCOV
179
static PyObject *TraceReaderIterator_iter(TraceReaderIteratorObject *self) {
×
180
    Py_INCREF(self);
UNCOV
181
    return (PyObject *)self;
×
182
}
183

UNCOV
184
static PyObject *TraceReaderIterator_next(TraceReaderIteratorObject *self) {
×
UNCOV
185
    if (self->mode == IteratorMode::JSON_DICT) {
×
186
        while (true) {
UNCOV
187
            if (self->json_dict_current_batch) {
×
UNCOV
188
                auto &events = self->json_dict_current_batch->events;
×
UNCOV
189
                Py_ssize_t n = static_cast<Py_ssize_t>(events.size());
×
UNCOV
190
                if (self->json_dict_index < n) {
×
191
                    JsonDictValueObject *obj =
UNCOV
192
                        (JsonDictValueObject *)JsonDictValueType.tp_alloc(
×
193
                            &JsonDictValueType, 0);
UNCOV
194
                    if (!obj) return NULL;
×
UNCOV
195
                    new (&obj->batch) std::shared_ptr<JsonDictBatch>(
×
UNCOV
196
                        self->json_dict_current_batch);
×
UNCOV
197
                    obj->event_index =
×
UNCOV
198
                        static_cast<std::size_t>(self->json_dict_index);
×
UNCOV
199
                    obj->is_args = false;
×
UNCOV
200
                    self->json_dict_index++;
×
UNCOV
201
                    return (PyObject *)obj;
×
202
                }
UNCOV
203
                self->json_dict_current_batch.reset();
×
UNCOV
204
                self->json_dict_index = 0;
×
205
            }
206

UNCOV
207
            auto *js = self->json_dict_state.get();
×
UNCOV
208
            std::optional<JsonDictBatch> batch;
×
UNCOV
209
            Py_BEGIN_ALLOW_THREADS batch = js->channel->blocking_receive();
×
UNCOV
210
            Py_END_ALLOW_THREADS
×
211

UNCOV
212
                if (!batch.has_value()) {
×
UNCOV
213
                std::lock_guard<std::mutex> lock(js->error_mtx);
×
UNCOV
214
                if (js->error) {
×
215
                    try {
216
                        std::rethrow_exception(js->error);
×
217
                    } catch (const std::exception &e) {
×
218
                        dftracer::utils::python::set_typed_py_error(e);
×
219
                        return NULL;
×
220
                    } catch (...) {
×
UNCOV
221
                        PyErr_SetString(PyExc_RuntimeError,
×
222
                                        "Unknown error in json dict iterator");
223
                        return NULL;
×
UNCOV
224
                    }
×
225
                }
UNCOV
226
                return NULL;
×
UNCOV
227
            }
×
228

UNCOV
229
            auto dequeued_bytes = dftracer::utils::python::byte_size(*batch);
×
UNCOV
230
            js->bytes_in_queue.fetch_sub(dequeued_bytes,
×
231
                                         std::memory_order_acq_rel);
232
            self->json_dict_current_batch =
UNCOV
233
                std::make_shared<JsonDictBatch>(std::move(*batch));
×
UNCOV
234
            self->json_dict_index = 0;
×
UNCOV
235
        }
×
236
    }
237

238
#ifdef DFTRACER_UTILS_ENABLE_ARROW
UNCOV
239
    if (self->mode == IteratorMode::ARROW) {
×
UNCOV
240
        auto *astate = self->arrow_state.get();
×
UNCOV
241
        std::optional<ArrowExportResult> batch;
×
UNCOV
242
        Py_BEGIN_ALLOW_THREADS batch = astate->channel->blocking_receive();
×
UNCOV
243
        Py_END_ALLOW_THREADS
×
244

UNCOV
245
            if (!batch.has_value()) {
×
UNCOV
246
            std::lock_guard<std::mutex> lock(astate->error_mtx);
×
UNCOV
247
            if (astate->error) {
×
248
                try {
249
                    std::rethrow_exception(astate->error);
×
250
                } catch (const std::exception &e) {
×
251
                    dftracer::utils::python::set_typed_py_error(e);
×
252
                    return NULL;
×
253
                } catch (...) {
×
UNCOV
254
                    PyErr_SetString(PyExc_RuntimeError,
×
255
                                    "Unknown error in Arrow iterator");
256
                    return NULL;
×
UNCOV
257
                }
×
258
            }
UNCOV
259
            return NULL;
×
UNCOV
260
        }
×
261

UNCOV
262
        auto dequeued_bytes = dftracer::utils::python::byte_size(*batch);
×
UNCOV
263
        astate->bytes_in_queue.fetch_sub(dequeued_bytes,
×
264
                                         std::memory_order_acq_rel);
265

266
        ArrowBatchCapsuleObject *obj =
UNCOV
267
            (ArrowBatchCapsuleObject *)ArrowBatchCapsuleType.tp_alloc(
×
268
                &ArrowBatchCapsuleType, 0);
UNCOV
269
        if (!obj) return NULL;
×
UNCOV
270
        obj->result = new ArrowExportResult(std::move(*batch));
×
UNCOV
271
        return (PyObject *)obj;
×
UNCOV
272
    }
×
273
#endif
274

275
    using namespace dftracer::utils::python;
276
    while (true) {
UNCOV
277
        if (self->current_batch) {
×
UNCOV
278
            auto *batch_obj = (MemoryViewBatchObject *)self->current_batch;
×
279
            Py_ssize_t n =
UNCOV
280
                static_cast<Py_ssize_t>(batch_obj->data->num_entries());
×
UNCOV
281
            if (self->batch_index < n) {
×
282
                PyObject *mv =
UNCOV
283
                    MemoryViewBatch_item(batch_obj, self->batch_index);
×
UNCOV
284
                self->batch_index++;
×
UNCOV
285
                return mv;
×
286
            }
UNCOV
287
            Py_DECREF(self->current_batch);
×
UNCOV
288
            self->current_batch = NULL;
×
UNCOV
289
            self->batch_index = 0;
×
290
        }
291

UNCOV
292
        auto *bs = self->batch_state.get();
×
UNCOV
293
        std::optional<MemoryViewBatchData> batch_data;
×
UNCOV
294
        Py_BEGIN_ALLOW_THREADS batch_data = bs->channel->blocking_receive();
×
UNCOV
295
        Py_END_ALLOW_THREADS
×
296

UNCOV
297
            if (!batch_data.has_value()) {
×
UNCOV
298
            std::lock_guard<std::mutex> lock(bs->error_mtx);
×
UNCOV
299
            if (bs->error) {
×
300
                try {
UNCOV
301
                    std::rethrow_exception(bs->error);
×
UNCOV
302
                } catch (const std::exception &e) {
×
UNCOV
303
                    dftracer::utils::python::set_typed_py_error(e);
×
UNCOV
304
                    return NULL;
×
305
                } catch (...) {
×
UNCOV
306
                    PyErr_SetString(PyExc_RuntimeError,
×
307
                                    "Unknown error in batch iterator");
UNCOV
308
                    return NULL;
×
UNCOV
309
                }
×
310
            }
UNCOV
311
            return NULL;
×
UNCOV
312
        }
×
313

UNCOV
314
        auto dequeued_bytes = dftracer::utils::python::byte_size(*batch_data);
×
UNCOV
315
        bs->bytes_in_queue.fetch_sub(dequeued_bytes, std::memory_order_acq_rel);
×
316

UNCOV
317
        auto *obj = (MemoryViewBatchObject *)MemoryViewBatchType.tp_alloc(
×
318
            &MemoryViewBatchType, 0);
UNCOV
319
        if (!obj) return NULL;
×
UNCOV
320
        obj->data = new MemoryViewBatchData(std::move(*batch_data));
×
UNCOV
321
        self->current_batch = (PyObject *)obj;
×
UNCOV
322
        self->batch_index = 0;
×
UNCOV
323
    }
×
324
}
325

326
PyTypeObject TraceReaderIteratorType = {
327
    PyVarObject_HEAD_INIT(NULL, 0) "dftracer_utils_ext.TraceReaderIterator",
328
    sizeof(TraceReaderIteratorObject),       /* tp_basicsize */
329
    0,                                       /* tp_itemsize */
330
    (destructor)TraceReaderIterator_dealloc, /* tp_dealloc */
331
    0,                                       /* tp_vectorcall_offset */
332
    0,                                       /* tp_getattr */
333
    0,                                       /* tp_setattr */
334
    0,                                       /* tp_as_async */
335
    0,                                       /* tp_repr */
336
    0,                                       /* tp_as_number */
337
    0,                                       /* tp_as_sequence */
338
    0,                                       /* tp_as_mapping */
339
    0,                                       /* tp_hash */
340
    0,                                       /* tp_call */
341
    0,                                       /* tp_str */
342
    0,                                       /* tp_getattro */
343
    0,                                       /* tp_setattro */
344
    0,                                       /* tp_as_buffer */
345
    Py_TPFLAGS_DEFAULT,                      /* tp_flags */
346
    "Lazy iterator over TraceReader lines or raw chunks",
347
    0,                                       /* tp_traverse */
348
    0,                                       /* tp_clear */
349
    0,                                       /* tp_richcompare */
350
    0,                                       /* tp_weaklistoffset */
351
    (getiterfunc)TraceReaderIterator_iter,   /* tp_iter */
352
    (iternextfunc)TraceReaderIterator_next,  /* tp_iternext */
353
    0,                                       /* tp_methods */
354
    0,                                       /* tp_members */
355
    0,                                       /* tp_getset */
356
    0,                                       /* tp_base */
357
    0,                                       /* tp_dict */
358
    0,                                       /* tp_descr_get */
359
    0,                                       /* tp_descr_set */
360
    0,                                       /* tp_dictoffset */
361
    0,                                       /* tp_init */
362
    0,                                       /* tp_alloc */
363
    0,                                       /* tp_new */
364
};
365

366
int dftracer::utils::python::init_trace_reader_iterator(PyObject *m) {
2 ✔
367
    if (register_type(m, &TraceReaderIteratorType, "TraceReaderIterator") < 0)
2 ✔
UNCOV
368
        return -1;
×
369

370
#ifdef DFTRACER_UTILS_ENABLE_ARROW
371
    if (register_type(m, &ArrowBatchCapsuleType, "_ArrowBatchCapsule") < 0)
2 ✔
UNCOV
372
        return -1;
×
373
#endif
374

375
    return 0;
2 ✔
376
}
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