summaryrefslogtreecommitdiff
path: root/src/KingSystem/Utils/Thread/MessageProcessor.cpp
diff options
context:
space:
mode:
authorLéo Lam <leo@leolam.fr>2021-01-31 15:39:44 +0100
committerLéo Lam <leo@leolam.fr>2021-01-31 16:23:51 +0100
commit8ac67528668d9a7f74e20c29e854449609637f54 (patch)
treeaa5635f38ce4d756094d05490eea1d60587c70a0 /src/KingSystem/Utils/Thread/MessageProcessor.cpp
parent176d68769881ba3a80548c89dc02f2ea4aff814a (diff)
ksys: Add MessageProcessor
Diffstat (limited to 'src/KingSystem/Utils/Thread/MessageProcessor.cpp')
-rw-r--r--src/KingSystem/Utils/Thread/MessageProcessor.cpp59
1 files changed, 59 insertions, 0 deletions
diff --git a/src/KingSystem/Utils/Thread/MessageProcessor.cpp b/src/KingSystem/Utils/Thread/MessageProcessor.cpp
new file mode 100644
index 00000000..296d24a8
--- /dev/null
+++ b/src/KingSystem/Utils/Thread/MessageProcessor.cpp
@@ -0,0 +1,59 @@
+#include "KingSystem/Utils/Thread/MessageProcessor.h"
+#include <tuple>
+#include "KingSystem/Utils/Thread/Message.h"
+#include "KingSystem/Utils/Thread/MessageAck.h"
+#include "KingSystem/Utils/Thread/MessageReceiver.h"
+
+namespace ksys {
+
+MessageProcessor::MessageProcessor(Logger* logger) : mLogger(logger) {}
+
+MessageProcessor::~MessageProcessor() = default;
+
+static bool checkTransceiver(const MesTransceiverId& id) {
+ if (!id.next)
+ return false;
+
+ MesTransceiverId* next = *id.next;
+ if (!next)
+ return false;
+
+ const auto& fields = [](const MesTransceiverId& i) { return std::tie(i.queue_id, i.id); };
+ return fields(id) == fields(*next);
+}
+
+bool MessageProcessor::process(Message* message) {
+ message->decrementDelay();
+
+ if (!message->shouldBeProcessed())
+ return false;
+
+ bool success = false;
+ bool dest_valid = false;
+
+ const auto& dest = message->getDestination();
+ if (checkTransceiver(dest)) {
+ success = dest.receiver->receive(*message) & 1;
+ mLogger->log(*message, success);
+ dest_valid = true;
+ }
+
+ const auto& src = message->getSource();
+ if (!message->hasDelayer() || checkTransceiver(src)) {
+ if (message->shouldAck()) {
+ const auto& source = message->getSource();
+ if (checkTransceiver(source)) {
+ auto* receiver = source.receiver;
+ const MessageAck ack{dest_valid, success, message->getDestination(),
+ message->getType(), message->getUserData()};
+ receiver->receive(ack);
+ }
+ }
+ } else {
+ mLogger->log(*message, false);
+ }
+
+ return true;
+}
+
+} // namespace ksys