-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathMessageControllerAlgorithms.h
More file actions
115 lines (94 loc) · 4 KB
/
Copy pathMessageControllerAlgorithms.h
File metadata and controls
115 lines (94 loc) · 4 KB
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
#pragma once
#include "RxTransport/Message.h"
#include "RxTransport/MessageControllerPolicy.h"
#include "RxTransport/MessageControllerData.h"
#include "RxTransport/Export.h"
namespace BaseLib { namespace Concurrent {
struct DLL_STATE MessageControllerAlgorithms
{
template <typename T>
static bool RouteMessages(const std::list<Message<T>>& messages, const MessageControllerData<T>& data)
{
for(const Message<T>& message : messages)
{
switch(message.data()->AccessorKind())
{
case ChannelAccessorKind::READER:
{
auto writer = data.writers().Find(message.data()->To());
bool result = writer != nullptr
? writer->Next(message)
: false;
break;
}
case ChannelAccessorKind::WRITER:
{
auto reader = data.readers().Find(message.data()->To());
bool result = reader != nullptr
? reader->Next(message)
: false;
break;
}
case ChannelAccessorKind::PUBLISHER:
{
ASSERT(message.data()->To() == MessageAddress::multicast);
auto writer = data.writers().Find(message.data()->From());
bool result = writer != nullptr
? writer->Next(message)
: false;
break;
}
case ChannelAccessorKind::SUBSCRIBER:
{
ASSERT(message.data()->To() == MessageAddress::multicast);
auto reader = data.readers().Find(message.data()->From());
bool result = reader != nullptr
? reader->Next(message)
: false;
break;
}
case ChannelAccessorKind::NO:
IFATAL() << "Message routed to no one: " << message;
break;
}
}
return !messages.empty();
}
template <typename T>
static Message<T> PullRequest(ReaderId id, int n, const std::shared_ptr<details::FlowControllerAction<Message<T>>>& flowController, const ChannelHandle& channelHandle)
{
MessageBuilder builder;
Message<T> message = builder
.Accessor(ChannelAccessorKind::SUBSCRIBER)
.Handle(channelHandle)
.FromTo(id, MessageAddress::multicast)
.With(MessagePolicy::Default())
.Add(std::make_shared<Request>(id, n))
.template Create<T>();
return flowController->Next(message)
? message
: Message<T>::Empty();
}
template <typename T>
static MessageFuture<T> RouteData(const MessagePolicy& policy, DataHolder<T>& messageData, MessageControllerData<T>& data, SequenceNumberCounter& status)
{
// TODO: race condition atomic inc and get
SequenceNumber currentSequenceNumber = status.CurrentState();
++currentSequenceNumber;
messageData.status_.NextState(TransmissionStatusKind::ROUTED, 1);
std::shared_ptr<DataMessage<T>> dataMessage = std::make_shared<DataMessage<T>>(messageData, currentSequenceNumber);
MessageBuilder builder;
Message<T> message = builder
.Accessor(ChannelAccessorKind::PUBLISHER)
.Handle(data.handle())
.FromTo(messageData.writerId_, MessageAddress::multicast)
.With(policy)
.Add(dataMessage)
.template Create<T>();
status.Next<MessageControllerData<T>>(&data, currentSequenceNumber, 1);
return data.flowController()->Next(message)
? MessageFuture<T>(message, dataMessage->DeliveryPromise()->get_future())
: MessageFuture<T>::Empty();
}
};
}}