• 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

60.93
/src/dftracer/utils/python/trace_viewer.cpp
1
#define PY_SSIZE_T_CLEAN
2
#include <dftracer/utils/core/common/config.h>
3
#include <dftracer/utils/core/common/filesystem.h>
4
#include <dftracer/utils/core/runtime.h>
5
#include <dftracer/utils/dataframe/dataframe.h>
6
#include <dftracer/utils/dataframe/internal/lazy_plan.h>
7
#include <dftracer/utils/plugins/plugins.h>
8
#include <dftracer/utils/python/dataframe.h>
9
#include <dftracer/utils/python/lazyframe.h>
10
#include <dftracer/utils/python/plugin_host.h>
11
#include <dftracer/utils/python/py_dict_helpers.h>
12
#include <dftracer/utils/python/py_errors.h>
13
#include <dftracer/utils/python/py_list_helpers.h>
14
#include <dftracer/utils/python/py_method.h>
15
#include <dftracer/utils/python/py_runtime_mixin.h>
16
#include <dftracer/utils/python/py_seq_helpers.h>
17
#include <dftracer/utils/python/py_str_helpers.h>
18
#include <dftracer/utils/python/py_type_helpers.h>
19
#include <dftracer/utils/python/trace_viewer.h>
20
#include <dftracer/utils/query/query.h>
21
#include <dftracer/utils/trace/internal/utils.h>
22
#include <dftracer/utils/trace/time_metric.h>
23
#include <dftracer/utils/trace/views/view.h>
24
#include <dftracer/utils/trace/views/view_source.h>
25
#include <dftracer/utils/utilities/fileio/compress/libdeflate_gzip.h>
26

27
#include <cctype>
28
#include <cstdint>
29
#include <cstdio>
30
#include <cstring>
31
#include <functional>
32
#include <memory>
33
#include <optional>
34
#include <string>
35
#include <string_view>
36
#include <utility>
37
#include <vector>
38

39
namespace {
40

41
namespace py = dftracer::utils::python;
42
namespace views = dftracer::utils::trace::views;
43
using dftracer::utils::dataframe::DataFrame;
44
using dftracer::utils::dataframe::LazyFrame;
45
using views::AggOp;
46
using views::AggSpec;
47
using views::ExportStats;
48
using views::GroupKey;
49
using views::Phase;
50
using views::View;
51

52
View& tv_of(PyObject* self) {
3,132 ✔
53
    return *reinterpret_cast<TraceViewerObject*>(self)->tv;
3,132 ✔
54
}
55

56
PyObject* make_viewer(View&& t) {
1,172 ✔
57
    auto* self = reinterpret_cast<TraceViewerObject*>(
586 ✔
58
        TraceViewerType.tp_alloc(&TraceViewerType, 0));
1,172 ✔
59
    if (!self) return nullptr;
1,172 ✔
60
    self->tv = new View(std::move(t));
1,172 ✔
61
    return reinterpret_cast<PyObject*>(self);
1,172 ✔
62
}
586 ✔
63

64
// Runs `fn` with the GIL held, turning a C++ exception into the typed Python
65
// error.
66
template <class Fn>
67
PyObject* guarded(Fn&& fn) {
2,976 ✔
68
    try {
69
        return fn();
2,976 !
70
    } catch (const std::exception& e) {
10 !
71
        py::set_typed_py_error(e);
10 !
72
        return nullptr;
10 ✔
73
    }
5 !
74
}
1,493 ✔
75

76
template <class Fn>
77
PyObject* build(PyObject* self, Fn&& fn) {
1,174 ✔
78
    return guarded([&] { return make_viewer(fn(tv_of(self))); });
2,348 !
79
}
80

81
PyObject* stats_dict(const ExportStats& s) {
12 ✔
82
    PyObject* d = PyDict_New();
12 ✔
83
    if (!d) return nullptr;
12 ✔
84
    if (dict_set_i64(d, "events_matched",
30 ✔
85
                     static_cast<long long>(s.events_matched)) < 0 ||
18 !
86
        dict_set_i64(d, "events_scanned",
18 !
87
                     static_cast<long long>(s.events_scanned)) < 0 ||
18 !
88
        dict_set_i64(d, "chunks_scanned",
18 !
89
                     static_cast<long long>(s.chunks_scanned)) < 0 ||
18 !
90
        dict_set_i64(d, "chunks_skipped",
18 !
91
                     static_cast<long long>(s.chunks_skipped)) < 0 ||
18 !
92
        dict_set_i64(d, "chunks_covered",
18 !
93
                     static_cast<long long>(s.chunks_covered)) < 0 ||
18 !
94
        dict_set_bool(d, "artifacts_committed", s.artifacts_committed) < 0 ||
12 !
95
        dict_set_bool(d, "truncated", s.truncated) < 0 ||
24 !
96
        dict_set_bool(d, "served_from_mv", s.served_from_mv) < 0) {
12 !
97
        Py_DECREF(d);
NEW
98
        return nullptr;
×
99
    }
100
    return d;
12 ✔
101
}
6 ✔
102

103
// A Python (done, total) callback as a ProgressFn. It fires on a runtime
104
// worker, so each call takes the GIL.
105
views::ProgressFn progress_from(PyObject* obj) {
82 ✔
106
    if (!obj || obj == Py_None) return {};
82 !
107
    Py_INCREF(obj);
1 ✔
108
    std::shared_ptr<PyObject> cb(obj, [](PyObject* p) {
3 ✔
109
        PyGILState_STATE g = PyGILState_Ensure();
2 ✔
110
        Py_DECREF(p);
1 ✔
111
        PyGILState_Release(g);
2 ✔
112
    });
2 !
113
    return [cb](std::size_t done, std::size_t total) {
3 !
114
        PyGILState_STATE g = PyGILState_Ensure();
2 ✔
115
        PyObject* r =
1 ✔
116
            PyObject_CallFunction(cb.get(), "nn", static_cast<Py_ssize_t>(done),
3 ✔
117
                                  static_cast<Py_ssize_t>(total));
1 ✔
118
        if (r)
2 !
119
            Py_DECREF(r);
1 ✔
120
        else
NEW
121
            PyErr_Clear();
×
122
        PyGILState_Release(g);
2 ✔
123
    };
2 !
124
}
42 ✔
125

126
class FileSink : public views::ExportSink {
127
   public:
128
    explicit FileSink(FILE* f) : f_(f) {}
48 ✔
129
    ~FileSink() override {
48 ✔
130
        if (f_) std::fclose(f_);
32 !
131
    }
48 ✔
132
    void write(std::string_view data) override {
8,268 ✔
133
        std::fwrite(data.data(), 1, data.size(), f_);
8,268 ✔
134
    }
8,268 ✔
135
    void flush() override { std::fflush(f_); }
32 ✔
136

137
   private:
138
    FILE* f_;
139
};
140

141
// Buffers writes and flushes whole-line gzip members past MEMBER_TARGET, so
142
// the output is a multi-member re-indexable trace.
143
class GzipSink : public views::ExportSink {
144
   public:
145
    GzipSink(FILE* f, int level) : f_(f), comp_(level) {}
3 !
146
    ~GzipSink() override {
5 ✔
147
        if (!buf_.empty()) flush(buf_.size());
2 !
148
        std::fclose(f_);
2 !
149
    }
5 ✔
150
    void write(std::string_view data) override {
800 ✔
151
        buf_.append(data);
800 ✔
152
        while (buf_.size() >= MEMBER_TARGET) {
800 ✔
NEW
153
            std::size_t cut = buf_.rfind('\n', buf_.size());
×
NEW
154
            if (cut == std::string::npos || cut + 1 < MEMBER_TARGET) break;
×
NEW
155
            flush(cut + 1);
×
156
        }
157
    }
800 ✔
158

159
   private:
160
    static constexpr std::size_t MEMBER_TARGET = 4 * 1024 * 1024;
161
    void flush(std::size_t n) {
2 ✔
162
        if (comp_.compress_member_into(scratch_, buf_.data(), n))
2 ✔
163
            std::fwrite(scratch_.data(), 1, scratch_.size(), f_);
2 ✔
164
        buf_.erase(0, n);
2 ✔
165
    }
2 ✔
166
    FILE* f_;
167
    dftracer::utils::utilities::fileio::compress::GzipMemberCompressor comp_;
168
    std::string buf_;
169
    std::vector<std::uint8_t> scratch_;
170
};
171

172
// "fn(key)" or "fn(key, 'a', 'b')" -> transform + inner key text; false when
173
// `t` is not a call, leaving `t` to parse as a plain key.
174
bool split_transform(const std::string& t, GroupKey::Transform& tf,
444 ✔
175
                     std::vector<std::string>& targs, std::string& inner) {
176
    const auto lp = t.find('(');
444 ✔
177
    if (lp == std::string::npos || t.back() != ')') return false;
444 !
178
    const std::string fn = t.substr(0, lp);
14 !
179
    if (fn == "dirname")
14 ✔
180
        tf = GroupKey::Transform::Dirname;
6 ✔
181
    else if (fn == "basename")
8 ✔
182
        tf = GroupKey::Transform::Basename;
2 ✔
183
    else if (fn == "lower")
6 ✔
184
        tf = GroupKey::Transform::Lower;
2 ✔
185
    else if (fn == "bucket")
4 !
186
        tf = GroupKey::Transform::Bucket;
4 ✔
187
    else
NEW
188
        return false;
×
189
    std::vector<std::string> parts;
14 ✔
190
    std::string cur;
14 ✔
191
    for (char ch : t.substr(lp + 1, t.size() - lp - 2)) {
200 !
192
        if (ch == ',') {
186 ✔
193
            parts.push_back(cur);
8 !
194
            cur.clear();
8 ✔
195
        } else {
4 ✔
196
            cur += ch;
178 !
197
        }
198
    }
7 ✔
199
    parts.push_back(cur);
14 !
200
    auto trim = [](const std::string& x) {
29 ✔
201
        const auto b = x.find_first_not_of(" \t'\"");
22 ✔
202
        const auto e = x.find_last_not_of(" \t'\"");
22 ✔
203
        return b == std::string::npos ? std::string() : x.substr(b, e - b + 1);
22 ✔
204
    };
205
    inner = trim(parts[0]);
14 !
206
    for (std::size_t i = 1; i < parts.size(); ++i)
22 ✔
207
        targs.push_back(trim(parts[i]));
8 !
208
    return !inner.empty();
14 ✔
209
}
229 ✔
210

211
// name | cat | pid | tid | fhash | hhash | io_cat | acc_pat | rank |
212
// file_path | file_name | host_name | arg:<key> | any other field, each
213
// optionally wrapped in a transform call.
214
GroupKey parse_group_key(const std::string& s) {
444 ✔
215
    std::string t = s;
444 !
216
    GroupKey::Transform tf = GroupKey::Transform::None;
444 ✔
217
    std::vector<std::string> targs;
444 ✔
218
    std::string inner;
444 ✔
219
    if (split_transform(t, tf, targs, inner)) t = inner;
444 !
220
    GroupKey out;
444 ✔
221
    if (t == "name")
444 ✔
222
        out = GroupKey::name();
88 !
223
    else if (t == "cat")
356 ✔
224
        out = GroupKey::cat();
196 !
225
    else if (t == "pid")
160 ✔
226
        out = GroupKey::pid();
38 !
227
    else if (t == "tid")
122 ✔
228
        out = GroupKey::tid();
16 !
229
    else if (t == "fhash")
106 ✔
230
        out = GroupKey::fhash();
12 !
231
    else if (t == "hhash")
94 ✔
232
        out = GroupKey::hhash();
10 !
233
    else if (t == "io_cat")
84 ✔
234
        out = GroupKey::io_cat();
16 !
235
    else if (t == "acc_pat")
68 !
NEW
236
        out = GroupKey::acc_pat();
×
237
    else if (t == "file_path" || t == "resolved.fpath" || t == "r.fpath")
68 !
238
        out = GroupKey::file_path();
26 !
239
    else if (t == "file_name")
42 ✔
240
        out = GroupKey::file_name();
4 !
241
    else if (t == "host_name" || t == "resolved.hostname" ||
65 ✔
242
             t == "r.hostname" || t == "resolved.host" || t == "r.host")
59 !
243
        out = GroupKey::host_name();
12 !
244
    else if (t == "rank")
26 ✔
245
        out = GroupKey::rank();
4 !
246
    else if (t.rfind("arg:", 0) == 0)
22 !
NEW
247
        out = GroupKey::of_arg(t.substr(4));
×
248
    else
249
        out = GroupKey::field(t);
22 !
250
    out.transform = tf;
444 ✔
251
    out.transform_args = std::move(targs);
444 ✔
252
    return out;
666 ✔
253
}
444 !
254

255
// "count" | "op:field" | "argmax:field:by" | "pct:field:q" | "pNN:field"
256
std::optional<AggSpec> parse_agg_spec(const std::string& t) {
874 ✔
257
    auto c1 = t.find(':');
874 ✔
258
    std::string op = t.substr(0, c1);
874 !
259
    std::string field, by;
874 ✔
260
    if (c1 != std::string::npos) {
874 ✔
261
        std::string rest = t.substr(c1 + 1);
504 !
262
        auto c2 = rest.find(':');
504 ✔
263
        field = rest.substr(0, c2);
504 !
264
        if (c2 != std::string::npos) by = rest.substr(c2 + 1);
504 !
265
    }
504 ✔
266
    const std::string occ = field.empty() ? "dur" : field;
874 !
267
    if (op == "count") return AggSpec(AggOp::Count, "", "", "");
874 !
268
    if (op == "sum") return AggSpec(AggOp::Sum, field);
546 !
269
    if (op == "sumsq") return AggSpec(AggOp::SumSq, field);
432 !
270
    if (op == "min") return AggSpec(AggOp::Min, field);
406 !
271
    if (op == "max") return AggSpec(AggOp::Max, field);
318 !
272
    if (op == "mean") return AggSpec(AggOp::Mean, field);
228 !
273
    if (op == "var") return AggSpec(AggOp::Var, field);
118 !
274
    if (op == "std") return AggSpec(AggOp::Std, field);
116 !
275
    if (op == "skew") return AggSpec(AggOp::Skew, field);
66 !
276
    if (op == "kurt") return AggSpec(AggOp::Kurt, field);
60 !
277
    if (op == "hist") return AggSpec(AggOp::Hist, field);
54 !
278
    if (op == "argmax") return AggSpec(AggOp::ArgMax, field, "", by);
50 !
279
    if (op == "set_union" || op == "uniq")
50 !
280
        return AggSpec(AggOp::SetUnion, field);
2 !
281
    if (op == "busy") return AggSpec(AggOp::Busy, occ);
48 !
282
    if (op == "concurrency") return AggSpec(AggOp::Concurrency, occ);
32 !
283
    if (op == "utilization") return AggSpec(AggOp::Utilization, occ);
20 !
284
    if (op == "active") return AggSpec(AggOp::Active, occ);
18 !
285
    if (op == "pct")
16 !
NEW
286
        return AggSpec(AggOp::Pct, field, "", "",
×
NEW
287
                       by.empty() ? 0.0 : std::stod(by));
×
288
    if (op.size() >= 2 && op[0] == 'p') {
16 !
289
        double denom = 1.0;
16 ✔
290
        for (std::size_t i = 1; i < op.size(); ++i) {
48 ✔
291
            if (op[i] < '0' || op[i] > '9') return std::nullopt;
32 !
292
            denom *= 10.0;
32 ✔
293
        }
16 ✔
294
        return AggSpec(AggOp::Pct, field, op + "_" + field, "",
32 !
295
                       std::stod(op.substr(1)) / denom);
32 !
296
    }
NEW
297
    return std::nullopt;
×
298
}
874 ✔
299

300
bool parse_specs(PyObject* args, std::vector<AggSpec>& out) {
374 ✔
301
    const Py_ssize_t n = PyTuple_Size(args);
374 ✔
302
    for (Py_ssize_t i = 0; i < n; ++i) {
1,248 ✔
303
        const char* s = as_utf8(PyTuple_GetItem(args, i));
874 !
304
        if (!s) return false;
874 ✔
305
        std::optional<AggSpec> spec = parse_agg_spec(s);
1,311 !
306
        if (!spec) {
874 !
NEW
307
            PyErr_Format(PyExc_ValueError, "unknown agg spec: %s", s);
×
NEW
308
            return false;
×
309
        }
310
        out.push_back(std::move(*spec));
874 !
311
    }
874 !
312
    return true;
374 ✔
313
}
187 ✔
314

315
bool parse_containment(PyObject* args, PyObject* kwds, bool with_group,
48 ✔
316
                       views::ContainmentArgs& out) {
317
    static const char* kw_group[] = {"partition", "ts",    "dur",
318
                                     "name",      "group", nullptr};
319
    static const char* kw_plain[] = {"partition", "ts", "dur", "name", nullptr};
320
    PyObject* part = nullptr;
48 ✔
321
    PyObject* group = nullptr;
48 ✔
322
    const char* ts = "ts";
48 ✔
323
    const char* dur = "dur";
48 ✔
324
    const char* name = "name";
48 ✔
325
    if (!PyArg_ParseTupleAndKeywords(
48 !
326
            args, kwds, with_group ? "|OsssO" : "|Osss",
24 ✔
327
            const_cast<char**>(with_group ? kw_group : kw_plain), &part, &ts,
24 ✔
328
            &dur, &name, &group))
NEW
329
        return false;
×
330
    if (part && part != Py_None) {
48 !
331
        out.partition.clear();
48 ✔
332
        if (!py::parse_string_seq(part, "partition must be a sequence",
72 !
333
                                  out.partition))
48 !
NEW
334
            return false;
×
335
    }
24 ✔
336
    if (group && group != Py_None &&
69 !
337
        !py::parse_string_seq(group, "group must be a sequence", out.group))
42 !
NEW
338
        return false;
×
339
    out.ts = ts;
48 !
340
    out.dur = dur;
48 !
341
    out.name = name;
48 !
342
    return true;
48 ✔
343
}
24 ✔
344

345
// Runs the task `make` builds on the runtime `runtime_arg` names, with the GIL
346
// released.
347
template <class Make, class T>
348
bool run_on(PyObject* runtime_arg, Make&& make, T& out) {
88 ✔
349
    std::shared_ptr<dftracer::utils::Runtime> rt =
44 !
350
        runtime_from_arg(runtime_arg);
44 ✔
351
    if (!rt) return false;
88 !
352
    return run_blocking([&] { out = rt->submit(make()).get(); });
176 !
353
}
88 ✔
354

355
PyObject* tv_new(PyTypeObject* type, PyObject*, PyObject*) {
508 ✔
356
    auto* self = reinterpret_cast<TraceViewerObject*>(type->tp_alloc(type, 0));
508 ✔
357
    if (self) self->tv = nullptr;
508 ✔
358
    return reinterpret_cast<PyObject*>(self);
508 ✔
359
}
360

361
void tv_dealloc(TraceViewerObject* self) {
1,680 ✔
362
    delete self->tv;
1,680 ✔
363
    Py_TYPE(self)->tp_free(reinterpret_cast<PyObject*>(self));
1,680 ✔
364
}
1,680 ✔
365

366
int tv_init(TraceViewerObject* self, PyObject* args, PyObject* kwds) {
508 ✔
367
    static const char* kwlist[] = {"files", "index_path", nullptr};
368
    PyObject* files = nullptr;
508 ✔
369
    PyObject* index_obj = nullptr;
508 ✔
370
    if (!PyArg_ParseTupleAndKeywords(
508 !
371
            args, kwds, "O|O", const_cast<char**>(kwlist), &files, &index_obj))
254 ✔
372
        return -1;
×
373
    std::string index_dir;
508 ✔
374
    if (index_obj && index_obj != Py_None) {
508 !
375
        const char* s = as_utf8(index_obj);
172 !
376
        if (!s) return -1;
172 ✔
377
        index_dir = s;
172 !
378
    }
86 ✔
379
    std::vector<std::string> paths;
508 ✔
380
    std::optional<std::string> dir;
508 ✔
381
    if (PyUnicode_Check(files)) {
508 ✔
382
        const char* path = PyUnicode_AsUTF8(files);
392 !
383
        if (!path) return -1;
392 ✔
384
        std::error_code ec;
392 ✔
385
        if (fs::is_directory(path, ec))
392 !
386
            dir = path;
26 !
387
        else
388
            paths.emplace_back(path);
366 !
389
    } else if (!py::parse_string_seq(files,
312 !
390
                                     "files must be a path or a "
391
                                     "sequence of paths",
392
                                     paths)) {
NEW
393
        return -1;
×
394
    }
395
    View tv;
508 !
396
    if (dir) {
508 ✔
397
        if (!run_blocking([&] {
39 !
398
                tv = py::get_default_runtime()
39 ✔
399
                         ->submit(View::from_directory(*dir, index_dir))
39 !
400
                         .get();
26 !
401
            }))
26 ✔
NEW
402
            return -1;
×
403
    } else {
13 ✔
404
        std::vector<views::ViewFile> vfiles;
482 ✔
405
        vfiles.reserve(paths.size());
482 !
406
        for (const std::string& p : paths)
1,022 ✔
407
            vfiles.push_back(views::ViewFile{
540 !
408
                p, dftracer::utils::trace::internal::determine_index_path(
270 !
409
                       p, index_dir)});
270 ✔
410
        tv = View::from_files(std::move(vfiles));
482 !
411
    }
482 ✔
412
    delete self->tv;
508 ✔
413
    self->tv = new View(std::move(tv));
508 !
414
    return 0;
508 ✔
415
}
508 ✔
416

417
PyObject* tv_lazy(PyObject* self, PyObject*) {
1,678 ✔
418
    return guarded(
2,517 !
419
        [&] { return py::wrap_lazyframe(LazyFrame(tv_of(self).lazy())); });
4,195 !
420
}
421

422
PyObject* tv_with_lazy(PyObject* self, PyObject* arg) {
56 ✔
423
    const LazyFrame* lf = py::lazyframe_of(arg);
56 !
424
    if (!lf) return nullptr;
56 ✔
425
    return build(self, [&](const View& t) { return t.with_lazy(*lf); });
112 !
426
}
28 ✔
427

428
PyObject* tv_filter(PyObject* self, PyObject* arg) {
162 ✔
429
    PyObject* s = PyObject_Str(arg);
162 !
430
    if (!s) return nullptr;
162 ✔
431
    const char* dsl = PyUnicode_AsUTF8(s);
162 !
432
    if (!dsl) {
162 ✔
433
        Py_DECREF(s);
NEW
434
        return nullptr;
×
435
    }
436
    auto parsed = dftracer::utils::query::Query::from_string(dsl);
162 !
437
    Py_DECREF(s);
81 !
438
    if (!parsed) {
162 ✔
439
        PyErr_Format(PyExc_ValueError, "invalid filter query: %s",
2 !
440
                     parsed.error().message.c_str());
2 !
441
        return nullptr;
2 ✔
442
    }
443
    return build(self, [&](const View& t) {
320 !
444
        return t.filter(std::move(parsed.value()));
160 !
445
    });
80 ✔
446
}
162 ✔
447

448
PyObject* tv_select(PyObject* self, PyObject* arg) {
24 ✔
449
    std::vector<std::string> names;
24 ✔
450
    if (!py::parse_string_seq(arg, "select expects a sequence of names", names))
24 !
NEW
451
        return nullptr;
×
452
    return build(self,
48 !
453
                 [&](const View& t) { return t.select(std::move(names)); });
48 !
454
}
24 ✔
455

456
PyObject* tv_phase(PyObject* self, PyObject* arg) {
54 ✔
457
    const char* s = as_utf8(arg);
54 !
458
    if (!s) return nullptr;
54 ✔
459
    const std::string t(s);
54 !
460
    Phase ph;
461
    if (t == "events")
54 ✔
462
        ph = Phase::Events;
40 ✔
463
    else if (t == "counters")
14 ✔
464
        ph = Phase::Counters;
2 ✔
465
    else if (t == "aggregated")
12 ✔
466
        ph = Phase::Aggregated;
4 ✔
467
    else if (t == "metadata")
8 !
468
        ph = Phase::Metadata;
8 ✔
NEW
469
    else if (t == "any")
×
NEW
470
        ph = Phase::Any;
×
471
    else {
NEW
472
        PyErr_SetString(PyExc_ValueError,
×
473
                        "phase must be 'events', 'counters', 'aggregated', "
474
                        "'metadata', or 'any'");
NEW
475
        return nullptr;
×
476
    }
477
    return build(self, [&](const View& v) { return v.phase(ph); });
108 !
478
}
54 ✔
479

480
PyObject* tv_time_range(PyObject* self, PyObject* args) {
14 ✔
481
    double begin = 0, end = 0;
14 ✔
482
    if (!PyArg_ParseTuple(args, "dd", &begin, &end)) return nullptr;
14 !
483
    return build(self, [&](const View& t) { return t.time_range(begin, end); });
28 !
484
}
7 ✔
485

486
PyObject* tv_time_bucket(PyObject* self, PyObject* args, PyObject* kwds) {
32 ✔
487
    long long us = 0;
32 ✔
488
    PyObject* normalize_to = nullptr;
32 ✔
489
    static const char* kwlist[] = {"interval_us", "normalize_to", nullptr};
490
    if (!PyArg_ParseTupleAndKeywords(
32 !
491
            args, kwds, "L|O", const_cast<char**>(kwlist), &us, &normalize_to))
16 ✔
NEW
492
        return nullptr;
×
493
    if (us < 0) {
32 ✔
NEW
494
        PyErr_SetString(PyExc_ValueError, "interval_us must be >= 0");
×
NEW
495
        return nullptr;
×
496
    }
497
    const auto width = static_cast<std::uint64_t>(us);
32 ✔
498
    if (!normalize_to || normalize_to == Py_None)
32 !
499
        return build(self, [&](const View& t) { return t.time_bucket(width); });
52 !
500
    if (PyUnicode_Check(normalize_to)) {
6 ✔
501
        const char* s = PyUnicode_AsUTF8(normalize_to);
2 !
502
        if (!s) return nullptr;
2 ✔
503
        if (std::strcmp(s, "min") != 0) {
2 !
NEW
504
            PyErr_SetString(PyExc_ValueError,
×
505
                            "normalize_to must be an int origin or 'min'");
NEW
506
            return nullptr;
×
507
        }
508
        return build(self,
4 !
509
                     [&](const View& t) { return t.time_bucket_min(width); });
4 ✔
510
    }
511
    const long long origin = PyLong_AsLongLong(normalize_to);
4 !
512
    if (origin == -1 && PyErr_Occurred()) return nullptr;
4 !
513
    if (origin < 0) {
4 ✔
NEW
514
        PyErr_SetString(PyExc_ValueError, "normalize_to must be >= 0");
×
NEW
515
        return nullptr;
×
516
    }
517
    return build(self, [&](const View& t) {
8 !
518
        return t.time_bucket(width, static_cast<std::uint64_t>(origin));
4 ✔
519
    });
2 ✔
520
}
16 ✔
521

522
PyObject* tv_resolution(PyObject* self, PyObject* arg) {
4 ✔
523
    const long long us = PyLong_AsLongLong(arg);
4 !
524
    if (us == -1 && PyErr_Occurred()) return nullptr;
4 !
525
    if (us < 0) {
4 ✔
NEW
526
        PyErr_SetString(PyExc_ValueError, "resolution must be >= 0");
×
NEW
527
        return nullptr;
×
528
    }
529
    return build(self, [&](const View& t) {
8 !
530
        return t.resolution(static_cast<std::uint64_t>(us));
4 ✔
531
    });
2 ✔
532
}
2 ✔
533

NEW
534
PyObject* tv_time_scale(PyObject* self, PyObject* arg) {
×
NEW
535
    const double ratio = PyFloat_AsDouble(arg);
×
NEW
536
    if (ratio == -1.0 && PyErr_Occurred()) return nullptr;
×
NEW
537
    return build(self, [&](const View& t) { return t.time_scale(ratio); });
×
538
}
539

540
PyObject* tv_time_unit(PyObject* self, PyObject* arg) {
4 ✔
541
    namespace trace = dftracer::utils::trace;
542
    const char* s = as_utf8(arg);
4 !
543
    if (!s) return nullptr;
4 ✔
544
    const std::string t(s);
4 !
545
    trace::TimeMetric target;
546
    if (t == "ns")
4 !
NEW
547
        target = trace::TimeMetric::NS;
×
548
    else if (t == "us")
4 !
549
        target = trace::TimeMetric::US;
4 ✔
NEW
550
    else if (t == "ms")
×
NEW
551
        target = trace::TimeMetric::MS;
×
NEW
552
    else if (t == "sec" || t == "s")
×
NEW
553
        target = trace::TimeMetric::SEC;
×
554
    else {
NEW
555
        PyErr_SetString(PyExc_ValueError, "time_unit must be ns/us/ms/sec/s");
×
NEW
556
        return nullptr;
×
557
    }
558
    return build(self, [&](const View& v) {
8 !
559
        const double ratio =
2 ✔
560
            static_cast<double>(
2 ✔
561
                trace::time_metric_ns_per_unit(v.time_metric())) /
6 ✔
562
            static_cast<double>(trace::time_metric_ns_per_unit(target));
4 ✔
563
        return v.time_scale(ratio);
4 ✔
564
    });
2 ✔
565
}
4 ✔
566

567
PyObject* tv_group_by(PyObject* self, PyObject* args) {
370 ✔
568
    std::vector<GroupKey> keys;
370 ✔
569
    const Py_ssize_t n = PyTuple_Size(args);
370 !
570
    for (Py_ssize_t i = 0; i < n; ++i) {
814 ✔
571
        const char* s = as_utf8(PyTuple_GetItem(args, i));
444 !
572
        if (!s) return nullptr;
444 ✔
573
        keys.push_back(parse_group_key(s));
444 !
574
    }
222 ✔
575
    return build(self,
740 !
576
                 [&](const View& t) { return t.group_by(std::move(keys)); });
740 !
577
}
370 ✔
578

579
PyObject* tv_agg(PyObject* self, PyObject* args) {
362 ✔
580
    std::vector<AggSpec> specs;
362 ✔
581
    if (!parse_specs(args, specs)) return nullptr;
362 !
582
    return build(self, [&](const View& t) { return t.agg(std::move(specs)); });
724 !
583
}
362 ✔
584

585
PyObject* tv_agg_numeric_args(PyObject* self, PyObject* args) {
12 ✔
586
    std::vector<AggSpec> specs;
12 ✔
587
    if (!parse_specs(args, specs)) return nullptr;
12 !
588
    return build(self, [&](const View& t) {
24 !
589
        return specs.empty() ? t.agg_numeric_args()
15 !
590
                             : t.agg_numeric_args(std::move(specs));
9 !
591
    });
6 ✔
592
}
12 ✔
593

594
PyObject* tv_metadata(PyObject* self, PyObject* arg) {
54 ✔
595
    const int on = PyObject_IsTrue(arg);
54 !
596
    if (on < 0) return nullptr;
54 ✔
597
    return build(self, [&](const View& t) { return t.metadata(on != 0); });
108 !
598
}
27 ✔
599

600
PyObject* tv_rollup_root(PyObject* self, PyObject* arg) {
14 ✔
601
    const char* dir = as_utf8(arg);
14 !
602
    if (!dir) return nullptr;
14 ✔
603
    return build(self, [&](const View& t) { return t.rollup_root(dir); });
28 !
604
}
7 ✔
605

606
PyObject* tv_views_root(PyObject* self, PyObject* arg) {
14 ✔
607
    const char* dir = as_utf8(arg);
14 !
608
    if (!dir) return nullptr;
14 ✔
609
    return build(self, [&](const View& t) { return t.views_root(dir); });
28 !
610
}
7 ✔
611

NEW
612
PyObject* tv_memory_budget(PyObject* self, PyObject* arg) {
×
NEW
613
    const long long b = PyLong_AsLongLong(arg);
×
NEW
614
    if (b == -1 && PyErr_Occurred()) return nullptr;
×
NEW
615
    if (b < 0) {
×
NEW
616
        PyErr_SetString(PyExc_ValueError, "memory_budget must be >= 0");
×
NEW
617
        return nullptr;
×
618
    }
NEW
619
    return build(self, [&](const View& t) {
×
NEW
620
        return t.memory_budget(static_cast<std::uint64_t>(b));
×
621
    });
622
}
623

NEW
624
PyObject* tv_columns(PyObject* self, PyObject*) {
×
NEW
625
    std::vector<std::string> cols;
×
NEW
626
    if (!run_blocking([&] { cols = tv_of(self).columns(); })) return nullptr;
×
NEW
627
    return str_list_from(cols);
×
NEW
628
}
×
629

630
PyObject* tv_column_info(PyObject* self, PyObject*) {
6 ✔
631
    std::vector<views::ColumnInfo> info;
6 ✔
632
    if (!run_blocking([&] { info = tv_of(self).column_info(); }))
12 !
NEW
633
        return nullptr;
×
634
    PyObject* d = PyDict_New();
6 !
635
    if (!d) return nullptr;
6 !
636
    for (const auto& c : info) {
84 ✔
637
        PyObject* v = PyUnicode_FromString(c.type.c_str());
78 !
638
        if (!v || PyDict_SetItemString(d, c.name.c_str(), v) < 0) {
78 !
NEW
639
            Py_XDECREF(v);
×
640
            Py_DECREF(d);
×
NEW
641
            return nullptr;
×
642
        }
643
        Py_DECREF(v);
39 !
644
    }
645
    return d;
6 ✔
646
}
6 ✔
647

648
PyObject* tv_time_metric(PyObject* self, PyObject*) {
4 ✔
649
    std::string s(dftracer::utils::trace::time_metric_to_string(
6 !
650
        tv_of(self).time_metric()));
6 !
651
    for (char& c : s) c = static_cast<char>(std::tolower(c));
12 !
652
    return PyUnicode_FromString(s.c_str());
6 !
653
}
4 ✔
654

NEW
655
PyObject* tv_aggregates(PyObject* self, PyObject*) {
×
NEW
656
    return guarded([&] { return PyBool_FromLong(tv_of(self).aggregates()); });
×
657
}
658

659
// True while a filter still selects raw events: no trace aggregation and no
660
// op but filters on the plan. Reads no index.
661
PyObject* tv_filters_events(PyObject* self, PyObject*) {
28 ✔
662
    return PyBool_FromLong(tv_of(self).filters_events());
28 ✔
663
}
664

665
PyObject* tv_call_tree(PyObject* self, PyObject* args, PyObject* kwds) {
6 ✔
666
    views::ContainmentArgs a;
6 !
667
    if (!parse_containment(args, kwds, false, a)) return nullptr;
6 !
668
    return guarded([&] {
12 !
669
        return py::wrap_lazyframe(
6 !
670
            tv_of(self).call_tree(a.partition, a.ts, a.dur, a.name));
9 !
671
    });
3 ✔
672
}
6 ✔
673

674
PyObject* tv_flamegraph(PyObject* self, PyObject* args, PyObject* kwds) {
26 ✔
675
    views::ContainmentArgs a;
26 !
676
    if (!parse_containment(args, kwds, true, a)) return nullptr;
26 !
677
    return guarded([&] {
52 !
678
        return py::wrap_lazyframe(
24 !
679
            tv_of(self).flamegraph(a.partition, a.ts, a.dur, a.name, a.group));
45 !
680
    });
15 ✔
681
}
26 ✔
682

683
PyObject* tv_containment(PyObject* self, PyObject* args, PyObject* kwds) {
10 ✔
684
    views::ContainmentArgs a;
10 !
685
    if (!parse_containment(args, kwds, true, a)) return nullptr;
10 !
686
    return guarded([&]() -> PyObject* {
20 !
687
        auto r =
688
            tv_of(self).containment(a.partition, a.ts, a.dur, a.name, a.group);
15 !
689
        PyObject* ct = py::wrap_lazyframe(LazyFrame(r.plans()[0]));
10 !
690
        if (!ct) return nullptr;
10 ✔
691
        PyObject* fg = py::wrap_lazyframe(LazyFrame(r.plans()[1]));
10 !
692
        if (!fg) {
10 !
693
            Py_DECREF(ct);
×
NEW
694
            return nullptr;
×
695
        }
696
        return Py_BuildValue("(NN)", ct, fg);
10 !
697
    });
15 ✔
698
}
10 ✔
699

700
PyObject* tv_flamegraph_partial(PyObject* self, PyObject* args,
6 ✔
701
                                PyObject* kwds) {
702
    views::ContainmentArgs a;
6 !
703
    if (!parse_containment(args, kwds, true, a)) return nullptr;
6 !
704
    return guarded([&] {
12 !
705
        return py::wrap_lazyframe(LazyFrame(
9 !
706
            tv_of(self)
6 ✔
707
                .flamegraph_partial(a.partition, a.ts, a.dur, a.name, a.group)
9 !
708
                .plans()
6 !
709
                .front()));
12 ✔
710
    });
3 ✔
711
}
6 ✔
712

713
PyObject* tv_aggregate_partial(PyObject* self, PyObject*) {
16 ✔
714
    return guarded([&] {
31 !
715
        return py::wrap_lazyframe(
14 !
716
            LazyFrame(tv_of(self).aggregate_partial().plans().front()));
22 !
717
    });
16 ✔
718
}
719

720
PyObject* tv_sink_json(PyObject* self, PyObject* arg) {
32 ✔
721
    const char* path = as_utf8(arg);
32 !
722
    if (!path) return nullptr;
32 ✔
723
    FILE* f = std::fopen(path, "wb");
32 !
724
    if (!f) {
32 ✔
NEW
725
        PyErr_SetFromErrnoWithFilename(PyExc_OSError, path);
×
NEW
726
        return nullptr;
×
727
    }
728
    auto sink = std::make_shared<FileSink>(f);
32 !
729
    return guarded([&] {
48 !
730
        return py::wrap_lazyframe(LazyFrame(
48 !
731
            tv_of(self).sink_json(sink, views::LAZY).plans().front()));
64 !
732
    });
16 ✔
733
}
32 ✔
734

735
PyObject* tv_typed(PyObject* self, PyObject* args, PyObject* kwds) {
76 ✔
736
    int shard_begin = 0, shard_end = 0;
76 ✔
737
    PyObject* progress = nullptr;
76 ✔
738
    PyObject* runtime_arg = nullptr;
76 ✔
739
    static const char* kwlist[] = {"shard_begin", "shard_end", "progress",
740
                                   "runtime", nullptr};
741
    if (!PyArg_ParseTupleAndKeywords(args, kwds, "|iiOO",
76 !
742
                                     const_cast<char**>(kwlist), &shard_begin,
743
                                     &shard_end, &progress, &runtime_arg))
NEW
744
        return nullptr;
×
745
    views::TypedResult typed;
76 ✔
746
    views::ProgressFn fn = progress_from(progress);
76 !
747
    if (!run_on(
76 !
748
            runtime_arg,
38 ✔
749
            [&] {
114 ✔
750
                return tv_of(self).collect_typed(shard_begin, shard_end, fn);
76 !
751
            },
752
            typed))
NEW
753
        return nullptr;
×
754
    PyObject* d = PyDict_New();
76 !
755
    if (!d) return nullptr;
76 !
756
    const std::pair<const char*, DataFrame*> parts[] = {
76 ✔
757
        {"regular", &typed.regular},
114 ✔
758
        {"aggregated", &typed.aggregated},
114 ✔
759
        {"counters", &typed.counters}};
76 ✔
760
    for (const auto& [key, frame] : parts) {
304 ✔
761
        PyObject* v = py::wrap_dataframe(std::move(*frame));
228 !
762
        if (!v || PyDict_SetItemString(d, key, v) < 0) {
228 !
NEW
763
            Py_XDECREF(v);
×
764
            Py_DECREF(d);
×
NEW
765
            return nullptr;
×
766
        }
767
        Py_DECREF(v);
114 !
768
    }
769
    return d;
76 ✔
770
}
76 ✔
771

772
PyObject* tv_materialize(PyObject* self, PyObject* args, PyObject* kwds) {
6 ✔
773
    long long checkpoint_size = 0, part_size = 0;
6 ✔
774
    PyObject* progress = nullptr;
6 ✔
775
    PyObject* runtime_arg = nullptr;
6 ✔
776
    static const char* kwlist[] = {"checkpoint_size", "part_size", "progress",
777
                                   "runtime", nullptr};
778
    if (!PyArg_ParseTupleAndKeywords(
6 !
779
            args, kwds, "|LLOO", const_cast<char**>(kwlist), &checkpoint_size,
3 ✔
780
            &part_size, &progress, &runtime_arg))
NEW
781
        return nullptr;
×
782
    ExportStats stats;
6 ✔
783
    views::ProgressFn fn = progress_from(progress);
6 !
784
    if (!run_on(
6 !
785
            runtime_arg,
3 ✔
786
            [&] {
9 ✔
787
                return tv_of(self).materialize(
9 !
788
                    static_cast<std::uint64_t>(checkpoint_size),
6 ✔
789
                    static_cast<std::uint64_t>(part_size), fn);
6 !
790
            },
791
            stats))
NEW
792
        return nullptr;
×
793
    return stats_dict(stats);
6 !
794
}
6 ✔
795

796
// A trace of the aggregation (counter events) when the viewer aggregates, else
797
// of the selected events.
798
PyObject* tv_export_trace(PyObject* self, PyObject* args, PyObject* kwds) {
6 ✔
799
    static const char* kwlist[] = {"path",        "compress", "index",
800
                                   "member_size", "level",    "part_size",
801
                                   "runtime",     nullptr};
802
    const char* path = nullptr;
6 ✔
803
    int compress = 1;
6 ✔
804
    int index = 0;
6 ✔
805
    long long member_size = 0;
6 ✔
806
    int level = 6;
6 ✔
807
    long long part_size = 0;
6 ✔
808
    PyObject* runtime_arg = nullptr;
6 ✔
809
    if (!PyArg_ParseTupleAndKeywords(
6 !
810
            args, kwds, "s|ppLiLO", const_cast<char**>(kwlist), &path,
3 ✔
811
            &compress, &index, &member_size, &level, &part_size, &runtime_arg))
NEW
812
        return nullptr;
×
813
    const View& t = tv_of(self);
6 ✔
814
    ExportStats stats;
6 ✔
815
    bool aggregates = false;
6 ✔
816
    try {
817
        aggregates = t.aggregates();
6 !
818
    } catch (const std::exception& e) {
3 !
NEW
819
        py::set_typed_py_error(e);
×
NEW
820
        return nullptr;
×
NEW
821
    }
×
822
    if (aggregates) {
6 ✔
823
        FILE* f = std::fopen(path, "wb");
2 !
824
        if (!f) {
2 ✔
NEW
825
            PyErr_SetFromErrnoWithFilename(PyExc_OSError, path);
×
NEW
826
            return nullptr;
×
827
        }
828
        std::unique_ptr<views::ExportSink> sink;
2 ✔
829
        if (compress)
2 ✔
830
            sink = std::make_unique<GzipSink>(f, level);
2 !
831
        else
NEW
832
            sink = std::make_unique<FileSink>(f);
×
833
        const bool ok =
1 ✔
834
            run_on(runtime_arg, [&] { return t.sink_counters(*sink); }, stats);
4 !
835
        sink.reset();
2 ✔
836
        if (!ok) return nullptr;
2 ✔
837
        return stats_dict(stats);
2 !
838
    }
2 ✔
839
    views::TraceWriteOptions opts;
4 ✔
840
    opts.output_path = path;
4 !
841
    opts.member_size = static_cast<std::size_t>(member_size);
4 ✔
842
    opts.compress = compress != 0;
4 ✔
843
    opts.level = level;
4 ✔
844
    opts.build_index = index != 0;
4 ✔
845
    opts.part_size = static_cast<std::size_t>(part_size);
4 ✔
846
    if (!run_on(
4 !
847
            runtime_arg, [&] { return t.sink_trace(std::move(opts)); }, stats))
6 !
NEW
848
        return nullptr;
×
849
    return stats_dict(stats);
4 !
850
}
5 ✔
851

852
PyObject* tv_merge_partials(PyObject* self, PyObject* arg) {
8 ✔
853
    std::vector<std::string> owned;
8 ✔
854
    if (!py::parse_bytes_seq(arg, "merge_partials expects a sequence", owned))
8 !
NEW
855
        return nullptr;
×
856
    std::vector<std::string_view> parts(owned.begin(), owned.end());
8 !
857
    DataFrame out;
8 ✔
858
    if (!run_blocking([&] { out = tv_of(self).merge_partials(parts); }))
16 !
NEW
859
        return nullptr;
×
860
    return py::wrap_dataframe(std::move(out));
8 !
861
}
8 ✔
862

863
PyObject* tv_materialize_partials(PyObject* self, PyObject* args,
2 ✔
864
                                  PyObject* kwds) {
865
    PyObject* arg = nullptr;
2 ✔
866
    PyObject* runtime_arg = nullptr;
2 ✔
867
    static const char* kwlist[] = {"partials", "runtime", nullptr};
868
    if (!PyArg_ParseTupleAndKeywords(
2 !
869
            args, kwds, "O|O", const_cast<char**>(kwlist), &arg, &runtime_arg))
1 ✔
NEW
870
        return nullptr;
×
871
    std::vector<std::string> owned;
2 ✔
872
    if (!py::parse_bytes_seq(arg, "materialize_partials expects a sequence",
2 !
873
                             owned))
NEW
874
        return nullptr;
×
875
    std::vector<std::string_view> parts(owned.begin(), owned.end());
2 !
876
    std::shared_ptr<dftracer::utils::Runtime> rt =
877
        runtime_from_arg(runtime_arg);
2 !
878
    if (!rt) return nullptr;
2 ✔
879
    if (!run_blocking(
2 !
880
            [&] { rt->submit(tv_of(self).materialize_partials(parts)).get(); }))
3 !
881
        return nullptr;
2 ✔
NEW
882
    Py_RETURN_NONE;
×
883
}
2 ✔
884

885
PyObject* tv_reconstruct_if_cached(PyObject* self, PyObject*) {
8 ✔
886
    std::optional<DataFrame> out;
8 ✔
887
    if (!run_blocking([&] { out = tv_of(self).reconstruct_if_cached(); }))
16 !
888
        return nullptr;
2 ✔
889
    if (!out) Py_RETURN_NONE;
6 ✔
890
    return py::wrap_dataframe(std::move(*out));
2 !
891
}
8 ✔
892

893
PyObject* tv_mv_source(PyObject* self, PyObject*) {
6 ✔
894
    std::vector<std::string> out;
6 ✔
895
    if (!run_blocking([&] { out = tv_of(self).mv_source(); })) return nullptr;
12 !
896
    return str_list_from(out);
6 !
897
}
6 ✔
898

NEW
899
PyObject* tv_materialize_dir(PyObject* self, PyObject*) {
×
NEW
900
    std::string out;
×
NEW
901
    if (!run_blocking([&] { out = tv_of(self).materialize_dir(); }))
×
NEW
902
        return nullptr;
×
NEW
903
    return PyUnicode_FromStringAndSize(out.data(),
×
NEW
904
                                       static_cast<Py_ssize_t>(out.size()));
×
NEW
905
}
×
906

NEW
907
PyObject* tv_register_materialized(PyObject* self, PyObject* arg) {
×
NEW
908
    const char* dir = as_utf8(arg);
×
NEW
909
    if (!dir) return nullptr;
×
NEW
910
    const std::string d(dir);
×
NEW
911
    if (!run_blocking([&] { tv_of(self).register_materialized(d); }))
×
NEW
912
        return nullptr;
×
NEW
913
    Py_RETURN_NONE;
×
NEW
914
}
×
915

916
PyObject* tv_compare(PyObject* self, PyObject* arg) {
10 ✔
917
    if (!PyObject_TypeCheck(arg, &TraceViewerType)) {
10 ✔
NEW
918
        PyErr_SetString(PyExc_TypeError, "compare() needs a _TraceViewer");
×
NEW
919
        return nullptr;
×
920
    }
921
    return guarded(
15 !
922
        [&] { return py::wrap_lazyframe(tv_of(self).compare(tv_of(arg))); });
20 !
923
}
5 ✔
924

925
// A branch that folds `host`'s plugin set over the scan and leaves the named
926
// results in the host, where plugin_results() reads them after the collect.
927
PyObject* tv_plugins(PyObject* self, PyObject* host) {
14 ✔
928
    if (!PyObject_TypeCheck(host, &PluginHostType)) {
14 !
NEW
929
        PyErr_SetString(PyExc_TypeError, "plugins() needs a Plugins instance");
×
NEW
930
        return nullptr;
×
931
    }
932
    const dftracer::utils::plugins::Plugins* set =
7 ✔
933
        py::plugin_host_plugins(host);
14 !
934
    if (!set) return nullptr;
14 ✔
935
    dftracer::utils::plugins::NamedResultRegistry* results =
7 ✔
936
        py::plugin_host_results(host);
14 !
937
    if (!results) return nullptr;
14 ✔
938
    Py_INCREF(host);
7 ✔
939
    std::shared_ptr<PyObject> keep(host, [](PyObject* p) {
21 ✔
940
        PyGILState_STATE g = PyGILState_Ensure();
14 ✔
941
        Py_DECREF(p);
7 ✔
942
        PyGILState_Release(g);
14 ✔
943
    });
14 !
944
    views::SessionBranch attach = [set, results, keep](views::ViewSession& s) {
28 !
945
        views::Deferred<dftracer::utils::plugins::PluginRun> h = set->attach(s);
14 !
946
        return std::function<void(const ExportStats&)>(
7 !
947
            [h, results, keep](const ExportStats&) mutable {
35 ✔
948
                *results = std::move(h.get().results);
14 ✔
949
            });
21 !
950
    };
21 !
951
    return guarded([&] {
21 !
952
        return py::wrap_lazyframe(tv_of(self).branch(std::move(attach)));
14 !
953
    });
7 ✔
954
}
14 ✔
955

956
PyObject* merge_flamegraph_partials_py(PyObject*, PyObject* arg) {
4 ✔
957
    std::vector<std::string> owned;
4 ✔
958
    if (!py::parse_bytes_seq(
4 !
959
            arg, "merge_flamegraph_partials expects a sequence", owned))
2 ✔
NEW
960
        return nullptr;
×
961
    std::vector<std::string_view> parts(owned.begin(), owned.end());
4 !
962
    return guarded([&] {
6 !
963
        return py::wrap_dataframe(View::merge_flamegraph_partials(parts));
4 !
964
    });
2 ✔
965
}
4 ✔
966

967
PyObject* plugin_results_py(PyObject*, PyObject* host) {
14 ✔
968
    if (!PyObject_TypeCheck(host, &PluginHostType)) {
14 ✔
NEW
969
        PyErr_SetString(PyExc_TypeError,
×
970
                        "plugin_results() needs a Plugins instance");
NEW
971
        return nullptr;
×
972
    }
973
    return py::plugin_host_results_dict(host);
14 ✔
974
}
7 ✔
975

976
PyMethodDef tv_methods[] = {
977
    {"lazy", tv_lazy, METH_NOARGS, "The plan as a _LazyFrame."},
978
    {"with_lazy", tv_with_lazy, METH_O,
979
     "This scan with `plan` (a _LazyFrame over it) as its plan."},
980
    {"filter", tv_filter, METH_O,
981
     "Keep events matching a query-DSL predicate."},
982
    {"select", tv_select, METH_O,
983
     "Fields the scan reads (raw events), or a projection of the plan."},
984
    {"phase", tv_phase, METH_O,
985
     "Select 'events', 'counters', 'aggregated', 'metadata', or 'any'."},
986
    {"time_range", tv_time_range, METH_VARARGS,
987
     "Restrict to a [begin, end) timestamp window."},
988
    {"time_bucket", DFTU_PYCFUNCTION(tv_time_bucket),
989
     METH_VARARGS | METH_KEYWORDS,
990
     "time_bucket(interval_us, normalize_to=None); normalize_to is an int "
991
     "origin or 'min'."},
992
    {"resolution", tv_resolution, METH_O,
993
     "Grid in microseconds the occupancy aggregates snap to; 0 is exact."},
994
    {"time_scale", tv_time_scale, METH_O, "Multiply ts/dur by this ratio."},
995
    {"time_unit", tv_time_unit, METH_O,
996
     "Normalize ts/dur to ns/us/ms/sec from the trace's own unit."},
997
    {"group_by", tv_group_by, METH_VARARGS, "Trace group keys."},
998
    {"agg", tv_agg, METH_VARARGS, "Trace aggregate specs."},
999
    {"agg_numeric_args", tv_agg_numeric_args, METH_VARARGS,
1000
     "Aggregate every discovered numeric arg."},
1001
    {"metadata", tv_metadata, METH_O, "Include metadata records."},
1002
    {"rollup_root", tv_rollup_root, METH_O, "Rollup root directory."},
1003
    {"views_root", tv_views_root, METH_O, "Materialized-view root directory."},
1004
    {"memory_budget", tv_memory_budget, METH_O, "Spill budget in bytes."},
1005
    {"columns", tv_columns, METH_NOARGS,
1006
     "Columns discoverable from the index (no scan)."},
1007
    {"column_info", tv_column_info, METH_NOARGS,
1008
     "Index columns mapped to their type name (no scan)."},
1009
    {"time_metric", tv_time_metric, METH_NOARGS, "The trace's time unit."},
1010
    {"aggregates", tv_aggregates, METH_NOARGS,
1011
     "True when the plan absorbs into an aggregating scan."},
1012
    {"filters_events", tv_filters_events, METH_NOARGS,
1013
     "True while a filter still selects raw events (reads no index)."},
1014
    {"call_tree", DFTU_PYCFUNCTION(tv_call_tree), METH_VARARGS | METH_KEYWORDS,
1015
     "Plan of the events plus level/parent_id."},
1016
    {"flamegraph", DFTU_PYCFUNCTION(tv_flamegraph),
1017
     METH_VARARGS | METH_KEYWORDS, "Plan of the folded node frame."},
1018
    {"containment", DFTU_PYCFUNCTION(tv_containment),
1019
     METH_VARARGS | METH_KEYWORDS,
1020
     "(call_tree plan, flamegraph plan) sharing one buffered fold."},
1021
    {"flamegraph_partial", DFTU_PYCFUNCTION(tv_flamegraph_partial),
1022
     METH_VARARGS | METH_KEYWORDS, "Plan of a one-row 'partial' frame."},
1023
    {"aggregate_partial", tv_aggregate_partial, METH_NOARGS,
1024
     "Plan of a one-row 'partial' frame."},
1025
    {"sink_json", tv_sink_json, METH_O,
1026
     "Plan writing the selected events to `path`; yields one stats row."},
1027
    {"typed", DFTU_PYCFUNCTION(tv_typed), METH_VARARGS | METH_KEYWORDS,
1028
     "The aggregation index's record families (runs now)."},
1029
    {"materialize", DFTU_PYCFUNCTION(tv_materialize),
1030
     METH_VARARGS | METH_KEYWORDS, "Persist this query (runs now)."},
1031
    {"export_trace", DFTU_PYCFUNCTION(tv_export_trace),
1032
     METH_VARARGS | METH_KEYWORDS, "Write a trace file (runs now)."},
1033
    {"merge_partials", tv_merge_partials, METH_O,
1034
     "Merge aggregate partials with this aggregation."},
1035
    {"materialize_partials", DFTU_PYCFUNCTION(tv_materialize_partials),
1036
     METH_VARARGS | METH_KEYWORDS, "Write the rollup from partials."},
1037
    {"reconstruct_if_cached", tv_reconstruct_if_cached, METH_NOARGS,
1038
     "The rollup as a DataFrame, or None on a miss."},
1039
    {"mv_source", tv_mv_source, METH_NOARGS,
1040
     "Materialized-view files that would serve this query."},
1041
    {"materialize_dir", tv_materialize_dir, METH_NOARGS,
1042
     "Create and return the shared materialized-view directory."},
1043
    {"register_materialized", tv_register_materialized, METH_O,
1044
     "Write the materialized-view manifest at `dir`."},
1045
    {"compare", tv_compare, METH_O,
1046
     "Plan comparing this aggregation against another viewer's events."},
1047
    {"plugins", tv_plugins, METH_O,
1048
     "Plan folding a Plugins set over the scan."},
1049
    {nullptr, nullptr, 0, nullptr}};
1050

1051
PyMethodDef module_methods[] = {
1052
    {"merge_flamegraph_partials", merge_flamegraph_partials_py, METH_O,
1053
     "merge_flamegraph_partials(partials) -> node DataFrame (no scan)."},
1054
    {"plugin_results", plugin_results_py, METH_O,
1055
     "plugin_results(plugins) -> {name: result} of its last run."},
1056
    {nullptr, nullptr, 0, nullptr}};
1057

1058
}  // namespace
1059

1060
PyTypeObject TraceViewerType = [] {
1 ✔
1061
    PyTypeObject t{PyVarObject_HEAD_INIT(nullptr, 0)};
1 ✔
1062
    t.tp_name = "dftracer_utils_ext._TraceViewer";
1 ✔
1063
    t.tp_basicsize = sizeof(TraceViewerObject);
1 ✔
1064
    t.tp_dealloc = reinterpret_cast<destructor>(tv_dealloc);
1 ✔
1065
    t.tp_flags = Py_TPFLAGS_DEFAULT;
1 ✔
1066
    t.tp_doc = "A trace scan and the plan over it.";
1 ✔
1067
    t.tp_methods = tv_methods;
1 ✔
1068
    t.tp_init = reinterpret_cast<initproc>(tv_init);
1 ✔
1069
    t.tp_new = tv_new;
1 ✔
1070
    return t;
1 ✔
1071
}();
1072

1073
int dftracer::utils::python::init_trace_viewer(PyObject* m) {
2 ✔
1074
    if (register_type(m, &TraceViewerType, "_TraceViewer") < 0) return -1;
2 ✔
1075
    return PyModule_AddFunctions(m, module_methods);
2 ✔
1076
}
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