1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
|
#pragma once
#include <container/seadBuffer.h>
#include <container/seadObjList.h>
#include <heap/seadDisposer.h>
#include <prim/seadRuntimeTypeInfo.h>
#include <prim/seadTypedBitFlag.h>
#include <thread/seadAtomic.h>
#include <thread/seadCriticalSection.h>
#include "KingSystem/Utils/Container/UniqueArrayPtr.h"
#include "KingSystem/Utils/Thread/Event.h"
#include "KingSystem/Utils/Thread/MessageDispatcherBase.h"
#include "KingSystem/Utils/Thread/MessageProcessor.h"
namespace sead {
class Thread;
}
namespace ksys {
class Message;
class MessageProcessor;
struct MesTransceiverId;
class MessageQueue {
public:
MessageQueue();
virtual ~MessageQueue();
virtual bool addMessage(const Message& message);
virtual void processQueue(MessageProcessor& processor);
virtual void clear();
private:
Message* findUnusedEntry() const;
util::UniqueArrayPtr<Message, 3000> mMessages;
};
class MessageDispatcher : public MessageDispatcherBase {
SEAD_SINGLETON_DISPOSER(MessageDispatcher)
SEAD_RTTI_OVERRIDE(MessageDispatcher, MessageDispatcherBase)
MessageDispatcher() = default;
~MessageDispatcher() override;
public:
struct InitArg {
// TODO: rename
int num1;
int num2;
int num_bools;
bool set_instance;
};
void init(const InitArg& arg, sead::Heap* heap);
bool isProcessingOnCurrentThread() const;
void registerTransceiver(MessageReceiverEx& receiver) override;
void deregisterTransceiver(MessageReceiverEx& receiver) override;
bool sendMessage(const MesTransceiverId& src, const MesTransceiverId& dest,
const MessageType& type, void* user_data, bool ack, bool) override;
bool sendMessageOnProcessingThread(const MesTransceiverId& src, const MesTransceiverId& dest,
const MessageType& type, void* user_data, bool ack,
bool) override;
bool sendMessage(const MesTransceiverId& src, IMessageBrokerRegister& reg,
const MessageType& type, void* user_data, bool ack, bool) override;
bool sendMessageOnProcessingThread(const MesTransceiverId& src, IMessageBrokerRegister& reg,
const MessageType& type, void* user_data, bool ack,
bool) override;
void update() override;
private:
friend struct AddMessageMainContext;
class DoubleBufferedQueue {
public:
DoubleBufferedQueue();
~DoubleBufferedQueue();
bool addMessage(const Message& message);
void clear();
void processQueue(MessageProcessor& processor);
void swapBuffer() { mActiveIdx ^= 1; }
MessageQueue* getQueue() { return &mBuffer[mActiveIdx ^ 1]; }
private:
u32 mActiveIdx = 1;
MessageQueue mBuffer[2];
};
class MainQueue {
public:
MainQueue();
virtual ~MainQueue();
virtual bool addMessage(const Message& message);
virtual void clear();
virtual void processQueue(MessageProcessor& processor);
private:
DoubleBufferedQueue mQueue;
bool mHasMessageToProcess = false;
};
class Queues {
public:
explicit Queues(MessageProcessor::Logger* logger);
~Queues();
const u32& getId() const { return mId; }
sead::CriticalSection& getCritSection() { return mCritSection; }
const auto& getIdPointers() const { return mTransceiverIdPtrs.mBuffer; }
DoubleBufferedQueue& getQueue() { return mQueue; }
MainQueue& getMainQueue() { return mMainQueue; }
bool isProcessing() const { return mIsProcessing; }
void process();
bool sendMessageOnProcessingThread(const MesTransceiverId& src,
const MesTransceiverId& dest, const MessageType& type,
void* user_data, bool ack);
private:
struct DummyLogger : public MessageProcessor::Logger {
~DummyLogger() override;
void log(const Message& message, bool success) override {}
};
struct TransceiverIdBuffer {
TransceiverIdBuffer();
~TransceiverIdBuffer();
util::UniqueArrayPtr<MesTransceiverId*, 10000> mBuffer;
};
sead::CriticalSection mCritSection;
u32 mId = 0xffffffff;
DummyLogger mDummyLogger;
TransceiverIdBuffer mTransceiverIdPtrs;
DoubleBufferedQueue mQueue;
MainQueue mMainQueue;
MessageProcessor mProcessor;
bool mIsProcessing = false;
};
enum class Flag {
Initialized = 1 << 0,
};
struct Logger : MessageProcessor::Logger {
~Logger() override;
void log(const Message& message, bool success) override;
};
sead::Thread* mProcessingThread{};
Logger mLogger{};
Queues* mQueues{};
sead::TypedBitFlag<Flag> mFlags;
sead::Buffer<u8> mBoolBuffer;
sead::ObjList<u8*> mBools;
sead::CriticalSection mCritSection;
util::Event mUpdateEndEvent;
sead::Atomic<int> mNumEntries = 0;
};
} // namespace ksys
|