Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -114,20 +114,21 @@ class ThreadInfo
};

private:
using TaskAddressCallback = std::function<void(TaskObj*)>;

void reset_cycle_state() noexcept;
void render_unwound_stacks(EchionSampler&);
[[nodiscard]] Result<void> unwind_tasks(EchionSampler&, PyThreadState*, microsecond_t wall_time_us);
void unwind_greenlets(EchionSampler&, PyThreadState*, unsigned long, microsecond_t wall_time_us);
[[nodiscard]] Result<std::vector<TaskInfo::Ptr>> get_all_tasks(EchionSampler&, PyThreadState* tstate);
[[nodiscard]] Result<void> for_each_task_address(EchionSampler&,
PyThreadState* tstate,
const TaskAddressCallback& callback);
#if PY_VERSION_HEX >= 0x030e0000
[[nodiscard]] Result<void> get_tasks_from_thread_linked_list(EchionSampler& echion,
std::vector<TaskInfo::Ptr>& tasks);
[[nodiscard]] Result<void> get_tasks_from_interpreter_linked_list(EchionSampler& echion,
PyThreadState* tstate,
std::vector<TaskInfo::Ptr>& tasks);
[[nodiscard]] Result<void> get_tasks_from_linked_list(EchionSampler& echion,
uintptr_t head_addr,
std::vector<TaskInfo::Ptr>& tasks);
[[nodiscard]] Result<void> get_tasks_from_thread_linked_list(const TaskAddressCallback& callback);
[[nodiscard]] Result<void> get_tasks_from_interpreter_linked_list(PyThreadState* tstate,
const TaskAddressCallback& callback);
[[nodiscard]] Result<void> get_tasks_from_linked_list(uintptr_t head_addr, const TaskAddressCallback& callback);
#endif
};

Expand Down
88 changes: 37 additions & 51 deletions ddtrace/internal/datadog/profiling/stack/src/echion/threads.cc
Original file line number Diff line number Diff line change
Expand Up @@ -411,7 +411,7 @@ ThreadInfo::unwind_tasks(EchionSampler& echion, PyThreadState* tstate, microseco
// ----------------------------------------------------------------------------
#if PY_VERSION_HEX >= 0x030e0000
Result<void>
ThreadInfo::get_tasks_from_thread_linked_list(EchionSampler& echion, std::vector<TaskInfo::Ptr>& tasks)
ThreadInfo::get_tasks_from_thread_linked_list(const TaskAddressCallback& callback)
{
if (this->tstate_addr == 0 || this->asyncio_loop == 0) {
return ErrorKind::TaskInfoError;
Expand All @@ -426,13 +426,11 @@ ThreadInfo::get_tasks_from_thread_linked_list(EchionSampler& echion, std::vector
constexpr size_t asyncio_tasks_head_offset = offsetof(_PyThreadStateImpl, asyncio_tasks_head);
uintptr_t head_addr = this->tstate_addr + asyncio_tasks_head_offset;

return get_tasks_from_linked_list(echion, head_addr, tasks);
return get_tasks_from_linked_list(head_addr, callback);
}

Result<void>
ThreadInfo::get_tasks_from_interpreter_linked_list(EchionSampler& echion,
PyThreadState* tstate,
std::vector<TaskInfo::Ptr>& tasks)
ThreadInfo::get_tasks_from_interpreter_linked_list(PyThreadState* tstate, const TaskAddressCallback& callback)
{
if (tstate == nullptr || tstate->interp == nullptr || this->asyncio_loop == 0) {
return ErrorKind::TaskInfoError;
Expand All @@ -441,11 +439,11 @@ ThreadInfo::get_tasks_from_interpreter_linked_list(EchionSampler& echion,
constexpr size_t asyncio_tasks_head_offset = offsetof(PyInterpreterState, asyncio_tasks_head);
uintptr_t head_addr = reinterpret_cast<uintptr_t>(tstate->interp) + asyncio_tasks_head_offset;

return get_tasks_from_linked_list(echion, head_addr, tasks);
return get_tasks_from_linked_list(head_addr, callback);
}

Result<void>
ThreadInfo::get_tasks_from_linked_list(EchionSampler& echion, uintptr_t head_addr, std::vector<TaskInfo::Ptr>& tasks)
ThreadInfo::get_tasks_from_linked_list(uintptr_t head_addr, const TaskAddressCallback& callback)
{
if (head_addr == 0 || this->asyncio_loop == 0) {
return ErrorKind::TaskInfoError;
Expand Down Expand Up @@ -489,14 +487,7 @@ ThreadInfo::get_tasks_from_linked_list(EchionSampler& echion, uintptr_t head_add
size_t task_node_offset_val = offsetof(TaskObj, task_node);
uintptr_t task_addr_uint = next_node_addr - task_node_offset_val;

// Create TaskInfo for the task
auto maybe_task_info = TaskInfo::create(echion, reinterpret_cast<TaskObj*>(task_addr_uint));
if (maybe_task_info) {
auto& task_info = *maybe_task_info;
if (task_info->loop == reinterpret_cast<PyObject*>(this->asyncio_loop)) {
tasks.push_back(std::move(task_info));
}
}
callback(reinterpret_cast<TaskObj*>(task_addr_uint));

// Read next node from current_node.next into current_node
if (copy_type(reinterpret_cast<void*>(next_node_addr), current_node)) {
Expand All @@ -507,12 +498,11 @@ ThreadInfo::get_tasks_from_linked_list(EchionSampler& echion, uintptr_t head_add
return Result<void>::ok();
}

Result<std::vector<TaskInfo::Ptr>>
ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState* tstate)
Result<void>
ThreadInfo::for_each_task_address(EchionSampler& echion, PyThreadState* tstate, const TaskAddressCallback& callback)
{
std::vector<TaskInfo::Ptr> tasks;
if (this->asyncio_loop == 0)
return tasks;
return Result<void>::ok();

// Python 3.14+: Native tasks are in linked-list per thread AND per interpreter
// CPython iterates over both:
Expand All @@ -521,10 +511,10 @@ ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState* tstate)
// First, get tasks from this thread's linked-list (if tstate_addr is set)
// Note: We continue processing even if one source fails to maximize partial results
if (tstate != nullptr && this->tstate_addr != 0) {
(void)get_tasks_from_thread_linked_list(echion, tasks);
(void)get_tasks_from_thread_linked_list(callback);

// Second, get tasks from interpreter's linked-list (lingering tasks)
(void)get_tasks_from_interpreter_linked_list(echion, tstate, tasks);
(void)get_tasks_from_interpreter_linked_list(tstate, callback);
}

// Handle third-party tasks from Python _scheduled_tasks WeakSet
Expand All @@ -540,11 +530,7 @@ ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState* tstate)
auto scheduled_tasks = std::move(*maybe_scheduled_tasks);
for (auto task_addr : scheduled_tasks) {
// In WeakSet.data (set), elements are the Task objects themselves
auto maybe_task_info = TaskInfo::create(echion, reinterpret_cast<TaskObj*>(task_addr));
if (maybe_task_info &&
(*maybe_task_info)->loop == reinterpret_cast<PyObject*>(this->asyncio_loop)) {
tasks.push_back(std::move(*maybe_task_info));
}
callback(reinterpret_cast<TaskObj*>(task_addr));
}
}
}
Expand All @@ -566,25 +552,19 @@ ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState* tstate)

auto eager_tasks = std::move(*maybe_eager_tasks);
for (auto task_addr : eager_tasks) {
auto maybe_task_info = TaskInfo::create(echion, reinterpret_cast<TaskObj*>(task_addr));
if (maybe_task_info) {
if ((*maybe_task_info)->loop == reinterpret_cast<PyObject*>(this->asyncio_loop)) {
tasks.push_back(std::move(*maybe_task_info));
}
}
callback(reinterpret_cast<TaskObj*>(task_addr));
}
}

return tasks;
return Result<void>::ok();
}
#else
// Pre-Python 3.14: get_all_tasks uses WeakSet approach
Result<std::vector<TaskInfo::Ptr>>
ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState*)
// Pre-Python 3.14: asyncio tracks tasks through Python weak sets.
Result<void>
ThreadInfo::for_each_task_address(EchionSampler& echion, PyThreadState*, const TaskAddressCallback& callback)
{
std::vector<TaskInfo::Ptr> tasks;
if (this->asyncio_loop == 0)
return tasks;
return Result<void>::ok();

auto asyncio_scheduled_tasks = echion.asyncio_scheduled_tasks();
auto maybe_scheduled_tasks_set = MirrorSet::create(asyncio_scheduled_tasks);
Expand All @@ -604,12 +584,7 @@ ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState*)
if (copy_type(task_wr_addr, task_wr))
continue;

auto maybe_task_info = TaskInfo::create(echion, reinterpret_cast<TaskObj*>(task_wr.wr_object));
if (maybe_task_info) {
if (reinterpret_cast<uintptr_t>((*maybe_task_info)->loop) == this->asyncio_loop) {
tasks.push_back(std::move(*maybe_task_info));
}
}
callback(reinterpret_cast<TaskObj*>(task_wr.wr_object));
}

auto asyncio_eager_tasks = echion.asyncio_eager_tasks();
Expand All @@ -628,19 +603,30 @@ ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState*)

auto eager_tasks = std::move(*maybe_eager_tasks);
for (auto task_addr : eager_tasks) {
auto maybe_task_info = TaskInfo::create(echion, reinterpret_cast<TaskObj*>(task_addr));
if (maybe_task_info) {
if (reinterpret_cast<uintptr_t>((*maybe_task_info)->loop) == this->asyncio_loop) {
tasks.push_back(std::move(*maybe_task_info));
}
}
callback(reinterpret_cast<TaskObj*>(task_addr));
}
}

return tasks;
return Result<void>::ok();
}
#endif // PY_VERSION_HEX >= 0x030e0000

Result<std::vector<TaskInfo::Ptr>>
ThreadInfo::get_all_tasks(EchionSampler& echion, PyThreadState* tstate)
{
std::vector<TaskInfo::Ptr> tasks;
auto result = for_each_task_address(echion, tstate, [&](TaskObj* task_addr) {
auto maybe_task_info = TaskInfo::create(echion, task_addr);
if (maybe_task_info && reinterpret_cast<uintptr_t>((*maybe_task_info)->loop) == this->asyncio_loop) {
tasks.push_back(std::move(*maybe_task_info));
}
});
if (!result) {
return result.error();
}
return tasks;
}

// ----------------------------------------------------------------------------
void
ThreadInfo::unwind_greenlets(EchionSampler& echion,
Expand Down
Loading