• 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

1.45
/src/dftracer/utils/python/arrow_stream_capsule.cpp
1
#include <dftracer/utils/core/common/config.h>
2
#include <dftracer/utils/python/py_type_helpers.h>
3
#ifdef DFTRACER_UTILS_ENABLE_ARROW
4

5
#define PY_SSIZE_T_CLEAN
6
#include <Python.h>
7
#include <dftracer/utils/python/arrow_stream_capsule.h>
8
#include <dftracer/utils/python/batch_byte_size.h>
9
#include <dftracer/utils/python/py_method.h>
10
#include <dftracer/utils/python/schema_reconcile.h>
11
#include <nanoarrow/nanoarrow.h>
12

13
#include <cerrno>
14
#include <cstring>
15
#include <deque>
16
#include <exception>
17
#include <mutex>
18
#include <optional>
19
#include <string>
20

21
using ArrowExportResult =
22
    dftracer::utils::utilities::common::arrow::ArrowExportResult;
23

24
namespace {
25

26
// Drain until K consecutive batches add no new columns, bounded by MAX.
27
constexpr int STABLE_BATCHES = 5;
28
constexpr int MAX_DRAIN = 128;
29

30
struct StreamPrivate {
31
    std::shared_ptr<ArrowIteratorState> state;
32
    dftracer::utils::python::SchemaReconciler reconciler;
33
    // Drained during discovery, emitted first from get_next.
34
    std::deque<ArrowExportResult> pending;
35
    std::string last_error;
36
    bool initialized = false;
37
    // Sticky: once set, all entry points short-circuit to EIO.
38
    bool error_set = false;
39
};
40

41
static void mark_error(StreamPrivate *p, std::string msg) {
×
42
    if (p->last_error.empty()) p->last_error = std::move(msg);
×
43
    p->error_set = true;
×
44
    p->initialized = true;
×
UNCOV
45
}
×
46

UNCOV
47
static int initialize_stream(StreamPrivate *p) {
×
UNCOV
48
    if (p->error_set) return EIO;
×
UNCOV
49
    if (p->initialized) return 0;
×
UNCOV
50
    auto *astate = p->state.get();
×
51

UNCOV
52
    int stable_run = 0;
×
UNCOV
53
    int drained = 0;
×
UNCOV
54
    while (stable_run < STABLE_BATCHES && drained < MAX_DRAIN) {
×
UNCOV
55
        auto batch = astate->channel->blocking_receive();
×
UNCOV
56
        if (!batch.has_value()) {
×
57
            // End-of-stream or producer error before discovery converged.
UNCOV
58
            std::lock_guard<std::mutex> lock(astate->error_mtx);
×
UNCOV
59
            if (astate->error) {
×
60
                try {
61
                    std::rethrow_exception(astate->error);
×
62
                } catch (const std::exception &e) {
×
63
                    mark_error(p, e.what());
×
64
                } catch (...) {
×
65
                    mark_error(p, "unknown error in Arrow stream");
×
66
                }
×
UNCOV
67
                return EIO;
×
68
            }
UNCOV
69
            break;  // clean early EOS; finalize with whatever we have
×
UNCOV
70
        }
×
UNCOV
71
        auto dequeued = dftracer::utils::python::byte_size(*batch);
×
UNCOV
72
        astate->bytes_in_queue.fetch_sub(dequeued, std::memory_order_acq_rel);
×
73

UNCOV
74
        bool added = p->reconciler.merge(batch->get_schema());
×
75
        if (!p->reconciler.last_error().empty()) {
×
76
            mark_error(p, p->reconciler.last_error());
×
UNCOV
77
            return EIO;
×
78
        }
UNCOV
79
        p->pending.push_back(std::move(*batch));
×
UNCOV
80
        stable_run = added ? 0 : (stable_run + 1);
×
UNCOV
81
        ++drained;
×
UNCOV
82
    }
×
83

84
    if (p->reconciler.finalize() != 0) {
×
85
        mark_error(p, p->reconciler.last_error().empty()
×
86
                          ? "failed to finalize schema union"
×
87
                          : p->reconciler.last_error());
×
UNCOV
88
        return EIO;
×
89
    }
UNCOV
90
    p->initialized = true;
×
UNCOV
91
    return 0;
×
92
}
93

UNCOV
94
static int stream_get_schema(struct ArrowArrayStream *s,
×
95
                             struct ArrowSchema *out) {
UNCOV
96
    auto *p = static_cast<StreamPrivate *>(s->private_data);
×
UNCOV
97
    int rc = initialize_stream(p);
×
UNCOV
98
    if (rc != 0) return rc;
×
UNCOV
99
    if (p->error_set) return EIO;
×
100
    if (p->reconciler.copy_schema(out) != 0) {
×
101
        mark_error(p, p->reconciler.last_error().empty()
×
102
                          ? "failed to copy locked schema"
×
103
                          : p->reconciler.last_error());
×
UNCOV
104
        return EIO;
×
105
    }
UNCOV
106
    return 0;
×
107
}
108

UNCOV
109
static int stream_get_next(struct ArrowArrayStream *s, struct ArrowArray *out) {
×
UNCOV
110
    auto *p = static_cast<StreamPrivate *>(s->private_data);
×
UNCOV
111
    if (p->error_set) return EIO;
×
112
    if (!p->initialized) {
×
113
        int rc = initialize_stream(p);
×
UNCOV
114
        if (rc != 0) return rc;
×
115
    }
116

117
    // Drain any discovery-phase batches first, then pull from the channel.
UNCOV
118
    std::optional<ArrowExportResult> batch;
×
UNCOV
119
    if (!p->pending.empty()) {
×
UNCOV
120
        batch = std::move(p->pending.front());
×
UNCOV
121
        p->pending.pop_front();
×
122
    } else {
UNCOV
123
        auto *astate = p->state.get();
×
UNCOV
124
        batch = astate->channel->blocking_receive();
×
UNCOV
125
        if (!batch.has_value()) {
×
UNCOV
126
            std::lock_guard<std::mutex> lock(astate->error_mtx);
×
UNCOV
127
            if (astate->error) {
×
128
                try {
129
                    std::rethrow_exception(astate->error);
×
130
                } catch (const std::exception &e) {
×
131
                    mark_error(p, e.what());
×
132
                } catch (...) {
×
133
                    mark_error(p, "unknown error in Arrow stream");
×
134
                }
×
UNCOV
135
                return EIO;
×
136
            }
137
            // End of stream per Arrow C spec: return success with
138
            // out->release == nullptr.
UNCOV
139
            out->release = nullptr;
×
UNCOV
140
            return 0;
×
141
        }
×
142
        auto dequeued = dftracer::utils::python::byte_size(*batch);
×
UNCOV
143
        astate->bytes_in_queue.fetch_sub(dequeued, std::memory_order_acq_rel);
×
144
    }
145

UNCOV
146
    if (p->reconciler.reconcile(batch->get_schema(), batch->get_array(), out) !=
×
147
        0) {
148
        mark_error(p, p->reconciler.last_error().empty()
×
149
                          ? "schema reconciliation failed"
×
150
                          : p->reconciler.last_error());
×
UNCOV
151
        return EIO;
×
152
    }
UNCOV
153
    return 0;
×
UNCOV
154
}
×
155

156
static const char *stream_get_last_error(struct ArrowArrayStream *s) {
×
157
    auto *p = static_cast<StreamPrivate *>(s->private_data);
×
158
    if (!p || p->last_error.empty()) return nullptr;
×
UNCOV
159
    return p->last_error.c_str();
×
160
}
161

UNCOV
162
static void stream_release(struct ArrowArrayStream *s) {
×
UNCOV
163
    auto *p = static_cast<StreamPrivate *>(s->private_data);
×
UNCOV
164
    if (p) {
×
UNCOV
165
        if (p->state) {
×
UNCOV
166
            p->state->cancelled.store(true, std::memory_order_release);
×
UNCOV
167
            if (p->state->channel) p->state->channel->close();
×
UNCOV
168
            if (p->state->task_future.valid()) {
×
169
                // Release the GIL if this callback was invoked from a
170
                // Python-holding context (e.g. capsule destructor during
171
                // GC). If the GIL is not held (pyarrow's C reader path),
172
                // _PyThreadState_UncheckedGet() returns null and we wait
173
                // without touching the Python thread state.
UNCOV
174
                if (Py_IsInitialized() && PyGILState_Check()) {
×
UNCOV
175
                    Py_BEGIN_ALLOW_THREADS p->state->task_future.wait();
×
UNCOV
176
                    Py_END_ALLOW_THREADS
×
177
                } else {
UNCOV
178
                    p->state->task_future.wait();
×
179
                }
180
            }
181
        }
UNCOV
182
        delete p;
×
183
    }
UNCOV
184
    s->private_data = nullptr;
×
UNCOV
185
    s->release = nullptr;
×
UNCOV
186
}
×
187

UNCOV
188
static void release_stream_capsule(PyObject *capsule) {
×
189
    auto *stream = static_cast<ArrowArrayStream *>(
UNCOV
190
        PyCapsule_GetPointer(capsule, "arrow_array_stream"));
×
UNCOV
191
    if (stream && stream->release) {
×
UNCOV
192
        stream->release(stream);
×
193
    }
UNCOV
194
    delete stream;
×
UNCOV
195
}
×
196

UNCOV
197
static PyObject *ArrowBatchStream_arrow_c_stream(ArrowBatchStreamObject *self,
×
198
                                                 PyObject *args) {
UNCOV
199
    PyObject *requested_schema = Py_None;
×
UNCOV
200
    if (!PyArg_ParseTuple(args, "|O", &requested_schema)) return NULL;
×
201

202
    // Per the PyCapsule protocol, a non-None `requested_schema` means the
203
    // caller wants the stream cast to that schema. We only emit our native
204
    // schema today; reject explicitly so misuse fails loudly instead of
205
    // silently returning arrays that don't match what the caller asked for.
206
    if (requested_schema != Py_None) {
×
UNCOV
207
        PyErr_SetString(PyExc_NotImplementedError,
×
208
                        "iter_arrow_stream does not support "
209
                        "requested_schema casting; pass None to use the "
210
                        "native schema.");
UNCOV
211
        return NULL;
×
212
    }
213

UNCOV
214
    if (self->consumed || !self->state) {
×
UNCOV
215
        PyErr_SetString(PyExc_RuntimeError,
×
216
                        "Arrow stream already exported via "
217
                        "__arrow_c_stream__; each stream can be "
218
                        "exported only once.");
UNCOV
219
        return NULL;
×
220
    }
221

UNCOV
222
    auto *priv = new StreamPrivate;
×
UNCOV
223
    priv->state = self->state;
×
UNCOV
224
    self->consumed = true;
×
UNCOV
225
    self->state.reset();
×
226

UNCOV
227
    auto *stream = new ArrowArrayStream;
×
UNCOV
228
    std::memset(stream, 0, sizeof(*stream));
×
UNCOV
229
    stream->get_schema = stream_get_schema;
×
UNCOV
230
    stream->get_next = stream_get_next;
×
UNCOV
231
    stream->get_last_error = stream_get_last_error;
×
UNCOV
232
    stream->release = stream_release;
×
UNCOV
233
    stream->private_data = priv;
×
234

235
    PyObject *capsule =
UNCOV
236
        PyCapsule_New(stream, "arrow_array_stream", release_stream_capsule);
×
237
    if (!capsule) {
×
238
        stream->release(stream);
×
239
        delete stream;
×
UNCOV
240
        return NULL;
×
241
    }
UNCOV
242
    return capsule;
×
243
}
244

UNCOV
245
static void ArrowBatchStream_dealloc(ArrowBatchStreamObject *self) {
×
UNCOV
246
    if (self->state) {
×
UNCOV
247
        self->state->cancelled.store(true, std::memory_order_release);
×
UNCOV
248
        if (self->state->channel) self->state->channel->close();
×
UNCOV
249
        Py_BEGIN_ALLOW_THREADS if (self->state->task_future.valid()) {
×
UNCOV
250
            self->state->task_future.wait();
×
251
        }
UNCOV
252
        Py_END_ALLOW_THREADS
×
253
    }
UNCOV
254
    self->state.~shared_ptr<ArrowIteratorState>();
×
UNCOV
255
    Py_TYPE(self)->tp_free((PyObject *)self);
×
UNCOV
256
}
×
257

258
static PyMethodDef ArrowBatchStream_methods[] = {
259
    {"__arrow_c_stream__", DFTU_PYCFUNCTION(ArrowBatchStream_arrow_c_stream),
260
     METH_VARARGS, "Export as Arrow C Data Interface stream PyCapsule"},
261
    {NULL}};
262

263
}  // namespace
264

265
PyTypeObject ArrowBatchStreamType = {
266
    PyVarObject_HEAD_INIT(NULL, 0) "dftracer_utils_ext._ArrowBatchStream",
267
    sizeof(ArrowBatchStreamObject), /* tp_basicsize */
268
    0,                              /* tp_itemsize */
269
    (destructor)ArrowBatchStream_dealloc,
270
    0,
271
    0,
272
    0,
273
    0,
274
    0,
275
    0,
276
    0,
277
    0,
278
    0,
279
    0,
280
    0,
281
    0,
282
    0,
283
    0,
284
    Py_TPFLAGS_DEFAULT,
285
    "Zero-iteration Arrow stream backed by a C++ coroutine channel",
286
    0,
287
    0,
288
    0,
289
    0,
290
    0,
291
    0,
292
    ArrowBatchStream_methods,
293
    0,
294
    0,
295
    0,
296
    0,
297
    0,
298
    0,
299
    0,
300
    0,
301
    0,
302
    0,
303
};
304

305
int dftracer::utils::python::init_arrow_batch_stream(PyObject *m) {
2 ✔
306
    if (register_type(m, &ArrowBatchStreamType, "_ArrowBatchStream") < 0)
2 ✔
UNCOV
307
        return -1;
×
308
    return 0;
2 ✔
309
}
1 ✔
310

UNCOV
311
PyObject *dftracer::utils::python::make_arrow_batch_stream(
×
312
    std::shared_ptr<ArrowIteratorState> state) {
UNCOV
313
    auto *obj = (ArrowBatchStreamObject *)ArrowBatchStreamType.tp_alloc(
×
314
        &ArrowBatchStreamType, 0);
UNCOV
315
    if (!obj) return NULL;
×
UNCOV
316
    new (&obj->state) std::shared_ptr<ArrowIteratorState>(std::move(state));
×
UNCOV
317
    obj->consumed = false;
×
UNCOV
318
    return (PyObject *)obj;
×
319
}
320

321
#endif
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