CMSDK 2.0.1
Cross-platform C++ base library and SDK for the Psyclone AIOS platform
Loading...
Searching...
No Matches
MemoryQueues.cpp
Go to the documentation of this file.
1
4
5#include "MemoryQueues.h"
6#include "PsyTime.h"
7#include "UnitTestFramework.h"
9
10namespace cmlabs {
11
12// static
13bool MemoryQueues::CreateMessageQueue(uint32& qid, const char* name) {
14 uint32 newQID;
15
16 uint32 newSize = DEFAULTQUEUESIZE;
17 uint32 pageID;
18 char* data = MemoryManager::CreateAndLockNewSystemPage(newSize, pageID);
19 if (data == NULL)
20 return false;
21 if (!MemoryMaps::CreateNewQueue(pageID, name, newQID)) {
22 MemoryManager::UnlockSystemBlock(pageID);
23 MemoryManager::DestroySystemPage(pageID);
24 return false;
25 }
26
28 header->size = newSize;
29 header->id = newQID;
30 header->pageID = pageID;
31 header->count = 0;
32 header->startPos = 0;
33 header->endPos = 0;
34 header->padding = 0;
35 header->id = 0;
36 memset(data + sizeof(MessageQueueHeader), 0, newSize - sizeof(MessageQueueHeader));
37
38 qid = newQID;
39 MemoryManager::UnlockSystemBlock(pageID);
40 return true;
41}
42
43// static
44bool MemoryQueues::GetMessageQueueByName(uint32& qid, const char* name) {
45 return MemoryMaps::GetQueueID(name, qid);
46}
47
48// static
49bool MemoryQueues::AddMessageToQueue(uint32 qid, DataMessage* msg) {
50
51 uint32 pageID = MemoryMaps::GetQueuePageID(qid);
52 uint32 dataSize = 0;
53 char* data = MemoryManager::GetAndLockSystemBlock(pageID, dataSize);
54 if (data == NULL)
55 return false;
56
57 uint32 msgSize = msg->getSize();
58 uint32 spaceToEnd;
60 char* qData = data + sizeof(MessageQueueHeader);
61 uint32 qDataSize = dataSize - sizeof(MessageQueueHeader);
62
63 if (header->endPos >= header->startPos) {
64 spaceToEnd = qDataSize - header->endPos;
65 if (spaceToEnd >= msgSize) {
66 // add msg here
67 memcpy(qData+header->endPos, msg->data, msgSize);
68 header->endPos += msgSize;
69 header->count++;
70 header->padding = 0;
71 }
72 else if (header->startPos >= msgSize) {
73 // fill end rest of space with 0
74 memset(qData+header->endPos, 0, spaceToEnd);
75 header->padding = spaceToEnd;
76 // add msg at buffer start
77 memcpy(qData, msg->data, msgSize);
78 header->endPos = msgSize;
79 header->count++;
80 }
81 else {
82 // we need to resize queue
83 MemoryManager::UnlockSystemBlock(pageID);
84 MemoryManager::ResizeSystemPage(pageID, dataSize*4);
85 data = MemoryManager::GetAndLockSystemBlock(pageID, dataSize);
86 header = (MessageQueueHeader*)data;
87 qData = data + sizeof(MessageQueueHeader);
88 qDataSize = dataSize - sizeof(MessageQueueHeader);
89 // rearrange data if needed
90 if (header->startPos > 0) {
91 memcpy(qData, qData+header->startPos, header->endPos - header->startPos);
92 header->endPos = header->endPos - header->startPos;
93 header->startPos = 0;
94 }
95 MemoryManager::UnlockSystemBlock(pageID);
96 // call this function again...
97 return AddMessageToQueue(qid, msg);
98 }
99 }
100 else {
101 spaceToEnd = header->startPos - header->endPos;
102 if (spaceToEnd >= msgSize) {
103 // add msg here
104 memcpy(qData+header->endPos, msg->data, msgSize);
105 header->endPos += msgSize;
106 header->count++;
107 }
108 else {
109 // we need to resize queue
110 MemoryManager::UnlockSystemBlock(pageID);
111 MemoryManager::ResizeSystemPage(pageID, dataSize*4);
112 uint32 newDataSize = 0;
113 uint32 oldQDataSize = qDataSize;
114 data = MemoryManager::GetAndLockSystemBlock(pageID, newDataSize);
115 header = (MessageQueueHeader*)data;
116 qData = data + sizeof(MessageQueueHeader);
117 qDataSize = newDataSize - sizeof(MessageQueueHeader);
118 // rearrange data
119 // move end data from start to after start data
120 memcpy(qData+oldQDataSize-header->padding, qData, header->endPos);
121 header->endPos = oldQDataSize + header->endPos - header->padding;
122
124 //memcpy(qData+oldQDataSize, qData, header->endPos);
126 //memcpy(qData, qData+header->startPos, oldQDataSize - header->startPos);
128 //memcpy(qData + oldQDataSize - header->startPos, qData+oldQDataSize, header->endPos);
129 //header->endPos = header->endPos + oldQDataSize - header->startPos;
130 //header->startPos = 0;
131 MemoryManager::UnlockSystemBlock(pageID);
132 // call this function again...
133 return AddMessageToQueue(qid, msg);
134 }
135 }
136
137 // Signal a named Semaphore
139 //utils::Semaphore* sem = GetQueueSemaphore(qid);
140 //if (sem) {
141 // sem->signal();
142 // delete(sem);
143 //}
144
145 MemoryManager::UnlockSystemBlock(pageID);
146 return true;
147}
148
149// static
150DataMessage* MemoryQueues::WaitForMessageQueue(uint32 qid, uint32 ms) {
151 uint32 pageID = MemoryMaps::GetQueuePageID(qid);
152 uint32 dataSize = 0;
153 uint64 now, end;
154 MessageQueueHeader* header = (MessageQueueHeader*) MemoryManager::GetAndLockSystemBlock(pageID, dataSize);
155 if (!header)
156 return NULL;
157
158 if (header->count == 0) {
159 // we no data and need to wait for the semaphore...
160 // utils::Semaphore* sem = GetQueueSemaphore(qid);
161 MemoryManager::UnlockSystemBlock(pageID);
162// if (!sem)
163 // return NULL;
164 end = GetTimeNow() + ms*1000;
165
166 // Wait for semaphore to be signalled
167 while (true) {
168 now = GetTimeNow();
169 // if ( (now >= end) || !sem->wait((uint32)((end-now)/1000)) ) {
170 if ( (now >= end) || !utils::WaitForSemaphore(qid, (uint32)((end-now)/1000)) ) {
171 // delete(sem);
172 return NULL;
173 }
174
175 if (!(header = (MessageQueueHeader*) MemoryManager::GetAndLockSystemBlock(pageID, dataSize)))
176 return NULL;
177 if (header->count > 0)
178 break;
179 MemoryManager::UnlockSystemBlock(pageID);
180 }
181 //delete(sem);
182 }
183 char* qData = ((char*)header) + sizeof(MessageQueueHeader);
184 uint32 qDataSize = dataSize - sizeof(MessageQueueHeader);
185
186 // We have data, get it
187 uint32 msgDataSize;
188 DataMessage* msg = NULL;
189 char* msgData = qData + header->startPos;
190 if (((DataMessageHeader*)msgData)->cid == DATAMESSAGEID) {
191 msgDataSize = ((DataMessageHeader*)msgData)->size;
192 // we have data here, get it
193 header->startPos += msgDataSize;
194 header->count--;
195 }
196 else {
197 // we have wrapped around, msg is now at the start of the data
198 msgData = qData;
199 if (((DataMessageHeader*)msgData)->cid == DATAMESSAGEID) {
200 msgDataSize = ((DataMessageHeader*)msgData)->size;
201 // we have data here, get it
202 header->startPos = msgDataSize;
203 header->count--;
204 header->padding = 0;
205 }
206 else {
207 MemoryManager::UnlockSystemBlock(pageID);
208 return NULL;
209 }
210 }
211
212 char* newData = (char*)malloc(msgDataSize);
213 memcpy(newData, msgData, msgDataSize);
214 msg = new DataMessage(newData);
215
216 MemoryManager::UnlockSystemBlock(pageID);
217 return msg;
218}
219
220// static
221uint32 MemoryQueues::GetMessageQueueSize(uint32 qid) {
222 uint32 pageID = MemoryMaps::GetQueuePageID(qid);
223 uint32 dataSize = 0;
224 char* data = MemoryManager::GetAndLockSystemBlock(pageID, dataSize);
225 if (data == NULL)
226 return false;
227
228 MessageQueueHeader* header = (MessageQueueHeader*)data;
229 uint32 qSize = header->count;
230
231 MemoryManager::UnlockSystemBlock(pageID);
232 return qSize;
233}
234
235// static
236bool MemoryQueues::DestroyMessageQueue(uint32 qid) {
237 uint32 pageID = MemoryMaps::GetQueuePageID(qid);
238 // Lock it so no-one else can use it
239 uint32 dataSize = 0;
240 char* data = MemoryManager::GetAndLockSystemBlock(pageID, dataSize);
241
242 MemoryMaps::DeleteQueue(qid);
243 MemoryManager::UnlockSystemBlock(pageID);
244
245 MemoryManager::DestroySystemPage(pageID);
246 return true;
247}
248
249
251
252// static
253bool MemoryQueues::AddRequest(DataMessage* msg, uint32& reqID) {
254 // Create request entry
255 RequestMapEntry* entry = MemoryMaps::CreateAndLockNewRequest(msg->getFrom(), msg->getTo(), reqID);
256 if (!entry)
257 return false;
258 // Add reqID to message
259 msg->setReference(reqID);
260 // Put message into appropriate request queue
261 // Find process for component
262 uint32 procID = ComponentData::GetComponentProcessID(msg->getTo());
263 if (procID == 0) {
264 // if not local, add request to Node Request Queue
266 }
267 // Find Queue ID for process
268 uint32 qID;
269 if (!MemoryMaps::GetProcessQueueID(procID, REQQ_ID, qID)) {
270 MemoryMaps::UnlockRequestMap();
271 return false;
272 }
273 // Add request message to queue
274 if (!AddMessageToQueue(qID, msg)) {
275 MemoryMaps::UnlockRequestMap();
276 return false;
277 }
278 entry->status = REQ_CREATED;
279 MemoryMaps::UnlockRequestMap();
280 return true;
281}
282
283// static
284bool MemoryQueues::AddReply(uint32 reqID, bool success, DataMessage* msg) {
285 uint8 status;
286 RequestMapEntry* entry = MemoryMaps::GetAndLockRequestEntry(reqID);
287 if (!entry)
288 return false;
289 if (msg) {
290 uint64 memID;
291 uint64 msgEOL = msg->getEOL();
292 if (!MemoryManager::InsertMemoryBlock((char*)msg->data, msg->getSize(), msgEOL, memID)) {
293 MemoryMaps::UnlockRequestMap();
294 return false;
295 }
296 entry->dataMessageID = memID;
297 entry->dataMessageEOL = msgEOL;
298 if (success)
299 status = REQ_SUCCESS_DATA;
300 else
301 status = REQ_FAILED_DATA;
302 }
303 else {
304 entry->dataMessageID = 0;
305 entry->dataMessageEOL = 0;
306 if (success)
307 status = REQ_SUCCESS_NODATA;
308 else
309 status = REQ_FAILED_NODATA;
310 }
311 entry->status = status;
312
313 utils::Semaphore* sem = GetRequestSemaphore(reqID);
314 if (sem) {
315 sem->signal();
316 delete(sem);
317 }
318 MemoryMaps::UnlockRequestMap();
319 return true;
320}
321
322// static
323bool MemoryQueues::WaitForReply(uint32 reqID, uint32 ms, uint8& status, DataMessage** outMsg) {
324
325 RequestMapEntry* entry = MemoryMaps::GetAndLockRequestEntry(reqID);
326 if (!entry)
327 return false;
328
329 uint64 now = GetTimeNow();
330 if (entry->status < REQ_REPLY_READY) {
331 // No answer yet, wait for it
332 utils::Semaphore* sem = GetRequestSemaphore(reqID);
333 if (!sem)
334 return NULL;
335 uint64 end = now + ms*1000;
336 MemoryMaps::UnlockRequestMap();
337
338 // Wait for semaphore to be signalled
339 do {
340 if (!sem->wait((uint32)((end-now)/1000))) {
341 delete(sem);
342 return false;
343 }
344
345 entry = MemoryMaps::GetAndLockRequestEntry(reqID);
346 if (!entry)
347 return false;
348
349 if (entry->status > REQ_REPLY_READY)
350 break;
351
352 } while ((now = GetTimeNow()) < end);
353
354 delete(sem);
355 }
356
357 if (entry->status < REQ_REPLY_READY) {
358 MemoryMaps::UnlockRequestMap();
359 return false;
360 }
361
362 // We have an answer, get it
363 status = entry->status;
364 *outMsg = NULL;
365 now = GetTimeNow();
366 if ( (status == REQ_SUCCESS_DATA) || (status == REQ_FAILED_DATA) ) {
367 if (entry->dataMessageEOL > now) {
368 uint32 size = 0;
369 uint64 eol = 0;
370 char* data = MemoryManager::GetCopyMemoryBlock(entry->dataMessageID, size, eol);
371 if (data != NULL)
372 *outMsg = new DataMessage(data);
373 }
374 if (*outMsg == NULL) {
375 // set status to _DATA_EOL
376 status += 2;
377 }
378 }
379
380 MemoryMaps::UnlockRequestMap();
381 MemoryMaps::DeleteRequest(reqID);
382 return true;
383}
384
385// static
386uint32 MemoryQueues::GetRequestCount() {
387 uint32 count;
388 if (MemoryMaps::GetRequestCount(count))
389 return count;
390 else
391 return 0;
392}
393
394
395
396// Get an existing Queue Mutex or creates it
397//utils::Semaphore* MemoryQueues::GetQueueSemaphore(uint32 qID) {
398// char* name = new char[MAXKEYNAMELEN];
399// sprintf(name, "QueueSemaphore_%u", qID);
400// utils::Semaphore* sem = new utils::Semaphore(name);
401// delete [] name;
402// return sem;
403//}
404
405// Get an existing Request Mutex or creates it
406utils::Semaphore* MemoryQueues::GetRequestSemaphore(uint32 reqID) {
407 char* name = new char[MAXKEYNAMELEN];
408 sprintf(name, "RequestSemaphore_%u", reqID);
409 utils::Semaphore* sem = new utils::Semaphore(name);
410 delete [] name;
411 return sem;
412}
413
414
415
417
418bool MemoryQueues::UnitTest() {
419 unittest::progress(0, "init MemoryManager");
420
421 // First create and initialise the MemoryManager
422 MemoryManager* manager = new MemoryManager();
423 uint32 maxPageCount = 100;
424 if (!manager->create(0, maxPageCount)) {
425 unittest::fail("MemoryManager init() failed");
426 delete(manager);
427 return false;
428 }
429
430 uint32 qid1;
431 if (!CreateMessageQueue(qid1)) {
432 unittest::fail("Could not create MessageQueue");
433 return false;
434 }
435 uint32 count = 10000;
436 uint32 writeCount = 0;
437 uint32 readCount = 0;
438 uint32 n;
439 DataMessage* msg;
440 char str[128];
441 uint64 time = GetTimeNow();
442
443 unittest::progress(10, "write phase 1");
444 uint64 wt0 = GetTimeNow();
445 for (n=0; n<count/4; n++) {
446 msg = new DataMessage(CTRL_TEST, writeCount);
447 sprintf(str, "Test%2d", writeCount);
448 msg->setString("TestEntry", str);
449 msg->setTime("TestTime", time);
450 msg->setString("TestEntry2", str);
451 msg->setTime("TestTime2", time);
452 if (!AddMessageToQueue(qid1, msg)) {
453 unittest::fail("Could not add Message %d to Queue", n);
454 return false;
455 }
456 delete(msg);
457 writeCount++;
458 }
459 double wt_us = (double)(GetTimeNow() - wt0);
460 uint32 wroteP1 = count/4;
461
462 if ( (n = GetMessageQueueSize(qid1)) != count/4) {
463 unittest::fail("Queue contains %u messages instead of %u", n, count/4);
464 return false;
465 }
466
467 unittest::progress(30, "read phase 1");
468 uint64 rt0 = GetTimeNow();
469 for (n=0; n<count/8; n++) {
470 if ( (msg = WaitForMessageQueue(qid1, 1000)) == NULL) {
471 unittest::fail("Could not get Message %u from Queue", n);
472 return false;
473 }
474 sprintf(str, "Test%2d", readCount);
475 if (strcmp(str, msg->getString("TestEntry")) != 0) {
476 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, msg->getString("TestEntry"));
477 return false;
478 }
479 if (msg->getTime("TestTime") != time) {
480 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime"));
481 return false;
482 }
483 if (strcmp(str, msg->getString("TestEntry2")) != 0) {
484 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, msg->getString("TestEntry2"));
485 return false;
486 }
487 if (msg->getTime("TestTime2") != time) {
488 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime2"));
489 return false;
490 }
491 readCount++;
492 delete(msg);
493 }
494 double rt_us = (double)(GetTimeNow() - rt0);
495 uint32 readP1 = count/8;
496
497 if ( (n = GetMessageQueueSize(qid1)) != count/8) {
498 unittest::fail("Queue contains %u messages instead of %u", n, count/8);
499 return false;
500 }
501
502 unittest::progress(50, "write phase 2");
503 uint64 wt0b = GetTimeNow();
504 for (n=0; n<3*count/4; n++) {
505 msg = new DataMessage(CTRL_TEST, n);
506 sprintf(str, "Test%2d", writeCount);
507 msg->setString("TestEntry", str);
508 msg->setTime("TestTime", time);
509 msg->setString("TestEntry2", str);
510 msg->setTime("TestTime2", time);
511 if (!AddMessageToQueue(qid1, msg)) {
512 unittest::fail("Could not add Message %d to Queue", n);
513 return false;
514 }
515 delete(msg);
516 writeCount++;
517 }
518 wt_us += (double)(GetTimeNow() - wt0b);
519 uint32 wroteTotal = wroteP1 + 3*count/4;
520
521 if ( (n = GetMessageQueueSize(qid1)) != 7*count/8) {
522 unittest::fail("Queue contains %u messages instead of %u", n, 7*count/8);
523 return false;
524 }
525
526 unittest::progress(70, "read phase 2");
527 uint64 rt0b = GetTimeNow();
528 for (n=0; n<7*count/8; n++) {
529 if ( (msg = WaitForMessageQueue(qid1, 1000)) == NULL) {
530 unittest::fail("Could not get Message %u from Queue", n);
531 return false;
532 }
533 sprintf(str, "Test%2d", readCount);
534 if (strcmp(str, msg->getString("TestEntry")) != 0) {
535 unittest::fail("Message %u from Queue corrupted string '%s' (queue size %u)", readCount, msg->getString("TestEntry"), GetMessageQueueSize(qid1));
536 return false;
537 }
538 unittest::detail("Message %u from Queue '%s'", readCount, msg->getString("TestEntry"));
539 if (msg->getTime("TestTime") != time) {
540 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime"));
541 return false;
542 }
543 if (strcmp(str, msg->getString("TestEntry2")) != 0) {
544 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, msg->getString("TestEntry2"));
545 return false;
546 }
547 if (msg->getTime("TestTime2") != time) {
548 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime2"));
549 return false;
550 }
551 readCount++;
552 delete(msg);
553 }
554 rt_us += (double)(GetTimeNow() - rt0b);
555 uint32 readTotal = readP1 + 7*count/8;
556
557 if ( (n = GetMessageQueueSize(qid1)) != 0) {
558 unittest::fail("Queue contains %u messages instead of 0", n);
559 return false;
560 }
561
562 if (!DestroyMessageQueue(qid1)) {
563 unittest::fail("Could not destroy MessageQueue");
564 return false;
565 }
566
567 if (wt_us > 0)
568 unittest::metric("queue_write_throughput", (double)wroteTotal / wt_us * 1e6, "msg/s", true);
569 if (rt_us > 0)
570 unittest::metric("queue_read_throughput", (double)readTotal / rt_us * 1e6, "msg/s", true);
571 if (wroteTotal)
572 unittest::metric("queue_write_latency", wt_us / (double)wroteTotal, "us", false);
573 if (readTotal)
574 unittest::metric("queue_read_latency", rt_us / (double)readTotal, "us", false);
575
576 unittest::progress(85, "request/reply queue");
577
578 // First create two components
579 uint16 procID1, procID2;
580 uint32 compID1 = 1, compID2 = 2;
581 uint32 qID;
582 //if (!MemoryMaps::CreateNewNode(1, "", 1)) {
583 // printf("Error: Testing Request Queue - CreateNewNode\n");
584 // return false;
585 //}
586 if (!MemoryMaps::CreateNewProcess("Proc1", procID1)) {
587 unittest::fail("Request Queue - CreateNewProcess");
588 return false;
589 }
590 if (!MemoryMaps::CreateNewProcess("Proc2", procID2)) {
591 unittest::fail("Request Queue - CreateNewProcess 2");
592 return false;
593 }
594 if (!MemoryQueues::CreateMessageQueue(qID)) {
595 unittest::fail("Request Queue - CreateMessageQueue");
596 return false;
597 }
598 if (!MemoryMaps::SetProcessQueueID(procID1, MSGQ_ID, qID)) {
599 unittest::fail("Request Queue - SetProcessQueueID");
600 return false;
601 }
602 if (!MemoryQueues::CreateMessageQueue(qID)) {
603 unittest::fail("Request Queue - CreateMessageQueue");
604 return false;
605 }
606 if (!MemoryMaps::SetProcessQueueID(procID2, REQQ_ID, qID)) {
607 unittest::fail("Request Queue - SetProcessQueueID");
608 return false;
609 }
610
611 ComponentData* compData1 = ComponentData::CreateComponent(compID1, "Comp1", 10*1024, 1, procID1);
612 if (compData1 == NULL) {
613 unittest::fail("Request Queue - CreateComponent1");
614 return false;
615 }
616 ComponentData* compData2 = ComponentData::CreateComponent(compID2, "Comp2", 10*1024, 1, procID2);
617 if (compData2 == NULL) {
618 unittest::fail("Request Queue - CreateComponent2");
619 return false;
620 }
621
622 unittest::progress(92, "request/reply roundtrip");
623 DataMessage* reqMsg = new DataMessage(CTRL_TEST, compID1, compID2);
624 uint32 reqID = 0;
625 if (!MemoryQueues::AddRequest(reqMsg, reqID)) {
626 unittest::fail("Request Queue - AddRequest");
627 return false;
628 }
629 delete(reqMsg);
630
631 if (!MemoryMaps::GetProcessQueueID(procID2, REQQ_ID, qID)) {
632 unittest::fail("Request Queue - GetProcessQueueID");
633 return false;
634 }
635 reqMsg = MemoryQueues::WaitForMessageQueue(qID, 1000);
636 if (reqMsg == NULL) {
637 unittest::fail("Request Queue - WaitForMessageQueue");
638 return false;
639 }
640
641 uint32 reqID2 = (uint32)reqMsg->getReference();
642 if (reqID != reqID2) {
643 unittest::fail("Request Queue - getReference mismatch");
644 return false;
645 }
646
647 DataMessage* replyMsg = new DataMessage(CTRL_TEST, compID2, 0, 1000000000);
648 if (!MemoryQueues::AddReply(reqID, true, replyMsg)) {
649 unittest::fail("Request Queue - AddReply");
650 return false;
651 }
652 delete(reqMsg);
653 delete(replyMsg);
654 replyMsg = NULL;
655
656 uint8 status;
657 if (!MemoryQueues::WaitForReply(reqID2, 1000, status, &replyMsg)) {
658 unittest::fail("Request Queue - WaitForReply");
659 return false;
660 }
661
662 if (status != REQ_SUCCESS_DATA) {
663 unittest::fail("Request Queue - Status mismatch");
664 return false;
665 }
666
667 if (replyMsg == NULL) {
668 unittest::fail("Request Queue - ReplyMsg NULL");
669 return false;
670 }
671
672 delete(replyMsg);
673 delete(compData1);
674 delete(compData2);
675
676 unittest::progress(100, "done");
677 delete(manager);
678 return true;
679}
680
681
682bool MemoryQueues::UnitTestQueues() {
683 printf("Testing Multithreaded Memory Queues...\n\n");
684
685 // First create and initialise the MemoryManager
686 MemoryManager* manager = new MemoryManager();
687 if (!manager->create(0, 30)) {
688 fprintf(stderr, "MemoryManager init() failed...\n");
689 delete(manager);
690 return false;
691 }
692
693 uint64 t1, t2, t3;
694 DataMessage* msg;
695 int64 c = 0;
696 uint32 q1, q2;
697
698 if (!MemoryQueues::CreateMessageQueue(q1, "Q1")) {
699 LogPrint(0, 0, 0, "[1] Could not create Test Queue 1...");
700 delete(manager);
701 return false;
702 }
703 if (!MemoryQueues::CreateMessageQueue(q2, "Q2")) {
704 LogPrint(0, 0, 0, "[1] Could not create Test Queue 2...");
705 delete(manager);
706 return false;
707 }
708
709 uint32 queueTestThreadID;
710 if (!ThreadManager::CreateThread(QueueTest, NULL, queueTestThreadID)) {
711 LogPrint(0, 0, 0, "[1] Could not create Queue Test thread...");
712 delete(manager);
713 return false;
714 }
715
716 msg = new DataMessage();
717 msg->setInt("Counter", 0);
718
719 if (!MemoryQueues::AddMessageToQueue(q1, msg)) {
720 LogPrint(0, 0, 0, "[1] Could not add initial message to Q1...");
721 delete(manager);
722 return false;
723 }
724
725 uint64 start = 0;
726
727 while (true) {
728 t1 = GetTimeNow();
729 if (msg = MemoryQueues::WaitForMessageQueue(q2, 50)) {
730 t2 = GetTimeNow();
731 msg->getInt("Counter", c);
732 if (c && (c % 99999 == 0)) {
733 if (start)
734 LogPrint(0,0,0,"[1] Average path time: %.3fus", (double)GetTimeAge(start)/99999.0);
735 start = GetTimeNow();
736 }
737 msg->setInt("Counter", c+1);
738 msg->setSendTime(GetTimeNow());
739 if (!MemoryQueues::AddMessageToQueue(q1, msg)) {
740 LogPrint(0, 0, 0, "[1] Could not add message %lld to Q1...", c+1);
741 delete(msg);
743 }
744 t3 = GetTimeNow();
745 if (t3-t1 > 10000)
746 LogPrint(0,0,0,"[1 - %lld] Wait: %lld Add: %lld Total: %lld Msg: %lld\n", c+1, t2-t1, t3-t2, t3-t1, t3-msg->getSendTime());
747 delete(msg);
748 }
749 else {
750 t3 = GetTimeNow();
751 LogPrint(0,0,0,"[1 - %lld] Timeout: %lld\n", c+1, t3-t1);
752 }
753 }
754
755 return true;
756}
757
761
762 uint64 t1, t2, t3;
763 DataMessage* msg;
764 int64 c = 0;
765 uint32 q1, q2;
766
767 if (!MemoryQueues::GetMessageQueueByName(q1, "Q1")) {
768 LogPrint(0, 0, 0, "[2] Could not Get Test Queue 1...");
770 }
771
772 if (!MemoryQueues::GetMessageQueueByName(q2, "Q2")) {
773 LogPrint(0, 0, 0, "[2] Could not Get Test Queue 2...");
775 }
776
777 while (true) {
778 t1 = GetTimeNow();
779 if (msg = MemoryQueues::WaitForMessageQueue(q1, 50)) {
780 t2 = GetTimeNow();
781 msg->getInt("Counter", c);
782 msg->setInt("Counter", c+1);
783 msg->setSendTime(GetTimeNow());
784 if (!MemoryQueues::AddMessageToQueue(q2, msg)) {
785 LogPrint(0, 0, 0, "[2] Could not add message %lld to Q2...", c+1);
786 delete(msg);
788 }
789 t3 = GetTimeNow();
790 if (t3-t1 > 10000)
791 LogPrint(0,0,0,"[2 - %lld] Wait: %lld Add: %lld Total: %lld Msg: %lld\n", c+1, t2-t1, t3-t2, t3-t1, t3-msg->getSendTime());
792 delete(msg);
793 }
794 else {
795 t3 = GetTimeNow();
796 LogPrint(0,0,0,"[2 - %lld] Timeout: %lld\n", c+1, t3-t1);
797 }
798 }
799
801}
802
803void Register_MemoryQueues_Tests() {
804 UnitTestRunner::instance().registerTest("memoryqueues", MemoryQueues::UnitTest,
805 "Message queue add/wait/size, request/reply roundtrip", "memory");
806 // NOTE: MemoryQueues::UnitTestQueues() is NOT registered — it runs an
807 // unbounded "while (true)" message ping-pong loop (and spawns QueueTest,
808 // which also loops forever) with no termination condition, so it never
809 // returns. Registering it would hang the test runner.
810}
811
812} // namespace cmlabs
#define REQ_FAILED_NODATA
#define REQ_REPLY_READY
#define REQ_CREATED
#define REQ_SUCCESS_NODATA
#define REQ_SUCCESS_DATA
#define REQ_FAILED_DATA
#define NODE_PROCESS_ID
Reserved process id of the node process itself.
#define MSGQ_ID
Data-message queue slot.
Definition MemoryMaps.h:43
#define REQQ_ID
Request queue slot.
Definition MemoryMaps.h:44
Ring-buffer message queues stored in shared memory, plus a static request/reply facility.
#define DATAMESSAGEID
Definition ObjectIDs.h:75
CMSDK time: µs-resolution 64-bit timestamps and the Time Mapping Constant (TMC).
Small, dependency-free unit test harness used by all CMSDK object tests.
#define MAXKEYNAMELEN
Definition Utils.h:85
#define thread_ret_val(ret)
Definition Utils.h:131
#define THREAD_RET
Definition Utils.h:127
#define THREAD_FUNCTION_CALL
Definition Utils.h:129
#define LogPrint
Definition Utils.h:313
#define THREAD_ARG
Definition Utils.h:130
Legacy accessor wrapping one component record inside a shared page.
static ComponentData * CreateComponent(uint32 id, const char *name, uint32 size, uint16 nodeID, uint16 procID)
static uint16 GetComponentProcessID(uint32 cid)
The central Psyclone data container: a self-contained binary message with typed, named user entries.
Top-level facade of the shared-memory subsystem for one process.
static MemoryManager * Singleton
Per-process singleton instance, set by the constructor.
static bool CreateThread(THREAD_FUNCTION func, void *args, uint32 &newID, uint32 reqID=0)
Create a new native thread and start it immediately.
static UnitTestRunner & instance()
Access the singleton (created on first use).
void registerTest(const char *name, UnitTestFunc func, const char *description="", const char *category="", bool inDefaultRun=true)
Register a test with the runner.
Counting semaphore, optionally named for cross-process use.
Definition Utils.h:502
uint64 GetTimeNow()
Return the current absolute time (µs since year 0) according to the TMC.
Definition PsyTime.cpp:69
int64 GetTimeAge(uint64 t)
Age of a timestamp relative to now.
Definition PsyTime.cpp:25
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
void fail(const char *fmt,...)
Set an explanatory reason shown on the FAIL line.
void metric(const char *name, double value, const char *unit="", bool higherIsBetter=true)
Record a performance metric.
void detail(const char *fmt,...)
Verbose-only indented diagnostic line (shown only when verbose=1).
void progress(int percent, const char *action)
Report progress with a short description of the current action.
static struct PsyType CTRL_TEST
Definition ObjectIDs.h:82
static THREAD_RET THREAD_FUNCTION_CALL QueueTest(THREAD_ARG arg)
The current (version 10) DataMessage wire/shared-memory header.
On-segment header of a circular message queue (legacy 32-bit layout).
uint32 size
Total size of the queue region in bytes (header + buffer).
uint32 count
Number of messages currently in the queue.
Request-map entry tracking one cross-component request/reply transaction.