• 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

37.99
/src/dftracer/utils/python/streaming_iterator.cpp
1
#include <dftracer/utils/core/common/config.h>
2
#include <dftracer/utils/python/py_errors.h>
3
#include <dftracer/utils/python/py_runtime_mixin.h>
4
#include <dftracer/utils/python/py_type_helpers.h>
5
#ifdef DFTRACER_UTILS_ENABLE_ARROW
6

7
#define PY_SSIZE_T_CLEAN
8
#include <Python.h>
9
#include <dftracer/utils/python/dataframe.h>
10
#include <dftracer/utils/python/py_method.h>
11
#include <dftracer/utils/python/streaming_iterator.h>
12
#include <dftracer/utils/python/trace_reader_iterator.h>
13

14
namespace dftracer::utils::python {
15

16
static PyObject* ArrowStreamingIterator_new(PyTypeObject* type,
30 ✔
17
                                            PyObject* /*args*/,
18
                                            PyObject* /*kwds*/) {
19
    ArrowStreamingIteratorObject* self =
15 ✔
20
        (ArrowStreamingIteratorObject*)type->tp_alloc(type, 0);
30 ✔
21
    if (self) {
30 ✔
22
        // Allocate C++ state separately to avoid layout issues
23
        self->cpp_state = new ArrowStreamingIteratorState();
30 ✔
24
    }
15 ✔
25
    return (PyObject*)self;
30 ✔
26
}
27

28
static void ArrowStreamingIterator_dealloc(ArrowStreamingIteratorObject* self) {
30 ✔
29
    if (self->cpp_state) {
30 ✔
30
        // Cancel the stream if still running
31
        if (self->cpp_state->cancel) {
30 ✔
32
            self->cpp_state->cancel();
30 ✔
33
        }
15 ✔
34
        delete self->cpp_state;
30 ✔
35
        self->cpp_state = nullptr;
30 ✔
36
    }
15 ✔
37
    Py_TYPE(self)->tp_free((PyObject*)self);
30 ✔
38
}
30 ✔
39

40
static PyObject* ArrowStreamingIterator_iter(PyObject* self) {
30 !
41
    Py_INCREF(self);
15 ✔
42
    return self;
30 ✔
43
}
44

45
static PyObject* ArrowStreamingIterator_next(
94 ✔
46
    ArrowStreamingIteratorObject* self) {
47
    if (!self->cpp_state ||
188 !
48
        (!self->cpp_state->pull_next && !self->cpp_state->pull_df)) {
94 !
UNCOV
49
        PyErr_SetString(PyExc_RuntimeError, "Iterator not initialized");
×
UNCOV
50
        return NULL;
×
51
    }
52

53
    // Native-DataFrame path: pull a DataFrame chunk (GIL released) and wrap it,
54
    // so the stream never crosses Arrow.
55
    if (self->cpp_state->pull_df) {
94 !
56
        std::optional<dftracer::utils::dataframe::DataFrame> df;
94 ✔
57
        if (!run_blocking_r([&] { return self->cpp_state->pull_df(); }, df))
188 !
UNCOV
58
            return NULL;
×
59
        if (!df.has_value()) {
94 ✔
60
            if (self->cpp_state->get_error) {
30 !
61
                if (auto ex = self->cpp_state->get_error()) {
30 !
62
                    try {
63
                        std::rethrow_exception(ex);
×
64
                    } catch (const std::exception& e) {
×
65
                        set_typed_py_error(e);
×
UNCOV
66
                        return NULL;
×
67
                    } catch (...) {
×
68
                        PyErr_SetString(PyExc_RuntimeError,
×
69
                                        "Unknown error in streaming iterator");
UNCOV
70
                        return NULL;
×
UNCOV
71
                    }
×
72
                }
15 !
73
            }
15 ✔
74
            return NULL;  // StopIteration
30 ✔
75
        }
76
        return wrap_dataframe(std::move(*df));
64 !
77
    }
94 ✔
78

UNCOV
79
    std::optional<ArrowExportResult> result;
×
UNCOV
80
    if (!run_blocking_r([&] { return self->cpp_state->pull_next(); }, result))
×
UNCOV
81
        return NULL;
×
82

UNCOV
83
    if (!result.has_value()) {
×
84
        // Check for error
UNCOV
85
        if (self->cpp_state->get_error) {
×
86
            auto ex = self->cpp_state->get_error();
×
87
            if (ex) {
×
88
                try {
89
                    std::rethrow_exception(ex);
×
UNCOV
90
                } catch (const std::exception& e) {
×
UNCOV
91
                    set_typed_py_error(e);
×
UNCOV
92
                    return NULL;
×
UNCOV
93
                } catch (...) {
×
UNCOV
94
                    PyErr_SetString(PyExc_RuntimeError,
×
95
                                    "Unknown error in streaming iterator");
UNCOV
96
                    return NULL;
×
UNCOV
97
                }
×
98
            }
UNCOV
99
        }
×
100
        // Normal completion
UNCOV
101
        return NULL;  // StopIteration
×
102
    }
103

104
    // Wrap the ArrowExportResult in an ArrowBatchCapsule
105
    ArrowBatchCapsuleObject* obj =
UNCOV
106
        (ArrowBatchCapsuleObject*)ArrowBatchCapsuleType.tp_alloc(
×
107
            &ArrowBatchCapsuleType, 0);
UNCOV
108
    if (!obj) return NULL;
×
UNCOV
109
    obj->result = new ArrowExportResult(std::move(*result));
×
UNCOV
110
    return (PyObject*)obj;
×
111
}
47 ✔
112

UNCOV
113
static PyObject* ArrowStreamingIterator_cancel(
×
114
    ArrowStreamingIteratorObject* self, PyObject* Py_UNUSED(args)) {
UNCOV
115
    if (self->cpp_state && self->cpp_state->cancel) {
×
UNCOV
116
        self->cpp_state->cancel();
×
117
    }
UNCOV
118
    Py_RETURN_NONE;
×
119
}
120

121
static PyMethodDef ArrowStreamingIterator_methods[] = {
122
    {"cancel", DFTU_PYCFUNCTION(ArrowStreamingIterator_cancel), METH_NOARGS,
123
     "Cancel the streaming iterator."},
124
    {NULL}};
125

126
PyTypeObject ArrowStreamingIteratorType = {
127
    PyVarObject_HEAD_INIT(NULL, 0) "dftracer_utils_ext._ArrowStreamingIterator",
128
    sizeof(ArrowStreamingIteratorObject),       /* tp_basicsize */
129
    0,                                          /* tp_itemsize */
130
    (destructor)ArrowStreamingIterator_dealloc, /* tp_dealloc */
131
    0,                                          /* tp_vectorcall_offset */
132
    0,                                          /* tp_getattr */
133
    0,                                          /* tp_setattr */
134
    0,                                          /* tp_as_async */
135
    0,                                          /* tp_repr */
136
    0,                                          /* tp_as_number */
137
    0,                                          /* tp_as_sequence */
138
    0,                                          /* tp_as_mapping */
139
    0,                                          /* tp_hash */
140
    0,                                          /* tp_call */
141
    0,                                          /* tp_str */
142
    0,                                          /* tp_getattro */
143
    0,                                          /* tp_setattro */
144
    0,                                          /* tp_as_buffer */
145
    Py_TPFLAGS_DEFAULT,                         /* tp_flags */
146
    "Streaming Arrow batch iterator.\n\n"
147
    "Yields ArrowBatch objects as they become available from the C++ "
148
    "pipeline.\n"
149
    "Call cancel() to stop the stream early.", /* tp_doc */
150
    0,                                         /* tp_traverse */
151
    0,                                         /* tp_clear */
152
    0,                                         /* tp_richcompare */
153
    0,                                         /* tp_weaklistoffset */
154
    ArrowStreamingIterator_iter,               /* tp_iter */
155
    (iternextfunc)ArrowStreamingIterator_next, /* tp_iternext */
156
    ArrowStreamingIterator_methods,            /* tp_methods */
157
    0,                                         /* tp_members */
158
    0,                                         /* tp_getset */
159
    0,                                         /* tp_base */
160
    0,                                         /* tp_dict */
161
    0,                                         /* tp_descr_get */
162
    0,                                         /* tp_descr_set */
163
    0,                                         /* tp_dictoffset */
164
    0,                                         /* tp_init */
165
    0,                                         /* tp_alloc */
166
    ArrowStreamingIterator_new,                /* tp_new */
167
};
168

169
int init_arrow_streaming_iterator(PyObject* m) {
2 ✔
170
    if (register_type(m, &ArrowStreamingIteratorType,
2 ✔
171
                      "_ArrowStreamingIterator") < 0)
2 ✔
UNCOV
172
        return -1;
×
173

174
    return 0;
2 ✔
175
}
1 ✔
176

177
}  // namespace dftracer::utils::python
178

179
#endif  // DFTRACER_UTILS_ENABLE_ARROW
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