CMSDK 2.0.1
Cross-platform C++ base library and SDK for the Psyclone AIOS platform
Loading...
Searching...
No Matches
MemoryRequestConnection.cpp
Go to the documentation of this file.
1
4
6
7namespace cmlabs {
8
10 conID = 0;
11 serverHeader = NULL;
12 conHeader = NULL;
13 queueHeader = NULL;
14 requestHeader = NULL;
15 mutex = NULL;
16 oneWaitingForQueue = false;
17}
18
22
23bool MemoryRequestConnection::connect(const char* memoryName) {
24 mutex = new utils::Mutex(utils::StringFormat("%s_Mutex", memoryName).c_str());
25 if (!mutex || !mutex->enter(1000))
26 return false;
27
28 // Allocate main shared memory segment
29 serverHeader = (RequestServerHeader*) utils::OpenSharedMemorySegment(memoryName, 0);
30 if (!serverHeader) {
31 mutex->leave();
32 return false;
33 }
34
35 serverHeader->lastUpdateTime = serverHeader->createdTime;
36
37 requestHeader = ((char*)serverHeader)+sizeof(RequestServerHeader)
38 + (serverHeader->maxConnectionCount * serverHeader->eachConnectionSize);
39 if (!MemoryRequestQueues::CheckRequestMap(requestHeader, serverHeader->maxRequestCount)) {
40 requestHeader = NULL;
41 mutex->leave();
42 shutdown();
43 return false;
44 }
45
46 // Now find the first available connection id
47 conID = 1;
48 conHeader = (RequestConnectionHeader*) (((char*)serverHeader)+sizeof(RequestServerHeader));
49 while (conID <= serverHeader->maxConnectionCount) {
50 if (!conHeader->size || !conHeader->conid || !conHeader->createdTime) {
51 conHeader->size = serverHeader->eachConnectionSize;
52 conHeader->conid = conID;
53 conHeader->createdTime = conHeader->lastUpdateTime = GetTimeNow();
54 conHeader->status = 1;
55 break;
56 }
57 else {
58 conID++;
59 conHeader = (RequestConnectionHeader*) (((char*)conHeader)+conHeader->size);
60 }
61 }
62 if (conID > serverHeader->maxConnectionCount) {
63 conHeader = NULL;
64 mutex->leave();
65 shutdown();
66 return false;
67 }
68 queueHeader = ((char*)conHeader)+sizeof(RequestConnectionHeader);
69 if (!MemoryRequestQueues::InitQueue(queueHeader, serverHeader->maxQueueSize)) {
70 mutex->leave();
71 shutdown();
72 return false;
73 }
74
75 mutex->leave();
76 return true;
77}
78
80 conID = 0;
81 if (!mutex)
82 return true;
83 if (!mutex->enter(1000))
84 return false;
85
86 utils::Mutex* tempMutex = mutex;
87 mutex = NULL;
88 requestHeader = NULL;
89 queueHeader = NULL;
90 conHeader = NULL;
91 if (serverHeader) {
92 utils::CloseSharedMemorySegment((char*) serverHeader, serverHeader->size);
93 serverHeader = NULL;
94 }
95
96 std::map<uint64, DataMessage*>::iterator i = replyCache.begin(), e = replyCache.end();
97 while (i != e) {
98 delete(i->second);
99 i++;
100 }
101 replyCache.clear();
102
103 tempMutex->leave();
104 delete(tempMutex);
105 return true;
106}
107
108
110 if (!mutex || !mutex->enter(1000))
111 return false;
112 uint32 to = msg->getTo();
113 uint64 reqID = MemoryRequestQueues::AddNewRequest(requestHeader, MEMORYREQUEST_INIT, conID, to);
114 if (!reqID) {
115 mutex->leave();
116 return 0;
117 }
118
119 char* destQueueHeader = getQueueHeader(to);
120 if (!destQueueHeader || !MemoryRequestQueues::CheckQueue(destQueueHeader, serverHeader->maxQueueSize)) {
121 MemoryRequestQueues::CloseRequest(requestHeader, reqID);
122 mutex->leave();
123 return 0;
124 }
125
126 msg->setFrom(conID);
127 msg->setReference(reqID);
128 if (!MemoryRequestQueues::AddToQueue(destQueueHeader, msg)) {
129 MemoryRequestQueues::CloseRequest(requestHeader, reqID);
130 mutex->leave();
131 return 0;
132 }
134 utils::SignalSemaphore(utils::StringFormat("%s_%u", serverHeader->name, to).c_str());
135 mutex->leave();
136 return reqID;
137}
138
139uint16 MemoryRequestConnection::getRequestStatus(uint64 id, uint64& createdTime, uint64& lastUpdateTime) {
140 if (!mutex || !mutex->enter(1000))
141 return MEMORYREQUEST_ERROR;
142 uint16 status = MemoryRequestQueues::GetRequestStatus(requestHeader, id, createdTime, lastUpdateTime);
143 mutex->leave();
144 return status;
145}
146
147
148DataMessage* MemoryRequestConnection::waitForQueue(uint32 ms) {
149 // We assume the mutex is entered
150 // If return NULL the mutex will be free, otherwise locked
152 if (!msg) {
153 mutex->leave();
154 if (!utils::WaitForSemaphore(utils::StringFormat("%s_%u", serverHeader->name, conID).c_str(), ms, true))
155 return NULL;
156 if (!mutex || !mutex->enter(1000))
157 return NULL;
158 if (!(msg = MemoryRequestQueues::GetNextQueueMessage(queueHeader))) {
159 mutex->leave();
160 return NULL;
161 }
162 }
163 return msg;
164}
165
167 if (!mutex || !mutex->enter(1000))
168 return NULL;
169
170 DataMessage* msg = waitForQueue(ms);
171 if (!msg) {
172 // mutex already unlocked
173 return NULL;
174 }
175 // mutex still locked
176
178 MemoryRequestQueues::CloseRequest(requestHeader, msg->getReference());
179
180 mutex->leave();
181 return msg;
182}
183
185 if (!mutex || !mutex->enter(1000))
186 return NULL;
187
188 DataMessage* msg = waitForQueue(ms);
189 if (!msg) {
190 // mutex already unlocked
191 return NULL;
192 }
193 // mutex still locked
194
196 mutex->leave();
197 return msg;
198}
199
200bool MemoryRequestConnection::setRequestStatus(uint64 id, uint16 status) {
201 if (!mutex || !mutex->enter(1000))
202 return false;
203 if (!MemoryRequestQueues::SetRequestStatus(requestHeader, id, status)) {
204 mutex->leave();
205 return false;
206 }
207 mutex->leave();
208 return true;
209}
210
212 if (!msg->getTo() || (msg->getTo() == conID))
213 return false;
214 if (!mutex || !mutex->enter(1000))
215 return false;
216 uint64 createdTime, lastUpdateTime;
217 uint16 status = MemoryRequestQueues::GetRequestStatus(requestHeader, id, createdTime, lastUpdateTime);
218 if (status < MEMORYREQUEST_SUBMITTED) {
219 mutex->leave();
220 return false;
221 }
222
223 uint32 to = msg->getTo();
224 char* destQueueHeader = getQueueHeader(to);
225 if (!destQueueHeader || !MemoryRequestQueues::CheckQueue(destQueueHeader, serverHeader->maxQueueSize)) {
227 mutex->leave();
228 return 0;
229 }
230
231 msg->setFrom(conID);
232 msg->setReference(id);
233 if (!MemoryRequestQueues::AddToQueue(destQueueHeader, msg)) {
235 mutex->leave();
236 return 0;
237 }
239 utils::SignalSemaphore(utils::StringFormat("%s_%u", serverHeader->name, to).c_str());
240 mutex->leave();
241
242 return true;
243}
244
245
246char* MemoryRequestConnection::getQueueHeader(uint32 otherConID) {
247 if (!otherConID || (otherConID > serverHeader->maxConnectionCount))
248 return NULL;
249 char* otherQueueHeader = ((char*)serverHeader)
250 + sizeof(RequestServerHeader)
251 + (otherConID-1) * (serverHeader->eachConnectionSize)
252 + sizeof(RequestConnectionHeader);
253 if (!MemoryRequestQueues::CheckQueue(otherQueueHeader, serverHeader->maxQueueSize))
254 return NULL;
255 return otherQueueHeader;
256}
257
258
259
261 if (!mutex || !mutex->enter(1000))
262 return NULL;
263
264 uint64 ref;
265 DataMessage* msg = NULL;
266 uint64 startTime = GetTimeNow();
267 int32 timeLeft = ms;
268
269 while (true) {
270 if (msg = replyCache[id]) {
271 replyCache.erase(id);
273 MemoryRequestQueues::CloseRequest(requestHeader, msg->getReference());
274 mutex->leave();
275 return msg;
276 }
277
278 if ((timeLeft = ms - GetTimeAgeMS(startTime)) < 0) {
279 mutex->leave();
280 return NULL;
281 }
282
283 // is anyone already waiting for the queue?
284 if (!oneWaitingForQueue) {
285 // then wait for the queue
286 oneWaitingForQueue = true;
287 msg = waitForQueue((uint32)timeLeft);
288 oneWaitingForQueue = false;
289 if (!msg) {
290 // mutex already unlocked
291 return NULL;
292 }
293 // mutex still locked
294 if ((ref = msg->getReference()) == id) {
295 // this is for me, thanks
297 MemoryRequestQueues::CloseRequest(requestHeader, msg->getReference());
298 mutex->leave();
299 return msg;
300 }
301 else {
302 // this is for someone else
303 replyCache[ref] = msg;
304 cacheEvent.signal();
305 continue;
306 }
307 }
308 else {
309 mutex->leave();
310 // wait for the event that someone else got a message not for them
311 if (!cacheEvent.waitNext(timeLeft))
312 return NULL;
313 else {
314 if (!mutex || !mutex->enter(1000))
315 return NULL;
316 continue;
317 }
318 }
319 };
320}
321
322
323
324
325
326
327
328
329} // namespace cmlabs
Client-side endpoint of the named shared-memory request/reply service.
#define MEMORYREQUEST_RECEIVED
Dequeued by the server.
#define MEMORYREQUEST_ERROR
Error/invalid.
#define MEMORYREQUEST_INIT
Request being composed.
#define MEMORYREQUEST_SUBMITTED
Placed on the server's queue.
#define MEMORYREQUEST_REPLIED
Reply enqueued to the requester.
#define MEMORYREQUEST_REPLYRECEIVED
Requester consumed the reply.
The central Psyclone data container: a self-contained binary message with typed, named user entries.
uint64 getReference()
getReference() Get and return message reference id
bool setFrom(uint32 from)
setFrom(uint32 from)
bool setReference(uint64 ref)
setReference(uint64 ref) Set message reference
uint32 getTo()
getTo()
bool replyToRequest(uint64 id, DataMessage *msg)
Server side: send the reply for id (copied).
bool setRequestStatus(uint64 id, uint16 status)
Server side: advance a request's status.
uint16 getRequestStatus(uint64 id, uint64 &createdTime, uint64 &lastUpdateTime)
Query a request's status.
DataMessage * waitForRequestReply(uint32 ms)
Wait for the reply to this connection's most recent request.
bool connect(const char *memoryName)
Attach to the named server segment and claim a connection slot.
bool shutdown()
Release the connection slot and detach.
uint64 makeRequest(DataMessage *msg)
Submit a request to the server.
DataMessage * waitForRequest(uint32 ms)
Server side: wait for the next incoming request.
static uint64 AddNewRequest(char *data, uint16 status, uint32 from, uint32 to)
Allocate a request slot.
static DataMessage * GetNextQueueMessage(char *data)
Pop the next message.
static bool AddToQueue(char *data, DataMessage *msg)
Append a copy of msg (caller keeps ownership).
static bool CloseRequest(char *data, uint64 id)
Release the request slot.
static bool SetRequestStatus(char *data, uint64 id, uint16 status)
Advance a request's MEMORYREQUEST_* status.
static uint32 GetRequestStatus(char *data, uint64 id, uint64 &createdTime, uint64 &lastUpdateTime)
static bool CheckQueue(char *data, uint32 size)
Validate an existing queue region.
static bool CheckRequestMap(char *data, uint32 count)
Validate an existing status map.
static bool InitQueue(char *data, uint32 size)
Initialise a queue region of size bytes.
Recursive mutual-exclusion lock, optionally named for cross-process use.
Definition Utils.h:463
bool leave()
Release the mutex.
Definition Utils.cpp:1194
bool enter()
Block until the mutex is acquired.
Definition Utils.cpp:1059
uint64 GetTimeNow()
Return the current absolute time (µs since year 0) according to the TMC.
Definition PsyTime.cpp:69
int32 GetTimeAgeMS(uint64 t)
Age of a timestamp relative to now, in milliseconds.
Definition PsyTime.cpp:35
bool SignalSemaphore(const char *name)
Signal a named global semaphore.
Definition Utils.cpp:804
bool WaitForSemaphore(const char *name, uint32 ms, bool autocreate=true)
Wait on a named global semaphore.
Definition Utils.cpp:765
char * OpenSharedMemorySegment(const char *name, uint64 size)
Open and map an existing named shared memory segment.
Definition Utils.cpp:2259
bool CloseSharedMemorySegment(char *data, uint64 size)
Unmap a segment previously created/opened here.
Definition Utils.cpp:2374
std::string StringFormat(const char *format,...)
printf into a std::string.
Definition Utils.cpp:6626
Header of one client connection area within a request-server segment.
Header of a request-server segment: identity, limits and layout parameters.