From 0fe03b649414c6280506deba484074005d7b012b Mon Sep 17 00:00:00 2001 From: German Andryeyev Date: Fri, 3 Feb 2023 10:33:55 -0500 Subject: [PATCH] SWDEV-353281 - Switch Graph to operate with stream MemPool was designed to use hip::Stream, but graph implementation uses amd::HostQueue. Hence switch graph to hip::Stream management. Change-Id: Ia319389de45e4c3c6043d17473279a6f27a13140 --- hipamd/src/hip_graph_internal.cpp | 43 ++++++++++++++----------------- hipamd/src/hip_graph_internal.hpp | 43 ++++++++++++++++--------------- 2 files changed, 42 insertions(+), 44 deletions(-) diff --git a/hipamd/src/hip_graph_internal.cpp b/hipamd/src/hip_graph_internal.cpp index ab7fc65d92..48a7bc3091 100644 --- a/hipamd/src/hip_graph_internal.cpp +++ b/hipamd/src/hip_graph_internal.cpp @@ -679,34 +679,29 @@ bool hipGraphExec::isGraphExecValid(hipGraphExec* pGraphExec) { return true; } -hipError_t hipGraphExec::CreateQueues(size_t numQueues) { - parallelQueues_.reserve(numQueues); - for (size_t i = 0; i < numQueues; i++) { - amd::HostQueue* queue; - queue = new amd::HostQueue( - *hip::getCurrentDevice()->asContext(), *hip::getCurrentDevice()->devices()[0], 0, - amd::CommandQueue::RealTimeDisabled, amd::CommandQueue::Priority::Normal); - - bool result = (queue != nullptr) ? queue->create() : false; - // Create a host queue - if (result) { - parallelQueues_.push_back(queue); - } else { - ClPrint(amd::LOG_ERROR, amd::LOG_CODE, "[hipGraph] Failed to create host queue\n"); +hipError_t hipGraphExec::CreateStreams(uint32_t num_streams) { + parallel_streams_.reserve(num_streams); + for (uint32_t i = 0; i < num_streams; ++i) { + auto stream = new hip::Stream(hip::getCurrentDevice(), + hip::Stream::Priority::Normal, hipStreamNonBlocking); + if (stream == nullptr || !stream->Create()) { + delete stream; + ClPrint(amd::LOG_ERROR, amd::LOG_CODE, "[hipGraph] Failed to create parallel stream!\n"); return hipErrorOutOfMemory; } + parallel_streams_.push_back(stream); } return hipSuccess; } hipError_t hipGraphExec::Init() { hipError_t status; - size_t reqNumQueues = 1; + size_t min_num_streams = 1; for (auto& node : levelOrder_) { - reqNumQueues += node->GetNumParallelQueues(); + min_num_streams += node->GetNumParallelStreams(); } - status = CreateQueues(parallelLists_.size() - 1 + reqNumQueues); + status = CreateStreams(parallelLists_.size() - 1 + min_num_streams); return status; } @@ -771,19 +766,19 @@ hipError_t FillCommands(std::vector>& parallelLists, return hipSuccess; } -void UpdateQueue(std::vector>& parallelLists, amd::HostQueue*& queue, +void UpdateStream(std::vector>& parallelLists, hip::Stream* stream, hipGraphExec* ptr) { int i = 0; for (const auto& list : parallelLists) { // first parallel list will be launched on the same queue as parent if (i == 0) { for (auto& node : list) { - node->SetQueue(queue, ptr); + node->SetStream(stream, ptr); } - } else { // New queue for parallel branches - amd::HostQueue* paralleQueue = ptr->GetAvailableQueue(); + } else { // New stream for parallel branches + hip::Stream* stream = ptr->GetAvailableStreams(); for (auto& node : list) { - node->SetQueue(paralleQueue, ptr); + node->SetStream(stream, ptr); } } i++; @@ -801,7 +796,9 @@ hipError_t hipGraphExec::Run(hipStream_t stream) { levelOrder_[0]->GetParentGraph()->FreeAllMemory(); } } - UpdateQueue(parallelLists_, queue, this); + auto hip_stream = (stream == nullptr) ? hip::getCurrentDevice()->GetNullStream() + : reinterpret_cast(stream); + UpdateStream(parallelLists_, hip_stream, this); std::vector rootCommands; amd::Command* endCommand = nullptr; status = diff --git a/hipamd/src/hip_graph_internal.hpp b/hipamd/src/hip_graph_internal.hpp index 536e0387fd..e2bee8a10f 100644 --- a/hipamd/src/hip_graph_internal.hpp +++ b/hipamd/src/hip_graph_internal.hpp @@ -39,8 +39,8 @@ hipError_t FillCommands(std::vector>& parallelLists, std::unordered_map>& nodeWaitLists, std::vector& levelOrder, std::vector& rootCommands, amd::Command*& endCommand, amd::HostQueue* queue); -void UpdateQueue(std::vector>& parallelLists, amd::HostQueue*& queue, - hipGraphExec* ptr); +void UpdateStream(std::vector>& parallelLists, hip::Stream* stream, + hipGraphExec* ptr); struct hipUserObject : public amd::ReferenceCountedObject { typedef void (*UserCallbackDestructor)(void* data); @@ -154,6 +154,7 @@ struct hipGraphNodeDOTAttribute { struct hipGraphNode : public hipGraphNodeDOTAttribute { protected: + hip::Stream* stream_ = nullptr; amd::HostQueue* queue_; uint32_t level_; unsigned int id_; @@ -223,7 +224,10 @@ struct hipGraphNode : public hipGraphNodeDOTAttribute { amd::HostQueue* GetQueue() { return queue_; } - virtual void SetQueue(amd::HostQueue* queue, hipGraphExec* ptr = nullptr) { queue_ = queue; } + virtual void SetStream(hip::Stream* stream, hipGraphExec* ptr = nullptr) { + stream_ = stream; + queue_ = stream->asHostQueue(); + } /// Create amd::command for the graph node virtual hipError_t CreateCommand(amd::HostQueue* queue) { commands_.clear(); @@ -337,7 +341,7 @@ struct hipGraphNode : public hipGraphNodeDOTAttribute { command->updateEventWaitList(waitList); } } - virtual size_t GetNumParallelQueues() { return 0; } + virtual size_t GetNumParallelStreams() { return 0; } /// Enqueue commands part of the node virtual void EnqueueCommands(hipStream_t stream) { // If the node is disabled it becomes empty node. To maintain ordering just enqueue marker. @@ -541,7 +545,7 @@ struct hipGraphExec { // level order of the graph doesn't include nodes embedded as part of the child graph std::vector levelOrder_; std::unordered_map> nodeWaitLists_; - std::vector parallelQueues_; + std::vector parallel_streams_; uint currentQueueIndex_; std::unordered_map clonedNodes_; amd::Command* lastEnqueuedCommand_; @@ -570,8 +574,8 @@ struct hipGraphExec { ~hipGraphExec() { // new commands are launched for every launch they are destroyed as and when command is // terminated after it complete execution - for (auto queue : parallelQueues_) { - queue->release(); + for (auto stream : parallel_streams_) { + delete stream; } for (auto it = clonedNodes_.begin(); it != clonedNodes_.end(); it++) delete it->second; amd::ScopedLock lock(graphExecSetLock_); @@ -596,10 +600,10 @@ struct hipGraphExec { std::vector& GetNodes() { return levelOrder_; } - amd::HostQueue* GetAvailableQueue() { return parallelQueues_[currentQueueIndex_++]; } + hip::Stream* GetAvailableStreams() { return parallel_streams_[currentQueueIndex_++]; } void ResetQueueIndex() { currentQueueIndex_ = 0; } hipError_t Init(); - hipError_t CreateQueues(size_t numQueues); + hipError_t CreateStreams(uint32_t num_streams); hipError_t Run(hipStream_t stream); }; @@ -628,20 +632,21 @@ struct hipChildGraphNode : public hipGraphNode { ihipGraph* GetChildGraph() { return childGraph_; } - size_t GetNumParallelQueues() { + size_t GetNumParallelStreams() { LevelOrder(childGraphlevelOrder_); size_t num = 0; for (auto& node : childGraphlevelOrder_) { - num += node->GetNumParallelQueues(); + num += node->GetNumParallelStreams(); } // returns total number of parallel queues required for child graph nodes to be launched // first parallel list will be launched on the same queue as parent return num + (parallelLists_.size() - 1); } - void SetQueue(amd::HostQueue* queue, hipGraphExec* ptr = nullptr) { - queue_ = queue; - UpdateQueue(parallelLists_, queue, ptr); + void SetStream(hip::Stream* stream, hipGraphExec* ptr = nullptr) { + stream_ = stream; + queue_ = stream->asHostQueue(); + UpdateStream(parallelLists_, stream, ptr); } // For nodes that are dependent on the child graph node waitlist is the last node of the first @@ -1876,9 +1881,7 @@ class hipGraphMemAllocNode : public hipGraphNode { virtual hipError_t CreateCommand(amd::HostQueue* queue) { auto error = hipGraphNode::CreateCommand(queue); - // Note: memory pool can work with hip::Streams only. It can't accept amd::HostQueue. - // Resource tracking is disabled! - auto ptr = Execute(); + auto ptr = Execute(stream_); return error; } @@ -1919,13 +1922,11 @@ class hipGraphMemFreeNode : public hipGraphNode { virtual hipError_t CreateCommand(amd::HostQueue* queue) { auto error = hipGraphNode::CreateCommand(queue); - // Note: memory pool can work with hip::Streams only. It can't accept amd::HostQueue. - // Resource tracking is disabled! - Execute(); + Execute(stream_); return error; } - void Execute(hip::Stream* stream = nullptr) { + void Execute(hip::Stream* stream) { auto graph = GetParentGraph(); if (graph != nullptr) { graph->FreeMemory(device_ptr_, stream);