Skip to content

Commit c81693b

Browse files
committed
refactor(profiling): visit task identities during discovery
1 parent 9155126 commit c81693b

2 files changed

Lines changed: 43 additions & 46 deletions

File tree

ddtrace/internal/datadog/profiling/stack/echion/echion/threads.h

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -121,17 +121,21 @@ class ThreadInfo
121121
};
122122

123123
private:
124+
using TaskAddressCallback = std::function<void(TaskObj*)>;
125+
124126
void reset_cycle_state() noexcept;
125127
void render_unwound_stacks(EchionSampler&);
126128
[[nodiscard]] Result<void> unwind_tasks(EchionSampler&, PyThreadState*, microsecond_t wall_time_us);
127129
void unwind_greenlets(EchionSampler&, PyThreadState*, unsigned long, microsecond_t wall_time_us);
128130
[[nodiscard]] Result<std::vector<TaskInfo::Ptr>> get_all_tasks(EchionSampler&, PyThreadState* tstate);
129-
[[nodiscard]] Result<std::vector<TaskObj*>> get_all_task_addresses(EchionSampler&, PyThreadState* tstate);
131+
[[nodiscard]] Result<void> for_each_task_address(EchionSampler&,
132+
PyThreadState* tstate,
133+
const TaskAddressCallback& callback);
130134
#if PY_VERSION_HEX >= 0x030e0000
131-
[[nodiscard]] Result<void> get_task_addresses_from_thread_linked_list(std::vector<TaskObj*>& tasks);
132-
[[nodiscard]] Result<void> get_task_addresses_from_interpreter_linked_list(PyThreadState* tstate,
133-
std::vector<TaskObj*>& tasks);
134-
[[nodiscard]] Result<void> get_task_addresses_from_linked_list(uintptr_t head_addr, std::vector<TaskObj*>& tasks);
135+
[[nodiscard]] Result<void> visit_thread_task_addresses(const TaskAddressCallback& callback);
136+
[[nodiscard]] Result<void> visit_interpreter_task_addresses(PyThreadState* tstate,
137+
const TaskAddressCallback& callback);
138+
[[nodiscard]] Result<void> visit_linked_task_addresses(uintptr_t head_addr, const TaskAddressCallback& callback);
135139
#endif
136140
};
137141

ddtrace/internal/datadog/profiling/stack/src/echion/threads.cc

Lines changed: 34 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -413,30 +413,30 @@ ThreadInfo::unwind_tasks(EchionSampler& echion, PyThreadState* tstate, microseco
413413
// ----------------------------------------------------------------------------
414414
#if PY_VERSION_HEX >= 0x030e0000
415415
Result<void>
416-
ThreadInfo::get_task_addresses_from_thread_linked_list(std::vector<TaskObj*>& tasks)
416+
ThreadInfo::visit_thread_task_addresses(const TaskAddressCallback& callback)
417417
{
418418
if (this->tstate_addr == 0 || this->asyncio_loop == 0) {
419419
return ErrorKind::TaskInfoError;
420420
}
421421

422422
constexpr size_t asyncio_tasks_head_offset = offsetof(_PyThreadStateImpl, asyncio_tasks_head);
423-
return get_task_addresses_from_linked_list(this->tstate_addr + asyncio_tasks_head_offset, tasks);
423+
return visit_linked_task_addresses(this->tstate_addr + asyncio_tasks_head_offset, callback);
424424
}
425425

426426
Result<void>
427-
ThreadInfo::get_task_addresses_from_interpreter_linked_list(PyThreadState* tstate, std::vector<TaskObj*>& tasks)
427+
ThreadInfo::visit_interpreter_task_addresses(PyThreadState* tstate, const TaskAddressCallback& callback)
428428
{
429429
if (tstate == nullptr || tstate->interp == nullptr || this->asyncio_loop == 0) {
430430
return ErrorKind::TaskInfoError;
431431
}
432432

433433
constexpr size_t asyncio_tasks_head_offset = offsetof(PyInterpreterState, asyncio_tasks_head);
434434
const uintptr_t head_addr = reinterpret_cast<uintptr_t>(tstate->interp) + asyncio_tasks_head_offset;
435-
return get_task_addresses_from_linked_list(head_addr, tasks);
435+
return visit_linked_task_addresses(head_addr, callback);
436436
}
437437

438438
Result<void>
439-
ThreadInfo::get_task_addresses_from_linked_list(uintptr_t head_addr, std::vector<TaskObj*>& tasks)
439+
ThreadInfo::visit_linked_task_addresses(uintptr_t head_addr, const TaskAddressCallback& callback)
440440
{
441441
if (head_addr == 0 || this->asyncio_loop == 0) {
442442
return ErrorKind::TaskInfoError;
@@ -462,7 +462,7 @@ ThreadInfo::get_task_addresses_from_linked_list(uintptr_t head_addr, std::vector
462462

463463
const uintptr_t next_node_addr = reinterpret_cast<uintptr_t>(current_node.next);
464464
const uintptr_t task_addr = next_node_addr - offsetof(TaskObj, task_node);
465-
tasks.push_back(reinterpret_cast<TaskObj*>(task_addr));
465+
callback(reinterpret_cast<TaskObj*>(task_addr));
466466

467467
if (copy_type(reinterpret_cast<void*>(next_node_addr), current_node)) {
468468
return ErrorKind::TaskInfoError;
@@ -472,28 +472,25 @@ ThreadInfo::get_task_addresses_from_linked_list(uintptr_t head_addr, std::vector
472472
return Result<void>::ok();
473473
}
474474

475-
Result<std::vector<TaskObj*>>
476-
ThreadInfo::get_all_task_addresses(EchionSampler& echion, PyThreadState* tstate)
475+
Result<void>
476+
ThreadInfo::for_each_task_address(EchionSampler& echion, PyThreadState* tstate, const TaskAddressCallback& callback)
477477
{
478-
std::vector<TaskObj*> tasks;
479478
if (this->asyncio_loop == 0) {
480-
return tasks;
479+
return Result<void>::ok();
481480
}
482481

483482
// Native tasks can appear in both lists. Consumers already tolerate duplicate task objects.
484-
if (this->tstate_addr != 0) {
485-
(void)get_task_addresses_from_thread_linked_list(tasks);
486-
}
487-
if (tstate != nullptr) {
488-
(void)get_task_addresses_from_interpreter_linked_list(tstate, tasks);
483+
if (tstate != nullptr && this->tstate_addr != 0) {
484+
(void)visit_thread_task_addresses(callback);
485+
(void)visit_interpreter_task_addresses(tstate, callback);
489486
}
490487

491488
// Python 3.14 stores only third-party Task implementations in _scheduled_tasks.
492489
if (auto scheduled = echion.asyncio_scheduled_tasks(); scheduled != nullptr) {
493490
if (auto maybe_set = MirrorSet::create(scheduled)) {
494491
if (auto maybe_tasks = maybe_set->as_unordered_set()) {
495492
for (auto task : *maybe_tasks) {
496-
tasks.push_back(reinterpret_cast<TaskObj*>(task));
493+
callback(reinterpret_cast<TaskObj*>(task));
497494
}
498495
}
499496
}
@@ -509,19 +506,18 @@ ThreadInfo::get_all_task_addresses(EchionSampler& echion, PyThreadState* tstate)
509506
return ErrorKind::TaskInfoError;
510507
}
511508
for (auto task : *maybe_tasks) {
512-
tasks.push_back(reinterpret_cast<TaskObj*>(task));
509+
callback(reinterpret_cast<TaskObj*>(task));
513510
}
514511
}
515512

516-
return tasks;
513+
return Result<void>::ok();
517514
}
518515
#else
519-
Result<std::vector<TaskObj*>>
520-
ThreadInfo::get_all_task_addresses(EchionSampler& echion, PyThreadState*)
516+
Result<void>
517+
ThreadInfo::for_each_task_address(EchionSampler& echion, PyThreadState*, const TaskAddressCallback& callback)
521518
{
522-
std::vector<TaskObj*> tasks;
523519
if (this->asyncio_loop == 0) {
524-
return tasks;
520+
return Result<void>::ok();
525521
}
526522

527523
auto maybe_set = MirrorSet::create(echion.asyncio_scheduled_tasks());
@@ -535,7 +531,7 @@ ThreadInfo::get_all_task_addresses(EchionSampler& echion, PyThreadState*)
535531
for (auto weakref_addr : *maybe_scheduled) {
536532
PyWeakReference weakref;
537533
if (!copy_type(weakref_addr, weakref)) {
538-
tasks.push_back(reinterpret_cast<TaskObj*>(weakref.wr_object));
534+
callback(reinterpret_cast<TaskObj*>(weakref.wr_object));
539535
}
540536
}
541537

@@ -549,28 +545,26 @@ ThreadInfo::get_all_task_addresses(EchionSampler& echion, PyThreadState*)
549545
return ErrorKind::TaskInfoError;
550546
}
551547
for (auto task : *maybe_eager) {
552-
tasks.push_back(reinterpret_cast<TaskObj*>(task));
548+
callback(reinterpret_cast<TaskObj*>(task));
553549
}
554550
}
555551

556-
return tasks;
552+
return Result<void>::ok();
557553
}
558554
#endif // PY_VERSION_HEX >= 0x030e0000
559555

560556
Result<std::vector<TaskInfo::Ptr>>
561557
ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState* tstate)
562558
{
563-
auto maybe_addresses = get_all_task_addresses(echion, tstate);
564-
if (!maybe_addresses) {
565-
return maybe_addresses.error();
566-
}
567-
568559
std::vector<TaskInfo::Ptr> tasks;
569-
for (TaskObj* task_addr : *maybe_addresses) {
560+
auto result = for_each_task_address(echion, tstate, [&](TaskObj* task_addr) {
570561
auto maybe_task = TaskInfo::create(echion, task_addr);
571562
if (maybe_task && reinterpret_cast<uintptr_t>((*maybe_task)->loop) == this->asyncio_loop) {
572563
tasks.push_back(std::move(*maybe_task));
573564
}
565+
});
566+
if (!result) {
567+
return result.error();
574568
}
575569
return tasks;
576570
}
@@ -1169,17 +1163,16 @@ ThreadInfo::sample_cpu_timer(EchionSampler& echion,
11691163
return;
11701164
}
11711165

1172-
auto maybe_addresses = get_all_task_addresses(echion, tstate);
1173-
if (maybe_addresses) {
1174-
std::vector<TaskIdentity> tasks;
1175-
tasks.reserve(maybe_addresses->size());
1176-
for (TaskObj* address : *maybe_addresses) {
1177-
TaskObj task;
1178-
if (!copy_type(address, task) && reinterpret_cast<uintptr_t>(task.task_loop) == asyncio_loop) {
1179-
tasks.push_back({ address, task.task_coro, task.task_fut_waiter });
1180-
}
1166+
// Snapshot each lightweight identity as its address is visited. Retaining bare addresses for a later pass
1167+
// would widen the race with task completion and address reuse.
1168+
std::vector<TaskIdentity> tasks;
1169+
auto visit_result = for_each_task_address(echion, tstate, [&](TaskObj* address) {
1170+
TaskObj task;
1171+
if (!copy_type(address, task) && reinterpret_cast<uintptr_t>(task.task_loop) == asyncio_loop) {
1172+
tasks.push_back({ address, task.task_coro, task.task_fut_waiter });
11811173
}
1182-
1174+
});
1175+
if (visit_result) {
11831176
if (TaskObj* address = find_captured_task(tasks, raw)) {
11841177
auto maybe_task = TaskInfo::create(echion, address);
11851178
if (maybe_task) {

0 commit comments

Comments
 (0)