summaryrefslogtreecommitdiff
path: root/src/KingSystem/Utils/Thread
diff options
context:
space:
mode:
authorLéo Lam <leo@leolam.fr>2020-09-15 19:13:52 +0200
committerLéo Lam <leo@leolam.fr>2020-09-16 17:49:37 +0200
commit3db6228dfc65a25da49bab26ab8bcf4abc69a98f (patch)
tree3946a9211c2adc7213b1ae94ccf4baca22b73020 /src/KingSystem/Utils/Thread
parent8b7369dffb1e816b2f27ab0133d2bac218383a45 (diff)
ksys: Add more task utilities
* TaskMgr, ManagedTask, ManagedTaskHandle * GameTaskThread: partial implementation because PhysicsMemSys / Havok stuff hasn't been decompiled yet and calc_() requires PhysicsMemSys
Diffstat (limited to 'src/KingSystem/Utils/Thread')
-rw-r--r--src/KingSystem/Utils/Thread/GameTaskThread.cpp15
-rw-r--r--src/KingSystem/Utils/Thread/GameTaskThread.h27
-rw-r--r--src/KingSystem/Utils/Thread/ManagedTask.cpp106
-rw-r--r--src/KingSystem/Utils/Thread/ManagedTask.h47
-rw-r--r--src/KingSystem/Utils/Thread/ManagedTaskHandle.cpp107
-rw-r--r--src/KingSystem/Utils/Thread/ManagedTaskHandle.h60
-rw-r--r--src/KingSystem/Utils/Thread/TaskMgr.cpp212
-rw-r--r--src/KingSystem/Utils/Thread/TaskMgr.h88
8 files changed, 662 insertions, 0 deletions
diff --git a/src/KingSystem/Utils/Thread/GameTaskThread.cpp b/src/KingSystem/Utils/Thread/GameTaskThread.cpp
new file mode 100644
index 00000000..eacd9e64
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/GameTaskThread.cpp
@@ -0,0 +1,15 @@
+#include "KingSystem/Utils/Thread/GameTaskThread.h"
+
+namespace ksys::util {
+
+GameTaskThread::GameTaskThread(const sead::SafeString& name, sead::Heap* heap, s32 priority,
+ sead::MessageQueue::BlockType block_type, long quit_msg,
+ s32 stack_size, s32 message_queue_size)
+ : TaskThread(name, heap, priority, block_type, quit_msg, stack_size, message_queue_size) {}
+
+void GameTaskThread::quit(bool) {
+ mMessageQueue.push(cMessage_GameThreadQuit, sead::MessageQueue::BlockType::Blocking);
+ Thread::quit(false);
+}
+
+} // namespace ksys::util
diff --git a/src/KingSystem/Utils/Thread/GameTaskThread.h b/src/KingSystem/Utils/Thread/GameTaskThread.h
new file mode 100644
index 00000000..a16e9f68
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/GameTaskThread.h
@@ -0,0 +1,27 @@
+#pragma once
+
+#include "KingSystem/Utils/Thread/TaskThread.h"
+#include "KingSystem/Utils/Types.h"
+
+namespace ksys::util {
+
+class GameTaskThread : public TaskThread {
+ SEAD_RTTI_OVERRIDE(GameTaskThread, TaskThread)
+public:
+ GameTaskThread(const sead::SafeString& name, sead::Heap* heap, s32 priority,
+ sead::MessageQueue::BlockType block_type, long quit_msg, s32 stack_size,
+ s32 message_queue_size);
+ ~GameTaskThread() override { ; }
+
+ void quit(bool is_jam) override;
+
+protected:
+ static constexpr s32 cMessage_GameThreadQuit = 4;
+
+ void calc_(sead::MessageQueue::Element msg) override;
+ u8 _1a0 = 0;
+ s32 _1a4 = -1;
+};
+KSYS_CHECK_SIZE_NX150(GameTaskThread, 0x1a8);
+
+} // namespace ksys::util
diff --git a/src/KingSystem/Utils/Thread/ManagedTask.cpp b/src/KingSystem/Utils/Thread/ManagedTask.cpp
new file mode 100644
index 00000000..8a0f3dfd
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/ManagedTask.cpp
@@ -0,0 +1,106 @@
+#include "KingSystem/Utils/Thread/ManagedTask.h"
+#include "KingSystem/Utils/Thread/ManagedTaskHandle.h"
+#include "KingSystem/Utils/Thread/TaskMgr.h"
+#include "KingSystem/Utils/Thread/TaskQueueLock.h"
+
+namespace ksys::util {
+
+ManagedTask::ManagedTask(sead::Heap* heap) : Task(heap) {}
+
+ManagedTask::ManagedTask(sead::IDisposer::HeapNullOption heap_null_option)
+ : Task(nullptr, heap_null_option) {}
+
+ManagedTask::~ManagedTask() {
+ if (mHandle) {
+ mHandle->finalize();
+ mHandle = nullptr;
+ }
+ finalize_();
+}
+
+void ManagedTask::prepare_(TaskRequest* request) {
+ mIsIdle = false;
+ prepareImpl_(request);
+}
+
+void ManagedTask::run_() {
+ Task::run_();
+ onRun_();
+}
+
+void ManagedTask::onRunFinished_() {}
+
+void ManagedTask::onFinish_() {
+ if (auto* handle = mHandle) {
+ handle->setIsSuccess(isSuccess());
+ handle->setStatus(ManagedTaskHandle::Status::TaskFinished);
+ }
+
+ if (!mHandle && mMgr)
+ mMgr->freeTask(this);
+}
+
+void ManagedTask::onPostFinish_() {
+ mIsIdle = true;
+}
+
+void ManagedTask::preRemove_() {
+ preRemoveImpl_();
+
+ if (mHandle)
+ mHandle->setStatus(ManagedTaskHandle::Status::TaskRemoved);
+
+ if (!mHandle && mMgr)
+ mMgr->freeTask(this);
+}
+
+void ManagedTask::postRemove_() {
+ mIsIdle = true;
+}
+
+void ManagedTask::onRun_() {}
+
+void ManagedTask::prepareImpl_(TaskRequest*) {}
+
+void ManagedTask::preRemoveImpl_() {}
+
+bool ManagedTask::isIdle() const {
+ return mIsIdle;
+}
+
+void ManagedTask::setMgr(TaskMgr* mgr) {
+ mMgr = mgr;
+}
+
+void ManagedTask::attachHandle(ManagedTaskHandle* handle, TaskQueueBase* queue) {
+ if (mHandle)
+ return;
+
+ if (handle) {
+ handle->attachTask({this, queue});
+ handle->setStatus(ManagedTaskHandle::Status::TaskAttached);
+ }
+ mHandle = handle;
+}
+
+// NON_MATCHING: switch
+void ManagedTask::detachHandle() {
+ TaskQueueLock lock;
+ lock.lock(mQueue);
+
+ if (mHandle) {
+ switch (mHandle->getStatus()) {
+ case ManagedTaskHandle::Status::TaskRemoved:
+ case ManagedTaskHandle::Status::TaskFinished:
+ mHandle = nullptr;
+ if (mMgr)
+ mMgr->freeTask(this);
+ break;
+ default:
+ mHandle = nullptr;
+ break;
+ }
+ }
+}
+
+} // namespace ksys::util
diff --git a/src/KingSystem/Utils/Thread/ManagedTask.h b/src/KingSystem/Utils/Thread/ManagedTask.h
new file mode 100644
index 00000000..67b8ef04
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/ManagedTask.h
@@ -0,0 +1,47 @@
+#pragma once
+
+#include "KingSystem/Utils/Thread/Task.h"
+#include "KingSystem/Utils/Types.h"
+
+namespace ksys::util {
+
+class ManagedTaskHandle;
+class TaskMgr;
+
+class ManagedTask : public Task {
+ SEAD_RTTI_OVERRIDE(ManagedTask, Task)
+public:
+ explicit ManagedTask(sead::Heap* heap);
+ explicit ManagedTask(sead::IDisposer::HeapNullOption heap_null_option);
+ ~ManagedTask() override;
+
+ bool isIdle() const;
+
+protected:
+ friend class ManagedTaskHandle;
+ friend class TaskMgr;
+
+ void setMgr(TaskMgr* mgr);
+ void attachHandle(ManagedTaskHandle* handle, TaskQueueBase* queue);
+ void detachHandle();
+
+ void prepare_(TaskRequest* request) override;
+ void run_() override;
+ void onRunFinished_() override;
+ void onFinish_() override;
+ void onPostFinish_() override;
+ void preRemove_() override;
+ void postRemove_() override;
+
+ virtual void onRun_();
+ virtual void prepareImpl_(TaskRequest* req);
+ virtual void preRemoveImpl_();
+
+ bool mIsIdle = true;
+ TaskMgr* mMgr = nullptr;
+ // FIXME: rename
+ ManagedTaskHandle* mHandle = nullptr;
+};
+KSYS_CHECK_SIZE_NX150(ManagedTask, 0xc0);
+
+} // namespace ksys::util
diff --git a/src/KingSystem/Utils/Thread/ManagedTaskHandle.cpp b/src/KingSystem/Utils/Thread/ManagedTaskHandle.cpp
new file mode 100644
index 00000000..874bf46f
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/ManagedTaskHandle.cpp
@@ -0,0 +1,107 @@
+#include "KingSystem/Utils/Thread/ManagedTaskHandle.h"
+#include "KingSystem/Utils/Thread/ManagedTask.h"
+#include "KingSystem/Utils/Thread/TaskQueueBase.h"
+#include "KingSystem/Utils/Thread/TaskQueueLock.h"
+
+namespace ksys::util {
+
+ManagedTaskHandle::ManagedTaskHandle() = default;
+
+ManagedTaskHandle::~ManagedTaskHandle() {
+ removeTaskFromQueue();
+ finalize();
+}
+
+void ManagedTaskHandle::removeTaskFromQueue() {
+ if (!mQueue)
+ return;
+
+ incrementRef_();
+
+ if (mTask)
+ mTask->removeFromQueue();
+
+ decrementRef_();
+}
+
+void ManagedTaskHandle::finalize() {
+ if (mQueue)
+ decrementRef_();
+}
+
+bool ManagedTaskHandle::hasTask() const {
+ return mTask != nullptr;
+}
+
+bool ManagedTaskHandle::wait() {
+ if (!mQueue)
+ return true;
+
+ incrementRef_();
+
+ if (mTask)
+ mTask->wait();
+
+ decrementRef_();
+ return true;
+}
+
+bool ManagedTaskHandle::wait(const sead::TickSpan& wait_duration) {
+ if (!mQueue)
+ return true;
+
+ incrementRef_();
+
+ if (!mTask) {
+ decrementRef_();
+ return true;
+ }
+
+ const bool ret = mTask->wait(wait_duration);
+ decrementRef_();
+ return ret;
+}
+
+bool ManagedTaskHandle::isTaskAttached() const {
+ return mTask && mStatus == Status::TaskAttached;
+}
+
+bool ManagedTaskHandle::didTaskCompleteSuccessfully() const {
+ return mStatus == Status::TaskFinished && mSuccess;
+}
+
+bool ManagedTaskHandle::attachTask(const SetTaskArg& arg) {
+ mTask = arg.task;
+ mQueue = arg.queue;
+ incrementRef_();
+ return true;
+}
+
+void ManagedTaskHandle::setStatus(Status status) {
+ mStatus = status;
+}
+
+void ManagedTaskHandle::setIsSuccess(bool success) {
+ mSuccess = success;
+}
+
+inline void ManagedTaskHandle::incrementRef_() {
+ TaskQueueLock lock;
+ mQueue->lock(&lock);
+ ++mRefCount;
+}
+
+inline void ManagedTaskHandle::decrementRef_() {
+ TaskQueueLock lock;
+ mQueue->lock(&lock);
+
+ if (mRefCount > 0)
+ --mRefCount;
+
+ if (mRefCount <= 0 && mTask) {
+ mTask->detachHandle();
+ mTask = nullptr;
+ }
+}
+
+} // namespace ksys::util
diff --git a/src/KingSystem/Utils/Thread/ManagedTaskHandle.h b/src/KingSystem/Utils/Thread/ManagedTaskHandle.h
new file mode 100644
index 00000000..61c41c51
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/ManagedTaskHandle.h
@@ -0,0 +1,60 @@
+#pragma once
+
+#include <basis/seadTypes.h>
+#include <time/seadTickSpan.h>
+#include "KingSystem/Utils/Types.h"
+
+namespace ksys::util {
+
+class ManagedTask;
+class TaskQueueBase;
+
+class ManagedTaskHandle {
+public:
+ enum class Status {
+ Uninitialized = 0,
+ TaskAttached = 1,
+ TaskFinished = 2,
+ TaskRemoved = 3,
+ };
+
+ ManagedTaskHandle();
+ virtual ~ManagedTaskHandle();
+
+ void finalize();
+
+ Status getStatus() const { return mStatus; }
+
+ void removeTaskFromQueue();
+ bool hasTask() const;
+
+ bool wait();
+ bool wait(const sead::TickSpan& wait_duration);
+
+ bool isTaskAttached() const;
+ bool didTaskCompleteSuccessfully() const;
+
+private:
+ friend class ManagedTask;
+
+ struct SetTaskArg {
+ ManagedTask* task;
+ TaskQueueBase* queue;
+ };
+
+ bool attachTask(const SetTaskArg& arg);
+ void setIsSuccess(bool success);
+ void setStatus(Status status);
+
+ void incrementRef_();
+ void decrementRef_();
+
+ bool mSuccess = false;
+ Status mStatus = Status::Uninitialized;
+ s32 mRefCount = 0;
+ ManagedTask* mTask = nullptr;
+ TaskQueueBase* mQueue = nullptr;
+};
+KSYS_CHECK_SIZE_NX150(ManagedTaskHandle, 0x28);
+
+} // namespace ksys::util
diff --git a/src/KingSystem/Utils/Thread/TaskMgr.cpp b/src/KingSystem/Utils/Thread/TaskMgr.cpp
new file mode 100644
index 00000000..b5eec563
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/TaskMgr.cpp
@@ -0,0 +1,212 @@
+#include "KingSystem/Utils/Thread/TaskMgr.h"
+#include <heap/seadHeap.h>
+#include <heap/seadHeapMgr.h>
+#include "KingSystem/Utils/Thread/ManagedTask.h"
+#include "KingSystem/Utils/Thread/TaskThread.h"
+
+namespace ksys::util {
+
+TaskMgr::TaskMgr(sead::Heap* heap)
+ : mTasksCS(heap), mCS2(heap), mNewFreeTaskEvent(heap, true), mEvent2(heap, true) {
+ mFreeTaskLists[0].initOffset(ManagedTask::getListNodeOffset());
+ mFreeTaskLists[1].initOffset(ManagedTask::getListNodeOffset());
+ mNewFreeTaskEvent.resetSignal();
+}
+
+TaskMgr::~TaskMgr() {
+ finalize();
+}
+
+void TaskMgr::finalize() {
+ if (mFlags.isOn(Flag::HeapIsFreeable)) {
+ if (mTask)
+ delete mTask;
+ mTask = nullptr;
+
+ for (auto*& task : mTasks) {
+ if (task)
+ delete task;
+ task = nullptr;
+ }
+
+ mTasks.freeBuffer();
+
+ } else {
+ if (mTask) {
+ mTask->~ManagedTask();
+ mTask = nullptr;
+ }
+
+ for (auto* task : mTasks)
+ task->~ManagedTask();
+ }
+}
+
+void TaskMgr::submitRequest(TaskMgrRequest& request) {
+ bool request_had_no_task = false;
+ if (!request.task) {
+ request_had_no_task = true;
+ fetchIdleTaskForRequest_(request, true);
+ }
+
+ if (!request.task)
+ return;
+
+ auto* task = request.task;
+ if (request.handle) {
+ auto* queue = request.request->mQueue;
+ if (!queue) {
+ auto* thread = request.request->mThread;
+ queue = thread ? thread->getTaskQueue() : nullptr;
+ }
+ task->attachHandle(request.handle, queue);
+ } else {
+ task->attachHandle(nullptr, nullptr);
+ }
+
+ const bool ok = task->submitRequest(*request.request);
+ if (!ok) {
+ if (auto* managed_task = sead::DynamicCast<ManagedTask>(request.task)) {
+ managed_task->attachHandle(nullptr, nullptr);
+ if (!request_had_no_task)
+ return;
+ freeTask(managed_task);
+ } else if (!request_had_no_task) {
+ return;
+ }
+ }
+
+ if (request_had_no_task || !ok)
+ request.task = nullptr;
+}
+
+// NON_MATCHING: reorderings
+bool TaskMgr::fetchIdleTaskForRequest_(TaskMgrRequest& request, bool retry_until_success) {
+ if (!hasTasks())
+ return false;
+
+ ManagedTask* task = [this, retry_until_success] {
+ const auto lock1 = sead::makeScopedLock(mTasksCS);
+ if (auto* task = fetchIdleTask_(retry_until_success))
+ return task;
+
+ swapLists_();
+ if (!retry_until_success)
+ return fetchIdleTask_(retry_until_success);
+
+ while (true) {
+ if (auto* task_1 = fetchIdleTask_(retry_until_success))
+ return task_1;
+ mNewFreeTaskEvent.wait();
+ swapLists_();
+ }
+ }();
+
+ if (!task)
+ return false;
+
+ request.task = task;
+ return true;
+}
+
+void TaskMgr::freeTask(ManagedTask* task) {
+ auto lock = sead::makeScopedLock(mCS2);
+ mFreeTaskLists[getListIndex2_()].pushBack(task);
+ mNewFreeTaskEvent.setSignal();
+}
+
+bool TaskMgr::trySubmitRequest(TaskMgrRequest& request) {
+ bool request_had_no_task = false;
+ if (!request.task) {
+ request_had_no_task = true;
+ if (!tryFetchTaskForRequest_(request, false))
+ return false;
+ }
+
+ if (!request.task)
+ return false;
+
+ auto* task = request.task;
+ if (request.handle) {
+ auto* queue = request.request->mQueue;
+ if (!queue) {
+ auto* thread = request.request->mThread;
+ queue = thread ? thread->getTaskQueue() : nullptr;
+ }
+ task->attachHandle(request.handle, queue);
+ } else {
+ task->attachHandle(nullptr, nullptr);
+ }
+
+ const bool ok = task->submitRequest(*request.request);
+ if (!ok) {
+ if (auto* managed_task = sead::DynamicCast<ManagedTask>(request.task)) {
+ managed_task->attachHandle(nullptr, nullptr);
+ if (!request_had_no_task)
+ return ok;
+ freeTask(managed_task);
+ } else if (!request_had_no_task) {
+ return ok;
+ }
+ }
+
+ if (request_had_no_task || !ok)
+ request.task = nullptr;
+ return ok;
+}
+
+// NON_MATCHING: the factory invoke function pointer is loaded earlier in the original code
+void TaskMgr::init(s32 num_tasks, sead::Heap* heap, ManagedTaskFactory& factory) {
+ if (!heap->isFreeable())
+ mFlags.reset(Flag::HeapIsFreeable);
+ else
+ mFlags.set(Flag::HeapIsFreeable);
+
+ const sead::ScopedCurrentHeapSetter heap_setter{heap};
+
+ mTasks.allocBufferAssert(num_tasks, heap);
+
+ if (mTasks.size() != 0) {
+ auto& list = mFreeTaskLists[mListIndex];
+ for (auto*& task : mTasks) {
+ factory(&task);
+ task->setMgr(this);
+ list.pushBack(task);
+ }
+ }
+
+ factory(&mTask);
+}
+
+bool TaskMgr::hasTasks() const {
+ return mTasks.size() > 0;
+}
+
+ManagedTask* TaskMgr::fetchIdleTask_(bool retry_until_success) {
+ auto lock = sead::makeScopedLock(mTasksCS);
+
+ if (mFreeTaskLists[mListIndex].isEmpty())
+ return nullptr;
+
+ ManagedTask* result = nullptr;
+ while (true) {
+ for (auto& task : mFreeTaskLists[mListIndex]) {
+ if (task.isIdle()) {
+ result = std::addressof(task);
+ break;
+ }
+ }
+
+ if (result) {
+ auto lock1 = sead::makeScopedLock(mTasksCS);
+ mFreeTaskLists[mListIndex].erase(result);
+ return result;
+ }
+
+ if (!retry_until_success)
+ return nullptr;
+ mEvent2.wait(sead::TickSpan::fromMilliSeconds(1));
+ }
+}
+
+} // namespace ksys::util
diff --git a/src/KingSystem/Utils/Thread/TaskMgr.h b/src/KingSystem/Utils/Thread/TaskMgr.h
new file mode 100644
index 00000000..76a91ad9
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/TaskMgr.h
@@ -0,0 +1,88 @@
+#pragma once
+
+#include <container/seadBuffer.h>
+#include <container/seadOffsetList.h>
+#include <container/seadSafeArray.h>
+#include <heap/seadDisposer.h>
+#include <prim/seadDelegate.h>
+#include <prim/seadRuntimeTypeInfo.h>
+#include <prim/seadScopedLock.h>
+#include <prim/seadTypedBitFlag.h>
+#include <thread/seadCriticalSection.h>
+#include "KingSystem/Utils/Thread/Event.h"
+#include "KingSystem/Utils/Types.h"
+
+namespace ksys::util {
+
+class ManagedTask;
+class ManagedTaskHandle;
+class TaskRequest;
+struct TaskMgrRequest;
+
+using ManagedTaskFactory = sead::IDelegate1<ManagedTask**>;
+
+struct TaskMgrRequest {
+ /// Optional. If null, a task from the internal buffer will be used.
+ ManagedTask* task;
+ /// Must not be null.
+ TaskRequest* request;
+ /// Optional.
+ ManagedTaskHandle* handle;
+};
+KSYS_CHECK_SIZE_NX150(TaskMgrRequest, 0x18);
+
+class TaskMgr {
+ SEAD_RTTI_BASE(TaskMgr)
+public:
+ explicit TaskMgr(sead::Heap* heap);
+ virtual ~TaskMgr();
+
+ void init(s32 num_tasks, sead::Heap* heap, ManagedTaskFactory& factory);
+ void finalize();
+
+ void submitRequest(TaskMgrRequest& request);
+ bool trySubmitRequest(TaskMgrRequest& request);
+
+ bool hasTasks() const;
+
+ void freeTask(ManagedTask* task);
+
+protected:
+ enum class Flag {
+ HeapIsFreeable = 0x1,
+ };
+
+ bool fetchIdleTaskForRequest_(TaskMgrRequest& request, bool retry_until_success);
+ ManagedTask* fetchIdleTask_(bool retry_until_success);
+
+ u8 getListIndex_() const { return mListIndex; }
+ u8 getListIndex2_() const { return ~mListIndex & 1; }
+
+ void swapLists_() {
+ auto lock = sead::makeScopedLock(mCS2);
+ mListIndex = getListIndex2_();
+ mNewFreeTaskEvent.resetSignal();
+ }
+
+ bool tryFetchTaskForRequest_(TaskMgrRequest& request, bool b) {
+ if (!mTasksCS.tryLock())
+ return false;
+
+ const bool ret = fetchIdleTaskForRequest_(request, b);
+ mTasksCS.unlock();
+ return ret;
+ }
+
+ sead::TypedBitFlag<Flag, u8> mFlags;
+ u8 mListIndex = 0;
+ ManagedTask* mTask = nullptr;
+ sead::CriticalSection mTasksCS;
+ sead::CriticalSection mCS2;
+ Event mNewFreeTaskEvent;
+ Event mEvent2;
+ sead::SafeArray<sead::OffsetList<ManagedTask>, 2> mFreeTaskLists;
+ sead::Buffer<ManagedTask*> mTasks;
+};
+KSYS_CHECK_SIZE_NX150(TaskMgr, 0x158);
+
+} // namespace ksys::util