diff options
| author | Léo Lam <leo@leolam.fr> | 2021-01-31 15:39:44 +0100 |
|---|---|---|
| committer | Léo Lam <leo@leolam.fr> | 2021-01-31 16:23:51 +0100 |
| commit | 8ac67528668d9a7f74e20c29e854449609637f54 (patch) | |
| tree | aa5635f38ce4d756094d05490eea1d60587c70a0 /src/KingSystem/Utils/Thread/MessageProcessor.cpp | |
| parent | 176d68769881ba3a80548c89dc02f2ea4aff814a (diff) | |
ksys: Add MessageProcessor
Diffstat (limited to 'src/KingSystem/Utils/Thread/MessageProcessor.cpp')
| -rw-r--r-- | src/KingSystem/Utils/Thread/MessageProcessor.cpp | 59 |
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 |
