CMSDK 2.0.1
Cross-platform C++ base library and SDK for the Psyclone AIOS platform
Loading...
Searching...
No Matches
ProcessMemory.cpp
Go to the documentation of this file.
1
4
5#include "ProcessMemory.h"
6#include "UnitTestFramework.h"
7#include "ThreadManager.h"
8
9namespace cmlabs {
10
11// Cooperative stop for the ProcessMemory unit-test slave threads (QueueTest /
12// ProcessMemoryPerfTest). The test sets this true and waits for the worker to
13// exit before freeing its MemoryManager, so the worker never touches freed
14// memory (which would otherwise hang the next test in the suite).
15static volatile bool g_processTestStop = false;
16static uint32 g_processTestThreadID = 0;
17
25
27 memset(qSemaphores, 0, MAXPROC4*sizeof(utils::Semaphore*));
28 mutex = NULL;
29 this->master = master;
30 header = NULL;
31 memorySize = 0;
32 port = 0;
33 serial = 0;
34}
35
37 if (mutex)
38 mutex->enter(5000, __FUNCTION__);
39 if (memorySize)
40 utils::CloseSharedMemorySegment((char*) header, memorySize);
41 memset(qSemaphores, 0, MAXPROC4*sizeof(utils::Semaphore*));
42 header = NULL;
43 if (mutex)
44 mutex->leave();
45 delete(mutex);
46 mutex = NULL;
47}
48
49bool ProcessMemory::getMemoryUsage(uint64& alloc, uint64& usage) {
50 if (!mutex || !mutex->enter(5000, __FUNCTION__))
51 return false;
53
54// printf(" --------- Get ProcessMap usage size: %u header: %u usage: %u serial %u --------\n", memorySize, header->size, header->usage, serial);
55
56 alloc = memorySize;
57 usage = header->usage;
58
59 mutex->leave();
60 return true;
61}
62
63bool ProcessMemory::checkProcessHeartbeats(uint32 timeoutMS, std::list<ProcessInfoStruct>& procIssues) {
64 if (!mutex || !mutex->enter(5000, __FUNCTION__))
65 return false;
67
68 bool issuesFound = false;
69 ProcessInfoStruct* pInfo = &(header->processes[1]);
70 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
71 if (pInfo->createdTime && (pInfo->nodeID == this->master->getNodeID())) {
72 //if (!pInfo->lastseen) {
73 // procIssues.push_back(*pInfo);
74 // issuesFound = true;
75 //}
76 if (pInfo->status > PSYPROC_ACTIVE) {
77 procIssues.push_back(*pInfo);
78 issuesFound = true;
79 }
80 else if ((pInfo->status > PSYPROC_IDLE) && (GetTimeAgeMS(pInfo->lastseen) > (int64)timeoutMS)) {
81 procIssues.push_back(*pInfo);
82 issuesFound = true;
83 }
84 }
85 }
86 mutex->leave();
87 return (!issuesFound);
88}
89
90std::vector<ProcessInfoStruct>* ProcessMemory::getAllProcesses() {
91 std::vector<ProcessInfoStruct>* processes = new std::vector<ProcessInfoStruct>;
92 if (!mutex || !mutex->enter(5000, __FUNCTION__))
93 return processes;
95
96 uint16 thisNodeID = this->master->getNodeID();
97 ProcessInfoStruct* pInfo = &(header->processes[1]);
98 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
99 if (pInfo->createdTime && (pInfo->nodeID == thisNodeID))
100 processes->push_back(*pInfo);
101 }
102 mutex->leave();
103 return processes;
104}
105
106bool ProcessMemory::addLocalPerformanceStats(std::list<PerfStats> &perfStats) {
107 if (!mutex || !mutex->enter(5000, __FUNCTION__))
108 return false;
110
111 uint16 thisNodeID = this->master->getNodeID();
112
113 PerfStats stats;
114 ProcessInfoStruct* pInfo = &(header->processes[1]);
115 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
116 if (pInfo->createdTime && (pInfo->nodeID == thisNodeID)) {
117 memset(&stats, 0, sizeof(PerfStats));
118 stats.spaceID = pInfo->id;
119 stats.nodeID = pInfo->nodeID;
120 stats.osID = pInfo->osID;
121 // stats.type = stats.space;
122 if (pInfo->type == 1) // inside node
123 stats.currentCPUTicks = 0;
124 else
126 // stats.currentMemoryBytes = pInfo->stats.procMemUsage;
128 stats.totalInputBytes = pInfo->stats.msgInBytes;
129 stats.totalInputCount = pInfo->stats.msgInCount;
130 stats.totalOutputBytes = pInfo->stats.msgOutBytes;
131 stats.totalOutputCount = pInfo->stats.msgOutCount;
132 // stats.runCount = 0;
133 getQueueSizes(pInfo->id, stats.totalQueueBytes, stats.totalQueueCount);
134 // stats.currentRunStartTime = 0;
135 stats.firstRunStartTime = pInfo->createdTime;
136 // stats.totalCycleCount = 0;
137 // stats.totalRunCount = 0;
138 // stats.migrationCount = 0;
139 perfStats.push_back(stats);
140 }
141 }
142 mutex->leave();
143
144 return true;
145}
146
147bool ProcessMemory::getQueueSizes(uint16 procID, uint64 &bytes, uint32 &count) {
148 if (!mutex || !mutex->enter(5000, __FUNCTION__))
149 return 0;
151
152 bytes = 0;
153 count = 0;
154
155 MessageQueueHeader* qHeader;
156 for (uint8 t = 1; t<5; t++) {
157 qHeader = getQHeader(procID, t);
158 if (qHeader) {
159 count += qHeader->count;
160 if (qHeader->count) {
161 if (qHeader->endPos > qHeader->startPos)
162 bytes = qHeader->endPos - qHeader->startPos;
163 else
164 bytes = qHeader->startPos - qHeader->endPos;
165 }
166 }
167 }
168 mutex->leave();
169 return true;
170}
171
172bool ProcessMemory::create(uint32 initialProcCount) {
173 serial = master->incrementProcessShmemSerial();
174 port = master->port;
175
176// LogPrint(0,LOG_MEMORY,0," --------- Creating ProcessMemory section serial %u --------", serial);
177 if (!mutex)
178 mutex = new utils::Mutex(utils::StringFormat("PsycloneProcessMemoryMutex_%u", port).c_str(), true);
179 if (!mutex->enter(5000, __FUNCTION__))
180 return false;
181
182 memorySize = sizeof(ProcessMemoryStruct) + (initialProcCount * INITIALQSIZE * 4);
183
184 //header = (ProcessMemoryStruct*) utils::OpenSharedMemorySegment(utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), memorySize);
185 //if (header) {
186 // utils::CloseSharedMemorySegment((char*)header, memorySize);
187 // LogPrint(0,LOG_MEMORY,2,"ProcessMemory removing stale shared memory (%u/%u)...", port, serial);
188 //}
189
190 header = (ProcessMemoryStruct*) utils::CreateSharedMemorySegment(utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), memorySize, true);
191 if (!header) {
192 LogPrint(0,LOG_MEMORY,0,"Cannot create memory");
193 mutex->leave();
194 return false;
195 }
196 memset(header, 0, (size_t)memorySize);
197 header->size = memorySize;
198 header->cid = PROCESSMEMORYID;
199 header->createdTime = GetTimeNow();
200 header->usage = sizeof(ProcessMemoryStruct);
201 // this->port = port;
202
203// printf(" --------- Creating new ProcessMap header: %u usage: %u serial %u --------\n", header->size, header->usage, serial);
204
205 // statically initialise base process id 0
206 ProcessInfoStruct* pInfo = &(header->processes[0]);
207 pInfo->createdTime = pInfo->lastseen = GetTimeNow();
208 pInfo->status = PSYPROC_ACTIVE;
209 pInfo->nodeID = master->getNodeID();
210 utils::strcpyavail(pInfo->name, "Node", MAXKEYNAMELEN, true);
211
212 uint16 cmdQID;
213 uint16 msgQID;
214 uint16 sigQID;
215 uint16 reqQID;
216
217 if (
218 !setupNextAvailableQ(cmdQID) ||
219 !setupNextAvailableQ(msgQID) ||
220 !setupNextAvailableQ(sigQID) ||
221 !setupNextAvailableQ(reqQID) ) {
222 mutex->leave();
223 LogPrint(0,LOG_MEMORY,0,"Cannot setup queues");
224 return false;
225 }
226
227 // reset as header might have changed...
228 pInfo = &(header->processes[0]);
229 pInfo->cmdQID = cmdQID;
230 pInfo->msgQID = msgQID;
231 pInfo->sigQID = sigQID;
232 pInfo->reqQID = reqQID;
233
234 master->setProcessShmemSize(memorySize);
235
236// printf(" --------- Created new ProcessMap header: %u usage: %u serial %u --------\n", header->size, header->usage, serial);
237 mutex->leave();
238 return true;
239}
240
242 serial = master->getProcessShmemSerial();
243 uint64 size = master->getProcessShmemSize();
244 this->port = master->port;
245// printf(" --------- Opening new ProcessMemory section serial %u --------\n", serial);
246
247 bool createdMutex = false;
248 if (!mutex) {
249 //LogPrint(0,LOG_MEMORY,0," --------- Opening new ProcessMemory section serial %u --------", serial);
250 mutex = new utils::Mutex(utils::StringFormat("PsycloneProcessMemoryMutex_%u", port).c_str());
251 createdMutex = true;
252 if (!mutex->enter(5000, __FUNCTION__))
253 return false;
254 }
255 //else
256 // LogPrint(0,LOG_MEMORY,0," --------- Reopening ProcessMemory section serial %u --------", serial);
257
258 ProcessMemoryStruct* newHeader = (ProcessMemoryStruct*) utils::OpenSharedMemorySegment(utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), size);
259 if (!newHeader) {
260 LogPrint(0,LOG_MEMORY,0,"Cannot open shared memory: %s size: %u", utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), size);
261 if (createdMutex)
262 mutex->leave();
263 return false;
264 }
265 if (newHeader->cid != PROCESSMEMORYID) {
266 LogPrint(0,LOG_MEMORY,0,"Cannot verify shared memory: %s size: %u", utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), size);
267 utils::CloseSharedMemorySegment((char*)newHeader, size);
268 if (createdMutex)
269 mutex->leave();
270 return false;
271 }
272
273 memorySize = size;
274 if (header)
275 utils::CloseSharedMemorySegment((char*)header, header->size);
276 header = newHeader;
277 if (createdMutex)
278 mutex->leave();
279 // else leave mutex locked as we are calling from within the object
280 return true;
281}
282
283
284
285
286
287
288
289
290// static
291bool ProcessMemory::createNewProcess(const char* name, uint16 &id) {
292 if (!mutex || !mutex->enter(5000, __FUNCTION__))
293 return false;
295
296 if (!header->size) {
297 printf(" --------- Process header error: %llu usage: %llu serial %u --------\n", header->size, header->usage, serial);
298 return false;
299 }
300
301 // Find first empty offset, 0 is the base and always reserved
302 id = 0;
303 ProcessInfoStruct* pInfo = &(header->processes[1]);
304 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
305 if (!pInfo->createdTime) {
306 id = n;
307 break;
308 }
309 }
310 if (!id) {
311 mutex->leave();
312 return false;
313 }
314 pInfo->id = id;
315 pInfo->createdTime = pInfo->lastseen = GetTimeNow();
316 pInfo->status = PSYPROC_CREATED;
317 pInfo->nodeID = master->getNodeID();
318 utils::strcpyavail(pInfo->name, name, MAXKEYNAMELEN, true);
319
320 uint16 cmdQID;
321 uint16 msgQID;
322 uint16 sigQID;
323 uint16 reqQID;
324
325 if (
326 !setupNextAvailableQ(cmdQID) ||
327 !setupNextAvailableQ(msgQID) ||
328 !setupNextAvailableQ(sigQID) ||
329 !setupNextAvailableQ(reqQID) ) {
330 mutex->leave();
331 return false;
332 }
333
334 // reset as header might have changed...
335 pInfo = &(header->processes[id]);
336 pInfo->cmdQID = cmdQID;
337 pInfo->msgQID = msgQID;
338 pInfo->sigQID = sigQID;
339 pInfo->reqQID = reqQID;
340
341 mutex->leave();
342 return true;
343}
344
345// static
347 if (!mutex || !mutex->enter(5000, __FUNCTION__))
348 return false;
350 // delete queues
351 deleteQ(header->processes[id].cmdQID);
352 deleteQ(header->processes[id].msgQID);
353 deleteQ(header->processes[id].sigQID);
354 deleteQ(header->processes[id].reqQID);
355 memset(&(header->processes[id]), 0, sizeof(ProcessInfoStruct));
356 return true;
357}
358
359// static
360bool ProcessMemory::getProcessName(uint16 id, char* name, uint32 maxSize) {
361 if (!mutex || !mutex->enter(5000, __FUNCTION__))
362 return false;
364 if (!header->processes[id].createdTime) {
365 mutex->leave();
366 return false;
367 }
368 utils::strcpyavail(name, header->processes[id].name, maxSize, true);
369 mutex->leave();
370 return true;
371}
372
373// static
374bool ProcessMemory::getProcessID(const char* name, uint16 &id) {
375 if (!mutex || !mutex->enter(5000, __FUNCTION__))
376 return false;
378 ProcessInfoStruct* pInfo = &(header->processes[1]);
379 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
380 if ((stricmp(name, pInfo->name) == 0) && pInfo->createdTime) {
381 id = n;
382 mutex->leave();
383 return true;
384 }
385 }
386 mutex->leave();
387 return false;
388}
389
390// static
392 if (!mutex || !mutex->enter(5000, __FUNCTION__))
393 return 0;
395 uint64 t = header->processes[id].createdTime;
396 mutex->leave();
397 return t;
398}
399
400// static
401bool ProcessMemory::getProcessCommandLine(uint16 id, char* cmdline, uint32 maxSize) {
402 if (!mutex || !mutex->enter(5000, __FUNCTION__))
403 return false;
405 if (!header->processes[id].createdTime) {
406 mutex->leave();
407 return false;
408 }
409 utils::strcpyavail(cmdline, header->processes[id].commandline, maxSize, true);
410 mutex->leave();
411 return true;
412}
413
414// static
415uint8 ProcessMemory::getProcessStatus(uint16 id, uint64& lastseen) {
416 uint64 createTime;
417 return getProcessStatus(id, lastseen, createTime);
418}
419
420uint8 ProcessMemory::getProcessStatus(uint16 id, uint64& lastseen, uint64& createTime) {
421 if (!mutex || !mutex->enter(5000, __FUNCTION__))
422 return 0;
424 if (!header->processes[id].createdTime) {
425 mutex->leave();
426 lastseen = 0;
427 return 0;
428 }
429 uint8 status = header->processes[id].status;
430 lastseen = header->processes[id].lastseen;
431 createTime = header->processes[id].createdTime;
432 mutex->leave();
433 return status;
434}
435
436// static
438 if (!mutex || !mutex->enter(5000, __FUNCTION__))
439 return 0;
441 ProcessInfoStruct* pInfo = &(header->processes[1]);
442 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
443 if (pInfo->osID == osid) {
444 uint32 id = n;
445 mutex->leave();
446 return id;
447 }
448 }
449 mutex->leave();
450 return 0;
451}
452
453// static
455 if (!mutex || !mutex->enter(5000, __FUNCTION__))
456 return 0;
458 if (!header->processes[id].createdTime) {
459 mutex->leave();
460 return 0;
461 }
462 uint32 osID = header->processes[id].osID;
463 mutex->leave();
464 return osID;
465}
466
467// static
468bool ProcessMemory::setProcessOSID(uint16 id, uint32 osid) {
469 if (!mutex || !mutex->enter(5000, __FUNCTION__))
470 return false;
472 if (!header->processes[id].createdTime) {
473 mutex->leave();
474 return false;
475 }
476 header->processes[id].osID = osid;
477 mutex->leave();
478 return true;
479}
480
481// static
482bool ProcessMemory::setProcessType(uint16 id, uint8 type) {
483 if (!mutex || !mutex->enter(5000, __FUNCTION__))
484 return false;
486 if (!header->processes[id].createdTime) {
487 mutex->leave();
488 return false;
489 }
490 header->processes[id].type = type;
491 mutex->leave();
492 return true;
493}
494
495// static
496bool ProcessMemory::setProcessCommandLine(uint16 id, const char* cmdline) {
497 if (!mutex || !mutex->enter(5000, __FUNCTION__))
498 return false;
500 if (!header->processes[id].createdTime) {
501 mutex->leave();
502 return false;
503 }
504 utils::strcpyavail(header->processes[id].commandline, cmdline, MAXCOMMANDLINELEN, false);
505 mutex->leave();
506 return true;
507}
508
509// static
510bool ProcessMemory::setProcessStatus(uint16 id, uint8 status, uint64 currentCPUTicks) {
511 if (!mutex || !mutex->enter(5000, __FUNCTION__))
512 return false;
514 ProcessInfoStruct* pInfo = &header->processes[id];
515 if (!pInfo->createdTime) {
516 mutex->leave();
517 return false;
518 }
519 pInfo->status = status;
520 pInfo->lastseen = GetTimeNow();
521 pInfo->stats.procMemUsage = utils::GetProcessMemoryUsage();
522 if (currentCPUTicks)
523 pInfo->stats.currentCPUTicks = currentCPUTicks;
524 mutex->leave();
525 return true;
526}
527
529 AveragePerfStats perfStats;
530 memset(&perfStats, 0, sizeof(AveragePerfStats));
531
532 if (!mutex || !mutex->enter(5000, __FUNCTION__))
533 return perfStats;
534 if (serial != master->getProcessShmemSerial()) {if (!open()) {mutex->leave();return perfStats;}}
535
536 ProcessInfoStruct* pInfo = &header->processes[procID];
537 if (!pInfo->createdTime) {
538 mutex->leave();
539 return perfStats;
540 }
541 memcpy(&perfStats, &pInfo->perfStats, sizeof(AveragePerfStats));
542 mutex->leave();
543 return perfStats;
544}
545
546bool ProcessMemory::setProcessPerfStats(uint16 procID, AveragePerfStats &perfStruct) {
547 if (!mutex || !mutex->enter(5000, __FUNCTION__))
548 return false;
550 ProcessInfoStruct* pInfo = &header->processes[procID];
551 if (!pInfo->createdTime) {
552 mutex->leave();
553 return false;
554 }
555 memcpy(&pInfo->perfStats, &perfStruct, sizeof(AveragePerfStats));
556 mutex->leave();
557 return true;
558}
559
560
561bool ProcessMemory::addToProcessStats(uint16 id, DataMessage* inputMsg, DataMessage* outputMsg) {
562 if (!mutex || !mutex->enter(5000, __FUNCTION__))
563 return false;
565 ProcessInfoStruct* pInfo = &header->processes[id];
566 if (!pInfo->createdTime) {
567 mutex->leave();
568 return false;
569 }
570
571 pInfo->stats.time = GetTimeNow();
572 DataMessage* draft;
573 if (inputMsg) {
574 pInfo->stats.msgInCount++;
575 pInfo->stats.msgInBytes += inputMsg->getSize();
576 // Circular ring: write the newest draft to the head slot and advance, instead
577 // of memmove-shifting the whole 10-deep ring (9*DRAFTMSGSIZE = 36KB) per message.
578 // Readers (toXML, PsyProbe) iterate all slots, so slot order doesn't matter.
579 char* inSlot = pInfo->stats.recentInMsg + ((pInfo->stats.recentInPos % 10) * DRAFTMSGSIZE);
580 if (inputMsg->data->size < DRAFTMSGSIZE)
581 memcpy(inSlot, inputMsg->data, inputMsg->data->size);
582 else {
583 draft = new DataMessage(*inputMsg, DRAFTMSGSIZE);
584 memcpy(inSlot, draft->data, draft->data->size);
585 delete(draft);
586 }
587 pInfo->stats.recentInPos = (pInfo->stats.recentInPos + 1) % 10;
588 }
589 if (outputMsg) {
590 pInfo->stats.msgOutCount++;
591 pInfo->stats.msgOutBytes += outputMsg->getSize();
592 char* outSlot = pInfo->stats.recentOutMsg + ((pInfo->stats.recentOutPos % 10) * DRAFTMSGSIZE);
593 if (outputMsg->data->size < DRAFTMSGSIZE)
594 memcpy(outSlot, outputMsg->data, outputMsg->data->size);
595 else {
596 draft = new DataMessage(*outputMsg, DRAFTMSGSIZE);
597 memcpy(outSlot, draft->data, draft->data->size);
598 delete(draft);
599 }
600 pInfo->stats.recentOutPos = (pInfo->stats.recentOutPos + 1) % 10;
601 }
602
603 mutex->leave();
604 return true;
605}
606
607bool ProcessMemory::addToCmdQ(uint16 procID, DataMessage* msg) {
608 if (!mutex || !mutex->enter(5000, __FUNCTION__))
609 return false;
611
612 if (!addToQ(procID, CMDQ_TYPE, msg)) {
613 mutex->leave();
614 return false;
615 }
616 mutex->leave();
617 return true;
618}
619
620bool ProcessMemory::addToMsgQ(uint16 procID, DataMessage* msg) {
621 if (!mutex || !mutex->enter(5000, __FUNCTION__))
622 return false;
624
625 if (!addToQ(procID, MSGQ_TYPE, msg)) {
626 mutex->leave();
627 return false;
628 }
629 mutex->leave();
630 return true;
631}
632
633bool ProcessMemory::addToSigQ(uint16 procID, DataMessage* msg) {
634 if (!mutex || !mutex->enter(5000, __FUNCTION__))
635 return false;
637
638 if (!addToQ(procID, SIGQ_TYPE, msg)) {
639 mutex->leave();
640 return false;
641 }
642 mutex->leave();
643 return true;
644}
645
647 if (!mutex || !mutex->enter(5000, __FUNCTION__))
648 return false;
650
651 ProcessInfoStruct* pInfo = &(header->processes[0]);
652 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
653 if (pInfo->createdTime)
654 addToQ(pInfo->id, SIGQ_TYPE, msg);
655 }
656
657 mutex->leave();
658 return true;
659}
660
662 if (!mutex || !mutex->enter(5000, __FUNCTION__))
663 return false;
665
666 ProcessInfoStruct* pInfo = &(header->processes[0]);
667 for (uint32 n=1; n<MAXPROC; n++, pInfo++) {
668 if (pInfo->createdTime && (pInfo->id != except))
669 addToQ(pInfo->id, SIGQ_TYPE, msg);
670 }
671
672 mutex->leave();
673 return true;
674}
675
676bool ProcessMemory::addToReqQ(uint16 procID, DataMessage* msg) {
677 if (!mutex || !mutex->enter(5000, __FUNCTION__))
678 return false;
680
681 if (!addToQ(procID, REQQ_TYPE, msg)) {
682 mutex->leave();
683 return false;
684 }
685 mutex->leave();
686 return true;
687}
688
689
690DataMessage* ProcessMemory::waitForCmdQ(uint16 procID, uint32 timeout) {
691 if (!mutex || !mutex->enter(5000, __FUNCTION__))
692 return NULL;
694 DataMessage* msg = waitForQ(procID, CMDQ_TYPE, timeout);
695 mutex->leave();
696 return msg;
697}
698
699DataMessage* ProcessMemory::waitForMsgQ(uint16 procID, uint32 timeout) {
700 if (!mutex || !mutex->enter(5000, __FUNCTION__))
701 return NULL;
703 DataMessage* msg = waitForQ(procID, MSGQ_TYPE, timeout);
704 mutex->leave();
705 return msg;
706}
707
708DataMessage* ProcessMemory::waitForSigQ(uint16 procID, uint32 timeout) {
709 if (!mutex || !mutex->enter(5000, __FUNCTION__))
710 return NULL;
712 DataMessage* msg = waitForQ(procID, SIGQ_TYPE, timeout);
713 mutex->leave();
714 return msg;
715}
716
717DataMessage* ProcessMemory::waitForReqQ(uint16 procID, uint32 timeout) {
718 if (!mutex || !mutex->enter(5000, __FUNCTION__))
719 return NULL;
721 DataMessage* msg = waitForQ(procID, REQQ_TYPE, timeout);
722 mutex->leave();
723 return msg;
724}
725
726uint32 ProcessMemory::getCmdQCount(uint16 procID) {
727 if (!mutex || !mutex->enter(5000, __FUNCTION__))
728 return 0;
730
731 MessageQueueHeader* qHeader = getQHeader(procID, CMDQ_TYPE);
732 if (!qHeader) {
733 mutex->leave();
734 return 0;
735 }
736 uint32 count = qHeader->count;
737 mutex->leave();
738 return count;
739}
740
741uint32 ProcessMemory::getMsgQCount(uint16 procID) {
742 if (!mutex || !mutex->enter(5000, __FUNCTION__))
743 return 0;
745
746 MessageQueueHeader* qHeader = getQHeader(procID, MSGQ_TYPE);
747 if (!qHeader) {
748 mutex->leave();
749 return 0;
750 }
751 uint32 count = qHeader->count;
752 mutex->leave();
753 return count;
754}
755
756uint32 ProcessMemory::getSigQCount(uint16 procID) {
757 if (!mutex || !mutex->enter(5000, __FUNCTION__))
758 return 0;
760
761 MessageQueueHeader* qHeader = getQHeader(procID, SIGQ_TYPE);
762 if (!qHeader) {
763 mutex->leave();
764 return 0;
765 }
766 uint32 count = qHeader->count;
767 mutex->leave();
768 return count;
769}
770
771uint32 ProcessMemory::getReqQCount(uint16 procID) {
772 if (!mutex || !mutex->enter(5000, __FUNCTION__))
773 return 0;
775
776 MessageQueueHeader* qHeader = getQHeader(procID, REQQ_TYPE);
777 if (!qHeader) {
778 mutex->leave();
779 return 0;
780 }
781 uint32 count = qHeader->count;
782 mutex->leave();
783 return count;
784}
785
786MessageQueueHeader* ProcessMemory::getQHeader(uint16 procID, uint8 qType) {
787 if (!header->processes[procID].createdTime)
788 return NULL;
789 uint16 qID;
790 switch(qType) {
791 case CMDQ_TYPE:
792 qID = header->processes[procID].cmdQID;
793 break;
794 case MSGQ_TYPE:
795 qID = header->processes[procID].msgQID;
796 break;
797 case SIGQ_TYPE:
798 qID = header->processes[procID].sigQID;
799 break;
800 case REQQ_TYPE:
801 qID = header->processes[procID].reqQID;
802 break;
803 default:
804 return NULL;
805 }
806 if (!header->qIndex[qID])
807 return NULL;
808
809 MessageQueueHeader* qHeader = (MessageQueueHeader*) (((char*)header) + header->qIndex[qID]);
810 if (qHeader->id != qID)
811 return NULL;
812 return qHeader;
813}
814
815
816
817bool ProcessMemory::setupNextAvailableQ(uint16& qID) {
818 qID = 0;
819 uint16 nextID = 0;
820 bool found = false;
821 for (uint32 n=0; n<MAXPROC4; n++) {
822 if (!found && !header->qIndex[n]) {
823 qID = n;
824 found = true;
825 }
826 else if (found && header->qIndex[n]) {
827 nextID = n;
828 break;
829 }
830 }
831 if (!found)
832 return false;
833
834 // Insert a new standard size queue in place
835 uint64 size = INITIALQSIZE;
836 // Do we have enough space free at the end?
837 if (header->size - header->usage < size) {
838 LogPrint(0,LOG_MEMORY,0,"SetupQ need more memory size: %u use: %u need: %u", header->size, header->usage, size);
839 // Resize whole memory block
840 if (!resize(memorySize * 2)) {
841 LogPrint(0,LOG_MEMORY,0,"Cannot resize ProcessMemory");
842 return false;
843 }
844 }
845
846 char* data;
847 // If any active queues follow the new queue
848 if (nextID) {
849 // Push every subsequent queue by size
850 header->qIndex[qID] = header->qIndex[nextID];
851 data = ((char*)header) + header->qIndex[nextID];
852 uint64 sizeMove = header->usage - header->qIndex[nextID];
853 memmove(data+size, data, (size_t)sizeMove);
854 // Adjust indexes
855 for (uint32 n=nextID; n<MAXPROC4; n++) {
856 if (header->qIndex[n])
857 header->qIndex[n] += size;
858 }
859 }
860 else {
861 header->qIndex[qID] = header->usage;
862 data = ((char*)header) + header->usage;
863 }
864
865 // We now have space, a start point and the indexes filled in
866 memset(data, 0, (size_t)size); // could just be the header
867 ((MessageQueueHeader*)data)->size = size;
868 ((MessageQueueHeader*)data)->id = qID;
869 snprintf(((MessageQueueHeader*)data)->name, MAXKEYNAMELEN, "PsycloneProcessMemoryQueue_%u_%u", port, qID);
870// snprintf(((MessageQueueHeader*)data)->name, MAXKEYNAMELEN, "%u_%u_PsycloneProcessMemoryQueue", port, qID);
871
872 header->usage += size;
873 return true;
874}
875
876bool ProcessMemory::deleteQ(uint16 id) {
877 if (!header->qIndex[id])
878 return false;
879 MessageQueueHeader* qHeader = (MessageQueueHeader*)(((char*)header) + header->qIndex[id]);
880 uint64 size = qHeader->size;
881
882 // Check if any data comes after this queue
883 uint16 nextID = 0;
884 if (id < MAXPROC4-2) {
885 for (uint32 n=id+1; n<MAXPROC4; n++) {
886 if (header->qIndex[n]) {
887 if (!nextID) {
888 nextID = n;
889 // move memory
890 memmove(qHeader, ((char*)qHeader)+size, (size_t)(header->usage - header->qIndex[nextID]));
891 }
892 header->qIndex[n] -= size;
893 }
894 }
895 }
896
897 header->qIndex[id] = 0;
898 header->usage -= size;
899 return true;
900}
901
902
903bool ProcessMemory::addToQ(uint16 procID, uint8 qType, DataMessage* msg) {
904 MessageQueueHeader* qHeader = getQHeader(procID, qType);
905 if (!qHeader)
906 return false;
907
908 uint32 msgSize = msg->getSize();
909 uint64 spaceToEnd, inUse = 0, inUse2 = 0;
910 char* qData = ((char*)qHeader) + sizeof(MessageQueueHeader);
911 uint64 qDataSize = qHeader->size - sizeof(MessageQueueHeader);
912
913 //if (qHeader->endPos >= qHeader->startPos)
914 // inUse = qHeader->endPos - qHeader->startPos;
915 //else
916 // inUse = qDataSize - qHeader->startPos - qHeader->padding + qHeader->endPos;
917
918 if (qHeader->endPos == qHeader->startPos) {
919 spaceToEnd = qDataSize;
920 if (spaceToEnd >= msgSize) {
921 // add msg at buffer start
922 memcpy(qData, msg->data, msgSize);
923 qHeader->startPos = 0;
924 qHeader->endPos = msgSize;
925 qHeader->count++;
926 qHeader->padding = 0;
927 }
928 else {
929 // we need to resize queue
930 if (!resizeQ(qHeader->id, qHeader->size*2))
931 return false;
932 // call this function again...
933 return addToQ(procID, qType, msg);
934 }
935 }
936 else if (qHeader->endPos > qHeader->startPos) {
937 spaceToEnd = qDataSize - qHeader->endPos;
938 if (spaceToEnd >= msgSize) {
939 // add msg here
940 memcpy(qData+qHeader->endPos, msg->data, msgSize);
941 qHeader->endPos += msgSize;
942 qHeader->count++;
943 qHeader->padding = 0;
944 }
945 else if (qHeader->startPos > msgSize) {
946 // fill end rest of space with 0
947 memset(qData+qHeader->endPos, 0, (size_t)spaceToEnd);
948 qHeader->padding = spaceToEnd;
949 // add msg at buffer start
950 memcpy(qData, msg->data, msgSize);
951 qHeader->endPos = msgSize;
952 qHeader->count++;
953 }
954 else {
955 // we need to resize queue
956 if (!resizeQ(qHeader->id, qHeader->size*2))
957 return false;
958 // call this function again...
959 return addToQ(procID, qType, msg);
960 }
961 }
962 else {
963 spaceToEnd = qHeader->startPos - qHeader->endPos;
964 if (spaceToEnd > msgSize) {
965 // add msg here
966 memcpy(qData+qHeader->endPos, msg->data, msgSize);
967 qHeader->endPos += msgSize;
968 qHeader->count++;
969 }
970 else {
971 // we need to resize queue
972 if (!resizeQ(qHeader->id, qHeader->size*2))
973 return false;
974 // call this function again...
975 return addToQ(procID, qType, msg);
976 }
977 }
978
979 //if (qHeader->endPos >= qHeader->startPos)
980 // inUse2 = qHeader->endPos - qHeader->startPos;
981 //else
982 // inUse2 = qDataSize - qHeader->startPos - qHeader->padding + qHeader->endPos;
983 //if (inUse2 < (qHeader->count * msgSize))
984 // int ee = 0;
985
986 utils::Semaphore** sem = &qSemaphores[qHeader->id];
987 if (!*sem)
988 if (!(*sem = utils::GetSemaphore(qHeader->name)))
989 return false;
990 (*sem)->signal();
991
992// utils::SignalSemaphore(qHeader->name);
993// printf("Signal %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(GetTimeNow(), true).c_str());
994 return true;
995}
996
997DataMessage* ProcessMemory::waitForQ(uint16 procID, uint8 qType, uint32 timeout) {
998 //char* name = NULL;
999 MessageQueueHeader* qHeader = getQHeader(procID, qType);
1000 if (!qHeader)
1001 return NULL;
1002
1003 uint64 start = GetTimeNow();
1004 int32 remain;
1005 while (!qHeader->count) {
1006 utils::Semaphore** sem = &qSemaphores[qHeader->id];
1007 if (!*sem) {
1008 if (!(*sem = utils::GetSemaphore(qHeader->name))) {
1009 return NULL;
1010 }
1011 }
1012
1013 //if (!name) {
1014 // name = new char[strlen(qHeader->name)+1];
1015 // utils::strcpyavail(name, qHeader->name, (uint32)strlen(qHeader->name)+1, true);
1016 //}
1017 if ((remain = timeout - GetTimeAgeMS(start)) <= 0) {
1018 // printf("StartWait1 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(start, true).c_str());
1019 // printf("Timeout1 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(GetTimeNow(), true).c_str());
1020 // delete [] name;
1021 return NULL;
1022 }
1023 mutex->leave();
1024 bool result = (*sem)->wait(remain);
1025 if (!mutex || !mutex->enter(5000, __FUNCTION__))
1026 return NULL;
1028
1029 if (!result) {
1030 // printf("StartWait2 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(start, true).c_str());
1031 // printf("Timeout2 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(GetTimeNow(), true).c_str());
1032 // delete [] name;
1033 // printf("Waited %.3fms (%ums)...\n", GetTimeAge(start)/1000.0, remain);
1034 return NULL;
1035 }
1036
1037 // header might be changed
1038 qHeader = getQHeader(procID, qType);
1039 if (!qHeader) {
1040 // mutex->leave();
1041 // delete [] name;
1042 return NULL;
1043 }
1044 }
1045 //delete [] name;
1046
1047 DataMessage* msg = NULL;
1048 if (qHeader->count) {
1049 msg = new DataMessage(((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos, true);
1050 qHeader->startPos += msg->getSize();
1051 if (qHeader->startPos >= qHeader->size - sizeof(MessageQueueHeader) - qHeader->padding) {
1052 qHeader->startPos = 0;
1053 qHeader->padding = 0;
1054 }
1055 qHeader->count--;
1056//if (qHeader->count && (GetObjID(((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos) != DATAMESSAGEID)) {
1057// printQ(qHeader);
1058// return false;
1059//}
1060 }
1061
1062 return msg;
1063}
1064
1065bool ProcessMemory::printQ(MessageQueueHeader* qHeader) {
1066 printf("[%u] Data Size: %llu, count: %u startpos: %llu, endpos: %llu, padding: %llu\n",
1067 qHeader->id, (uint64)(qHeader->size - sizeof(MessageQueueHeader)), qHeader->count, qHeader->startPos, qHeader->endPos, qHeader->padding);
1068 char* data = ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos;
1069 for (uint32 n=0; n<qHeader->count; n++) {
1070 if (GetObjID(data) == DATAMESSAGEID)
1071 printf(" [%u: %llu - %llu] Msg size: %u Time: %s Type: %u.%u.%u\n", n,
1072 (uint64)(data - ((char*)qHeader) - sizeof(MessageQueueHeader)),
1073 (uint64)(data - ((char*)qHeader) - sizeof(MessageQueueHeader) + ((DataMessageHeader*)data)->size),
1074 ((DataMessageHeader*)data)->size, PrintTimeString(((DataMessageHeader*)data)->time).c_str(),
1075 (((DataMessageHeader*)data)->type).levels[0],
1076 (((DataMessageHeader*)data)->type).levels[1],
1077 (((DataMessageHeader*)data)->type).levels[2]);
1078 else {
1079 printf(" [%u: %llu] Msg corrupted\n",
1080 n, (uint64)(data - ((char*)qHeader) - sizeof(MessageQueueHeader)));
1081 return false;
1082 }
1083 data += ((DataMessageHeader*)data)->size;
1084 if (data == ((char*)qHeader) + qHeader->size - qHeader->padding)
1085 data = ((char*)qHeader) + sizeof(MessageQueueHeader);
1086 }
1087 return true;
1088}
1089
1090bool ProcessMemory::resize(uint64 newMemorySize) {
1091
1092 if (newMemorySize > 100000000) {
1093 //LogPrint(0,LOG_MEMORY,0,"Memory Error: Process Memory requested size %s denied...", utils::BytifySize((double)newMemorySize).c_str());
1094 printf("Memory Error: Process Memory requested size %s denied...\n", utils::BytifySize((double)newMemorySize).c_str());
1095 fflush(stdout);
1096 return false;
1097 }
1098
1099 serial = master->incrementProcessShmemSerial();
1100 //LogPrint(0,LOG_MEMORY,0," --------- Resizing ProcessMemory section %u -> %u serial %u --------", memorySize, newMemorySize, serial);
1101
1102 if (!mutex || !mutex->enter(5000, __FUNCTION__))
1103 return false;
1104
1105 //ProcessMemoryStruct* newHeader = (ProcessMemoryStruct*) utils::OpenSharedMemorySegment(utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), newMemorySize);
1106 //if (newHeader) {
1107 // utils::CloseSharedMemorySegment((char*)newHeader, newMemorySize);
1108 // //LogPrint(0,LOG_MEMORY,2,"ProcessMemory removing stale shared memory (%u/%u)...", port, serial);
1109 //}
1110
1111 ProcessMemoryStruct* newHeader = (ProcessMemoryStruct*) utils::CreateSharedMemorySegment(utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), newMemorySize, true);
1112 if (!newHeader) {
1113 mutex->leave();
1114 return false;
1115 }
1116 memcpy(newHeader, header, (size_t)header->size);
1117 newHeader->size = newMemorySize;
1118
1119 utils::CloseSharedMemorySegment((char*)header, memorySize);
1120 header = newHeader;
1121 memorySize = newMemorySize;
1122 master->setProcessShmemSize(memorySize);
1123 mutex->leave();
1124
1125 LogPrint(0,LOG_MEMORY,2,"Process Memory resized to %s...", utils::BytifySize((double)newMemorySize).c_str());
1126 if (newMemorySize > 20000000) {
1127 LogPrint(0,LOG_MEMORY,2,"Memory Warning: Process Memory now %s...", utils::BytifySize((double)newMemorySize).c_str());
1128 }
1129
1130 return true;
1131}
1132
1133bool ProcessMemory::resizeQ(uint16 id, uint64 newSize) {
1134 if (!header->qIndex[id])
1135 return false;
1136 MessageQueueHeader* qHeader = (MessageQueueHeader*) (((char*)header) + header->qIndex[id]);
1137 if (qHeader->id != id)
1138 return false;
1139
1140 uint64 difSize = newSize - qHeader->size;
1141 if (header->size - header->usage < difSize) {
1142 if (!resize(memorySize * 2))
1143 return false;
1144 return resizeQ(id, newSize);
1145 }
1146
1147 uint16 nextID = 0;
1148 uint32 n;
1149 for (n=id+1; n<MAXPROC4; n++) {
1150 if (header->qIndex[n]) {
1151 nextID = n;
1152 break;
1153 }
1154 }
1155
1156 // shift qID > id down by difsize
1157 if (nextID) {
1158 // Push every subsequent queue by size
1159 char* data = ((char*)header) + header->qIndex[nextID];
1160 uint64 sizeMove = header->usage - header->qIndex[nextID];
1161 memmove(data+difSize, data, (size_t)sizeMove);
1162 // Adjust indexes
1163 for (n=nextID; n<MAXPROC4; n++) {
1164 if (header->qIndex[n])
1165 header->qIndex[n] += difSize;
1166 }
1167 }
1168
1169 // Now rearrange messages in queue
1170 int64 size = qHeader->endPos - qHeader->startPos;
1171 uint64 firstSize = 0;
1172 uint64 secondSize = 0;
1173 if (size == 0)
1174 qHeader->endPos = qHeader->startPos = 0;
1175 else if ((size > 0) && qHeader->startPos) {
1176 memmove(
1177 ((char*)qHeader) + sizeof(MessageQueueHeader),
1178 ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos,
1179 (size_t)size);
1180 qHeader->endPos = (uint64)size;
1181 qHeader->startPos = 0;
1182 }
1183 else if (size < 0) {
1184 firstSize = qHeader->size - sizeof(MessageQueueHeader) - qHeader->startPos - qHeader->padding;
1185 secondSize = qHeader->endPos;
1186
1187 // is there enough space to move things around without a safety copy involving a malloc?
1188 if (difSize > secondSize) {
1189 // move first part of memory to the newly available space at the end
1190 memmove(
1191 ((char*)qHeader) + qHeader->size,
1192 ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos,
1193 (size_t)firstSize);
1194 // move last part of memory to the right place
1195 memmove(
1196 ((char*)qHeader) + sizeof(MessageQueueHeader) + firstSize,
1197 ((char*)qHeader) + sizeof(MessageQueueHeader),
1198 (size_t)secondSize);
1199 // copy first part back
1200 memmove(
1201 ((char*)qHeader) + sizeof(MessageQueueHeader),
1202 ((char*)qHeader) + qHeader->size,
1203 (size_t)firstSize);
1204 }
1205 else {
1206 // make safety copy of first part
1207 char* copy = (char*)malloc((size_t)firstSize);
1208 memcpy(copy, ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos, (size_t)firstSize);
1209 // move last part of memory to the right place
1210 memmove(
1211 ((char*)qHeader) + sizeof(MessageQueueHeader) + firstSize,
1212 ((char*)qHeader) + sizeof(MessageQueueHeader),
1213 (size_t)secondSize);
1214 // copy first part back from the safety copy
1215 memcpy(((char*)qHeader) + sizeof(MessageQueueHeader), copy, (size_t)firstSize);
1216 free(copy);
1217 }
1218 qHeader->endPos = firstSize + secondSize;
1219 qHeader->startPos = 0;
1220 }
1221
1222 qHeader->size = newSize;
1223 qHeader->padding = 0;
1224 header->usage += difSize;
1225
1226//if (qHeader->count) {
1227// printQ(qHeader);
1228// return false;
1229//}
1230
1231 return true;
1232}
1233
1234
1236 unittest::progress(0, "init memory manager");
1237
1238 // First create and initialise the MemoryManager
1239 MemoryManager* manager = new MemoryManager();
1240 if (!manager->create(0)) {
1241 unittest::fail("MemoryManager create(0) failed");
1243 delete(manager);
1244 return false;
1245 }
1246
1247 uint32 count = 10000;
1248 uint32 writeCount = 0;
1249 uint32 readCount = 0;
1250 uint32 n;
1251 DataMessage* msg;
1252 char str[128];
1253 const char* cstr;
1254
1255 uint32 queueTestThreadID;
1256 g_processTestStop = false;
1257 if (!ThreadManager::CreateThread(QueueTest, manager, queueTestThreadID)) {
1258 unittest::fail("Could not create Queue Test thread");
1260 delete(manager);
1261 return false;
1262 }
1263 g_processTestThreadID = queueTestThreadID;
1264
1265 uint64 time = GetTimeNow();
1266
1267 // empty-queue wait must actually block for ~1s
1268 unittest::progress(10, "empty queue wait");
1269 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) != NULL) {
1270 unittest::fail("Got unexpected Message from empty Queue");
1272 delete(manager);
1273 return false;
1274 }
1275
1276 uint64 t = 0;
1277 if ( (t = GetTimeAge(time)) < 950000) {
1278 unittest::fail("Message Queue wait 1000ms didn't actually wait more than %.3fms", (t/1000.0));
1280 delete(manager);
1281 return false;
1282 }
1283
1284 // Single-threaded round-trip latency (same thread adds + reads sig queue)
1285 unittest::progress(25, "single-thread latency");
1286 const uint32 latIters = 20000;
1287 uint64 singleTotal = 0;
1288 t = 0;
1289 for (n=0; n<latIters; n++) {
1290 msg = new DataMessage(CTRL_TEST, 0);
1291 if (!manager->processMemory->addToSigQ(0, msg)) {
1292 unittest::fail("Could not add Message %u to SigQueue", n);
1293 delete(manager);
1294 return false;
1295 }
1296 delete(msg);
1297 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) == NULL) {
1298 unittest::fail("Could not get Message [%u] from SigQueue", n);
1299 delete(manager);
1300 return false;
1301 }
1302 t += GetTimeAge(msg->getCreatedTime());
1303 delete(msg);
1304 if (n && n % 10000 == 0) {
1305 unittest::detail("Single threaded message delay[%u]: %.3fus", n, ((double)t)/(2*10000.0));
1306 singleTotal += t;
1307 t = 0;
1308 }
1309 }
1310 singleTotal += t;
1311 // each iteration measures the created-time age twice (add + read)
1312 double singleLatency = (double)singleTotal / (2.0 * latIters);
1313 unittest::metric("single_thread_latency", singleLatency, "us", false);
1314
1315 // warm up the queue (cross-thread path via cmd queue -> slave -> sig queue)
1316 msg = new DataMessage(CTRL_TEST, 0);
1317 if (!manager->processMemory->addToCmdQ(0, msg)) {
1318 unittest::fail("Could not add startup Message to Queue");
1320 delete(manager);
1321 return false;
1322 }
1323 delete(msg);
1324 utils::Sleep(10);
1325 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) == NULL) {
1326 unittest::fail("Could not get startup Message from Queue");
1328 delete(manager);
1329 return false;
1330 }
1331 delete(msg);
1332
1333 // Multi-threaded round-trip latency (slave thread relays cmd -> sig)
1334 unittest::progress(45, "multi-thread latency");
1335 uint64 multiTotal = 0;
1336 t = 0;
1337 for (n=0; n<latIters; n++) {
1338 msg = new DataMessage(CTRL_TEST, 0);
1339 if (!manager->processMemory->addToCmdQ(0, msg)) {
1340 unittest::fail("Could not add Message %u to CmdQueue", n);
1341 delete(manager);
1342 return false;
1343 }
1344 delete(msg);
1345 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) == NULL) {
1346 unittest::fail("Could not get Message [%u] from SigQueue", n);
1347 delete(manager);
1348 return false;
1349 }
1350 t += GetTimeAge(msg->getCreatedTime());
1351 delete(msg);
1352 if (n && n % 10000 == 0) {
1353 unittest::detail("Multi threaded message delay[%u]: %.3fus", n, ((double)t)/(2*10000.0));
1354 multiTotal += t;
1355 t = 0;
1356 }
1357 }
1358 multiTotal += t;
1359 double multiLatency = (double)multiTotal / (2.0 * latIters);
1360 unittest::metric("multi_thread_latency", multiLatency, "us", false);
1361
1362
1363
1364
1365
1366
1367 // Message queue correctness + write throughput
1368 unittest::progress(65, "msg queue fill");
1369 time = GetTimeNow();
1370 uint64 writeStart = GetTimeNow();
1371 for (n=0; n<count/4; n++) {
1372 msg = new DataMessage(CTRL_TEST, writeCount);
1373 snprintf(str, 128, "Test%2u", writeCount);
1374 msg->setString("TestEntry", str);
1375 msg->setTime("TestTime", time);
1376 msg->setString("TestEntry2", str);
1377 msg->setTime("TestTime2", time);
1378 if (!manager->processMemory->addToMsgQ(0, msg)) {
1379 unittest::fail("Could not add Message %u to MsgQueue", n);
1380 delete(manager);
1381 return false;
1382 }
1383 // if (n == 1000)
1384 // utils::PrintBinary(msg->data, msg->data->size, false, "Msg1000");
1385 delete(msg);
1386 writeCount++;
1387 }
1388
1389 if ( (n = manager->processMemory->getMsgQCount(0)) != count/4) {
1390 unittest::fail("Queue contains %u messages instead of %u", n, count/4);
1392 delete(manager);
1393 return false;
1394 }
1395
1396 unittest::progress(75, "msg queue drain");
1397 for (n=0; n<count/8; n++) {
1398 if ( (msg = manager->processMemory->waitForMsgQ(0, 1000)) == NULL) {
1399 unittest::fail("Could not get Message %u from Queue", n);
1400 delete(manager);
1401 return false;
1402 }
1403 snprintf(str, 128, "Test%2u", readCount);
1404 if (strcmp(str, msg->getString("TestEntry")) != 0) {
1405 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, msg->getString("TestEntry"));
1406 delete(msg);
1407 delete(manager);
1408 return false;
1409 }
1410 if (msg->getTime("TestTime") != time) {
1411 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime"));
1412 delete(msg);
1413 delete(manager);
1414 return false;
1415 }
1416 cstr = msg->getString("TestEntry2");
1417 if (!cstr || strcmp(str, cstr) != 0) {
1418 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, cstr);
1419 delete(msg);
1420 delete(manager);
1421 return false;
1422 }
1423 if (msg->getTime("TestTime2") != time) {
1424 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime2"));
1425 delete(msg);
1426 delete(manager);
1427 return false;
1428 }
1429 readCount++;
1430 delete(msg);
1431 }
1432
1433 if ( (n = manager->processMemory->getMsgQCount(0)) != count/8) {
1434 unittest::fail("Queue contains %u messages instead of %u", n, count/8);
1436 delete(manager);
1437 return false;
1438 }
1439
1440 unittest::progress(85, "msg queue refill");
1441 for (n=0; n<3*count/4; n++) {
1442 msg = new DataMessage(CTRL_TEST, writeCount);
1443 snprintf(str, 128, "Test%2u", writeCount);
1444 msg->setString("TestEntry", str);
1445 msg->setTime("TestTime", time);
1446 msg->setString("TestEntry2", str);
1447 msg->setTime("TestTime2", time);
1448 if (!manager->processMemory->addToMsgQ(0, msg)) {
1449 unittest::fail("Could not add Message %u to MsgQueue", n);
1450 delete(manager);
1451 return false;
1452 }
1453 delete(msg);
1454 writeCount++;
1455 }
1456
1457 if ( (n = manager->processMemory->getMsgQCount(0)) != 7*count/8) {
1458 unittest::fail("Queue contains %u messages instead of %u", n, 7*count/8);
1460 delete(manager);
1461 return false;
1462 }
1463
1464 unittest::progress(92, "msg queue verify");
1465 for (n=0; n<7*count/8; n++) {
1466 if ( (msg = manager->processMemory->waitForMsgQ(0, 1000)) == NULL) {
1467 unittest::fail("Could not get Message %u from Queue", n);
1468 delete(manager);
1469 return false;
1470 }
1471 snprintf(str, 128, "Test%2u", readCount);
1472 if (strcmp(str, msg->getString("TestEntry")) != 0) {
1473 unittest::fail("Message %u from Queue corrupted string '%s' (queue has %u)", readCount, msg->getString("TestEntry"), manager->processMemory->getMsgQCount(0));
1474 delete(msg);
1475 delete(manager);
1476 return false;
1477 }
1478 if (msg->getTime("TestTime") != time) {
1479 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime"));
1480 delete(msg);
1481 delete(manager);
1482 return false;
1483 }
1484 cstr = msg->getString("TestEntry2");
1485 if (!cstr || strcmp(str, cstr) != 0) {
1486 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, cstr);
1487 delete(msg);
1488 delete(manager);
1489 return false;
1490 }
1491 if (msg->getTime("TestTime2") != time) {
1492 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime2"));
1493 delete(msg);
1494 delete(manager);
1495 return false;
1496 }
1497 readCount++;
1498 delete(msg);
1499 }
1500
1501 if ( (n = manager->processMemory->getMsgQCount(0)) != 0) {
1502 unittest::fail("Queue contains %u messages instead of 0", n);
1504 delete(manager);
1505 return false;
1506 }
1507
1508 // Aggregate write throughput over all enqueued messages (count/4 + 3*count/4)
1509 double writeUs = (double)(GetTimeNow() - writeStart);
1510 if (writeUs > 0.0)
1511 unittest::metric("msg_queue_throughput", (double)count / writeUs * 1e6, "msg/s", true);
1512
1513
1514// printf("Testing Request Queue...\n\n");
1515//
1516// // First create two components
1517// uint16 procID1, procID2;
1518// uint32 compID1 = 1, compID2 = 2;
1519// uint32 qID;
1520// //if (!MemoryMaps::CreateNewNode(1, "", 1)) {
1521// // printf("Error: Testing Request Queue - CreateNewNode\n");
1522// // return false;
1523// //}
1524// if (!MemoryMaps::CreateNewProcess("Proc1", procID1)) {
1525// printf("Error: Testing Request Queue - CreateNewProcess\n");
1526// return false;
1527// }
1528// if (!MemoryMaps::CreateNewProcess("Proc2", procID2)) {
1529// printf("Error: Testing Request Queue - CreateNewProcess 2\n");
1530// return false;
1531// }
1532// if (!MemoryQueues::CreateMessageQueue(qID)) {
1533// printf("Error: Testing Request Queue - CreateMessageQueue\n");
1534// return false;
1535// }
1536// if (!MemoryMaps::SetProcessQueueID(procID1, MSGQ_ID, qID)) {
1537// printf("Error: Testing Request Queue - SetProcessQueueID\n");
1538// return false;
1539// }
1540// if (!MemoryQueues::CreateMessageQueue(qID)) {
1541// printf("Error: Testing Request Queue - CreateMessageQueue\n");
1542// return false;
1543// }
1544// if (!MemoryMaps::SetProcessQueueID(procID2, REQQ_ID, qID)) {
1545// printf("Error: Testing Request Queue - SetProcessQueueID\n");
1546// return false;
1547// }
1548//
1549// ComponentData* compData1 = ComponentData::CreateComponent(compID1, "Comp1", 10*1024, 1, procID1);
1550// if (compData1 == NULL) {
1551// printf("Error: Testing Request Queue - CreateComponent1\n");
1552// return false;
1553// }
1554// ComponentData* compData2 = ComponentData::CreateComponent(compID2, "Comp2", 10*1024, 1, procID2);
1555// if (compData2 == NULL) {
1556// printf("Error: Testing Request Queue - CreateComponent2\n");
1557// return false;
1558// }
1559//
1560// DataMessage* reqMsg = new DataMessage(CTRL_TEST, compID1, compID2);
1561// uint32 reqID = 0;
1562// if (!MemoryQueues::AddRequest(reqMsg, reqID)) {
1563// printf("Error: Testing Request Queue - AddRequest\n");
1564// return false;
1565// }
1568// delete(reqMsg);
1569//
1570// if (!MemoryMaps::GetProcessQueueID(procID2, REQQ_ID, qID)) {
1571// printf("Error: Testing Request Queue - GetProcessQueueID\n");
1572// return false;
1573// }
1574// reqMsg = MemoryQueues::WaitForMessageQueue(qID, 1000);
1575// if (reqMsg == NULL) {
1576// printf("Error: Testing Request Queue - WaitForMessageQueue\n");
1577// return false;
1578// }
1579//
1580// uint32 reqID2 = (uint32)reqMsg->getReference();
1581// if (reqID != reqID2) {
1582// printf("Error: Testing Request Queue - getReference mismatch\n");
1583// return false;
1584// }
1585//
1586// DataMessage* replyMsg = new DataMessage(CTRL_TEST, compID2, 0, 1000000000);
1587// if (!MemoryQueues::AddReply(reqID, true, replyMsg)) {
1588// printf("Error: Testing Request Queue - AddReply\n");
1589// return false;
1590// }
1591// delete(reqMsg);
1592// delete(replyMsg);
1593// replyMsg = NULL;
1594//
1595// uint8 status;
1596// if (!MemoryQueues::WaitForReply(reqID2, 1000, status, &replyMsg)) {
1597// printf("Error: Testing Request Queue - WaitForReply\n");
1598// return false;
1599// }
1600//
1601// if (status != REQ_SUCCESS_DATA) {
1602// printf("Error: Testing Request Queue - Status mismatch\n");
1603// return false;
1604// }
1605//
1606// if (replyMsg == NULL) {
1607// printf("Error: Testing Request Queue - ReplyMsg NULL\n");
1608// return false;
1609// }
1610
1611 //delete(replyMsg);
1612 //delete(compData1);
1613 //delete(compData2);
1614
1615 unittest::progress(100, "done");
1617 delete(manager);
1618 return true;
1619}
1620
1621
1622//bool ProcessMemory::UnitTestQueues() {
1623// printf("Testing Multithreaded Memory Queues...\n\n");
1624//
1625// // First create and initialise the MemoryManager
1626// MemoryManager* manager = new MemoryManager();
1627// if (!manager->create(0, 30)) {
1628// fprintf(stderr, "MemoryManager init() failed...\n");
1629// delete(manager);
1630// return false;
1631// }
1632//
1633// uint64 t1, t2, t3;
1634// DataMessage* msg;
1635// int64 c = 0;
1636// uint32 q1, q2;
1637//
1638// if (!MemoryQueues::CreateMessageQueue(q1, "Q1")) {
1639// LogPrint(0, 0, 0, "[1] Could not create Test Queue 1...");
1640// delete(manager);
1641// return false;
1642// }
1643// if (!MemoryQueues::CreateMessageQueue(q2, "Q2")) {
1644// LogPrint(0, 0, 0, "[1] Could not create Test Queue 2...");
1645// delete(manager);
1646// return false;
1647// }
1648//
1649// uint32 queueTestThreadID;
1650// if (!ThreadManager::CreateThread(QueueTest, NULL, queueTestThreadID)) {
1651// LogPrint(0, 0, 0, "[1] Could not create Queue Test thread...");
1652// delete(manager);
1653// return false;
1654// }
1655//
1656// msg = new DataMessage();
1657// msg->setInt("Counter", 0);
1658//
1659// if (!MemoryQueues::AddMessageToQueue(q1, msg)) {
1660// LogPrint(0, 0, 0, "[1] Could not add initial message to Q1...");
1661// delete(manager);
1662// return false;
1663// }
1664//
1665// uint64 start = 0;
1666//
1667// while (true) {
1668// t1 = GetTimeNow();
1669// if (msg = MemoryQueues::WaitForMessageQueue(q2, 50)) {
1670// t2 = GetTimeNow();
1671// msg->getInt("Counter", c);
1672// if (c && (c % 99999 == 0)) {
1673// if (start)
1674// LogPrint(0,0,0,"[1] Average path time: %.3fus", (double)GetTimeAge(start)/99999.0);
1675// start = GetTimeNow();
1676// }
1677// msg->setInt("Counter", c+1);
1678// msg->setSendTime(GetTimeNow());
1679// if (!MemoryQueues::AddMessageToQueue(q1, msg)) {
1680// LogPrint(0, 0, 0, "[1] Could not add message %lld to Q1...", c+1);
1681// delete(msg);
1682// thread_ret_val(0);
1683// }
1684// t3 = GetTimeNow();
1685// if (t3-t1 > 10000)
1686// 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());
1687// delete(msg);
1688// }
1689// else {
1690// t3 = GetTimeNow();
1691// LogPrint(0,0,0,"[1 - %lld] Timeout: %lld\n", c+1, t3-t1);
1692// }
1693// }
1694//
1695// return true;
1696//}
1697//
1699
1700 MemoryManager* manager = (MemoryManager*) arg;
1701
1702 uint64 t = 0;
1703 uint32 c = 0;
1704 DataMessage* msg;
1705
1706 while (!g_processTestStop) {
1707
1708 if (msg = manager->processMemory->waitForCmdQ(0, 1000)) {
1709 c++;
1710 if (!manager->processMemory->addToSigQ(0, msg)) {
1711 fprintf(stderr, "Test Slave could not add Message %u to Queue...\n", c);
1712 delete(msg);
1713 thread_ret_val(0);
1714 }
1715 delete(msg);
1716 }
1717 }
1718 thread_ret_val(0);
1719
1720 //uint64 t1, t2, t3;
1721 //int64 c = 0;
1722 //uint32 q1, q2;
1723
1724 //if (!MemoryQueues::GetMessageQueueByName(q1, "Q1")) {
1725 // LogPrint(0, 0, 0, "[2] Could not Get Test Queue 1...");
1726 // thread_ret_val(0);
1727 //}
1728
1729 //if (!MemoryQueues::GetMessageQueueByName(q2, "Q2")) {
1730 // LogPrint(0, 0, 0, "[2] Could not Get Test Queue 2...");
1731 // thread_ret_val(0);
1732 //}
1733
1734 //while (true) {
1735 // t1 = GetTimeNow();
1736 // if (msg = MemoryQueues::WaitForMessageQueue(q1, 50)) {
1737 // t2 = GetTimeNow();
1738 // msg->getInt("Counter", c);
1739 // msg->setInt("Counter", c+1);
1740 // msg->setSendTime(GetTimeNow());
1741 // if (!MemoryQueues::AddMessageToQueue(q2, msg)) {
1742 // LogPrint(0, 0, 0, "[2] Could not add message %lld to Q2...", c+1);
1743 // delete(msg);
1744 // thread_ret_val(0);
1745 // }
1746 // t3 = GetTimeNow();
1747 // if (t3-t1 > 10000)
1748 // 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());
1749 // delete(msg);
1750 // }
1751 // else {
1752 // t3 = GetTimeNow();
1753 // LogPrint(0,0,0,"[2 - %lld] Timeout: %lld\n", c+1, t3-t1);
1754 // }
1755 //}
1756
1757 thread_ret_val(0);
1758}
1759
1761
1762 MemoryManager* manager = (MemoryManager*) arg;
1763
1764 uint64 t = 0;
1765 uint32 c = 0;
1766 DataMessage* msg;
1767
1768 while (!g_processTestStop) {
1769 if (msg = manager->processMemory->waitForCmdQ(0, 1000)) {
1770 c++;
1771 if (!manager->processMemory->addToMsgQ(0, msg)) {
1772 fprintf(stderr, "Test Slave could not add Message %u to Queue...\n", c);
1773 thread_ret_val(-1);
1774 }
1775 delete(msg);
1776 }
1777 }
1778 thread_ret_val(0);
1779}
1780
1782 unittest::progress(0, "init memory manager");
1783
1784 // First create and initialise the MemoryManager
1785 MemoryManager* manager = new MemoryManager();
1786 uint32 maxPageCount = 100;
1787 if (!manager->create(0)) {
1788 unittest::fail("MemoryManager create(0) failed");
1790 delete(manager);
1791 return false;
1792 }
1793
1794 uint32 count = 10000;
1795 uint32 writeCount = 0;
1796 uint32 readCount = 0;
1797 uint32 n;
1798 DataMessage* msg;
1799
1800 uint32 queueTestThreadID;
1801 g_processTestStop = false;
1802 if (!ThreadManager::CreateThread(ProcessMemoryPerfTest, manager, queueTestThreadID)) {
1803 unittest::fail("Could not create Queue Test thread");
1805 delete(manager);
1806 return false;
1807 }
1808 g_processTestThreadID = queueTestThreadID;
1809
1810 unittest::progress(10, "dual-thread round trips");
1811 uint64 perfStart = GetTimeNow();
1812 const uint32 perfIters = 20000;
1813 uint64 dualTotal = 0;
1814 uint64 t = 0;
1815 for (n=0; n<perfIters; n++) {
1816 msg = new DataMessage(CTRL_TEST, 0);
1817 if (!manager->processMemory->addToCmdQ(0, msg)) {
1818 unittest::fail("Could not add Message %u to CmdQueue", n);
1819 delete(manager);
1820 return false;
1821 }
1822 delete(msg);
1823 if ( (msg = manager->processMemory->waitForMsgQ(0, 10000)) == NULL) {
1824 unittest::fail("Could not get Message [%u] from MsgQueue", n);
1825 delete(manager);
1826 return false;
1827 }
1828 t += GetTimeAge(msg->getCreatedTime());
1829 delete(msg);
1830 if (n && n % 10000 == 0) {
1831 unittest::detail("Dual threaded message delay[%u]: %.3fus", n, ((double)t)/(2*10000.0));
1832 dualTotal += t;
1833 t = 0;
1834 unittest::progress(10 + (int)(80ULL * n / perfIters), "dual-thread round trips");
1835 }
1836 }
1837 dualTotal += t;
1838
1839 double elapsedUs = (double)(GetTimeNow() - perfStart);
1840 // each iteration measures created-time age twice (enqueue + dequeue)
1841 double dualLatency = (double)dualTotal / (2.0 * perfIters);
1842 unittest::metric("dual_thread_latency", dualLatency, "us", false);
1843 if (elapsedUs > 0.0)
1844 unittest::metric("dual_thread_throughput", (double)perfIters / elapsedUs * 1e6, "msg/s", true);
1845
1846 unittest::progress(100, "done");
1848 delete(manager);
1849 return true;
1850}
1851
1854 "Process memory queues: signal/command/message queue round-trips and integrity", "memory");
1856 "Process memory dual-threaded queue latency and throughput", "memory");
1857}
1858
1859} // namespace cmlabs
1860
1861
1862
1863
1864
1865
1866
1867
1868
1869
1870
1871
1872
#define DRAFTMSGSIZE
Fixed byte size of the per-message draft slots used in the recent-message rings of ComponentStats / P...
#define PROCESSMEMORYID
Definition ObjectIDs.h:49
#define DATAMESSAGEID
Definition ObjectIDs.h:75
#define GetObjID(data)
Extract the cid field from a binary object block: the uint32 at byte offset 4 (after the leading size...
Definition ObjectIDs.h:26
Shared-memory process ("space") table plus per-process message queues.
#define PSYPROC_ACTIVE
Process is actively running.
#define MSGQ_TYPE
Data-message queue.
#define CHECKPROCESSMEMORYSERIAL
Re-open the process segment if the master's resize serial no longer matches ours (another process re-...
#define SIGQ_TYPE
Signal queue.
#define PSYPROC_IDLE
Process is idle.
#define MAXPROC4
Maximum queues per node (MAXPROC x 4 queue types).
#define PSYPROC_CREATED
Table entry created; process not yet started.
#define INITIALQSIZE
Initial byte size of a newly created per-process queue.
#define CMDQ_TYPE
Command queue.
#define MAXPROC
Maximum processes (spaces) per node.
#define REQQ_TYPE
Request queue.
Process-wide thread registry and lifecycle manager: the concurrency core of CMSDK.
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 stricmp
Definition Utils.h:132
#define THREAD_FUNCTION_CALL
Definition Utils.h:129
#define LOG_MEMORY
Definition Utils.h:199
#define LogPrint
Definition Utils.h:313
#define MAXCOMMANDLINELEN
Definition Utils.h:88
#define THREAD_ARG
Definition Utils.h:130
The central Psyclone data container: a self-contained binary message with typed, named user entries.
bool setTime(const char *key, uint64 value)
setTime(const char* key, uint64 value)
bool setString(const char *key, const char *value)
setString(const char* key, const char* value)
DataMessageHeader * data
Pointer to the message's flat memory block (header + user entries).
uint32 getSize()
getSize() Get message size Many types of data of any size can be put into a message as user entries; ...
uint64 getTime(const char *key)
getTime(const char* key)
uint64 getCreatedTime()
getCreatedTime()
const char * getString(const char *key)
getString(const char* key)
Handle to the node's master shared-memory segment (MemoryMasterStruct).
Top-level facade of the shared-memory subsystem for one process.
ProcessMemory * processMemory
Accessor for the process table and per-process queues.
bool create(uint16 sysID, uint32 slotCount=100000, uint16 binCount=2, uint32 minBlockSize=1024, uint32 maxBlockSize=64 *1024, uint64 initSize=50000000L, uint64 maxSize=1000000000L, bool force=false)
Create all shared segments for a new node instance (master process only).
bool setProcessCommandLine(uint16 id, const char *cmdline)
Record the launch command line.
DataMessage * waitForCmdQ(uint16 procID, uint32 timeout)
Wait on the command queue.
static bool UnitTest()
Self-test.
bool addToCmdQ(uint16 procID, DataMessage *msg)
Enqueue on the command queue.
bool getProcessName(uint16 id, char *name, uint32 maxSize)
Copy the process name into name (max maxSize bytes).
uint32 getProcessOSID(uint16 id)
bool createNewProcess(const char *name, uint16 &id)
Register a new process (space) in the table.
uint32 getCmdQCount(uint16 procID)
bool addToReqQ(uint16 procID, DataMessage *msg)
Enqueue on the request queue.
bool addToMsgQ(uint16 procID, DataMessage *msg)
Enqueue on the data-message queue.
ProcessMemory(MasterMemory *master)
static bool PerfTest()
Queue throughput benchmark.
bool addLocalPerformanceStats(std::list< PerfStats > &perfStats)
Append PerfStats snapshots for all local processes to perfStats.
uint32 getSigQCount(uint16 procID)
AveragePerfStats getProcessPerfStats(uint16 procID)
bool addToAllSignalQs(DataMessage *msg)
Broadcast to every process's signal queue.
bool addToSigQ(uint16 procID, DataMessage *msg)
Enqueue on the signal queue.
bool addToAllSignalQsExcept(DataMessage *msg, uint16 except)
Broadcast to all signal queues except process except.
uint64 getProcessCreateTime(uint16 id)
uint32 getMsgQCount(uint16 procID)
DataMessage * waitForMsgQ(uint16 procID, uint32 timeout)
Wait on the data-message queue.
bool setProcessOSID(uint16 id, uint32 osid)
Record the OS pid.
bool setProcessStatus(uint16 id, uint8 status, uint64 currentCPUTicks=0)
Update status + heartbeat (and optionally CPU ticks) for process id.
bool getQueueSizes(uint16 procID, uint64 &bytes, uint32 &count)
Total queued bytes/messages across all four queues of procID.
uint8 getProcessStatus(uint16 id, uint64 &lastseen)
bool getMemoryUsage(uint64 &alloc, uint64 &usage)
Report allocation/usage of the process segment.
bool deleteProcess(uint16 id)
Remove a process entry and free its queues.
bool checkProcessHeartbeats(uint32 timeoutMS, std::list< ProcessInfoStruct > &procIssues)
Find processes whose heartbeat is older than timeoutMS.
bool getProcessID(const char *name, uint16 &id)
Look up a process id by name.
bool open()
Attach to the existing process segment.
bool addToProcessStats(uint16 id, DataMessage *inputMsg, DataMessage *outputMsg)
Add message traffic to the process's stats/rings.
DataMessage * waitForReqQ(uint16 procID, uint32 timeout)
Wait on the request queue.
uint32 getReqQCount(uint16 procID)
bool setProcessType(uint16 id, uint8 type)
Set type (0 = normal, 1 = inside node).
uint16 getProcessIDFromOSID(uint32 osid)
Map an OS pid to a Psyclone process id.
bool getProcessCommandLine(uint16 id, char *cmdline, uint32 maxSize)
Copy the launch command line.
std::vector< ProcessInfoStruct > * getAllProcesses()
Snapshot all process records.
DataMessage * waitForSigQ(uint16 procID, uint32 timeout)
Wait on the signal queue.
bool setProcessPerfStats(uint16 procID, AveragePerfStats &perfStruct)
Store new performance averages.
bool create(uint32 initialProcCount)
Create the process segment (master only).
static bool CreateThread(THREAD_FUNCTION func, void *args, uint32 &newID, uint32 reqID=0)
Create a new native thread and start it immediately.
static bool IsThreadRunning(uint32 id)
Check whether the thread is still alive at the OS level.
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.
Recursive mutual-exclusion lock, optionally named for cross-process use.
Definition Utils.h:463
Counting semaphore, optionally named for cross-process use.
Definition Utils.h:502
std::string PrintTimeString(uint64 t, bool local=true, bool us=true, bool ms=true)
Definition PsyTime.cpp:676
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
int64 GetTimeAge(uint64 t)
Age of a timestamp relative to now.
Definition PsyTime.cpp:25
uint64 GetProcessMemoryUsage()
Current resident memory usage of this process.
Definition Utils.cpp:4557
bool Sleep(uint32 ms)
Suspend the calling thread.
Definition Utils.cpp:2802
char * OpenSharedMemorySegment(const char *name, uint64 size)
Open and map an existing named shared memory segment.
Definition Utils.cpp:2259
Semaphore * GetSemaphore(const char *name, bool autocreate=true)
Look up (and optionally create) a named semaphore in the global registry.
Definition Utils.cpp:726
char * CreateSharedMemorySegment(const char *name, uint64 size, bool force=false)
Create a named shared memory segment and map it into this process.
Definition Utils.cpp:2156
bool CloseSharedMemorySegment(char *data, uint64 size)
Unmap a segment previously created/opened here.
Definition Utils.cpp:2374
std::string BytifySize(double val)
Format a byte count with binary units, e.g.
Definition Utils.cpp:7724
std::string StringFormat(const char *format,...)
printf into a std::string.
Definition Utils.cpp:6626
uint32 strcpyavail(char *dst, const char *src, uint32 maxlen, bool copyAvailable)
Bounded strcpy that always NUL-terminates.
Definition Utils.cpp:6056
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.
THREAD_RET THREAD_FUNCTION_CALL ProcessMemoryPerfTest(THREAD_ARG arg)
static volatile bool g_processTestStop
static uint32 g_processTestThreadID
static struct PsyType CTRL_TEST
Definition ObjectIDs.h:82
static THREAD_RET THREAD_FUNCTION_CALL QueueTest(THREAD_ARG arg)
void Register_ProcessMemory_Tests()
static void StopProcessTestThread()
Sliding-window averages derived from successive PerfStats snapshots.
uint32 size
Total size of the whole message block in bytes (header + all entries).
On-segment header of a circular message queue (legacy 32-bit layout).
uint32 count
Number of messages currently in the queue.
uint32 startPos
Read offset into the circular buffer.
uint32 endPos
Write offset into the circular buffer.
Raw per-component performance counters sampled inside shared memory.
uint16 spaceID
Id of the space (process slot) hosting the component.
uint32 osID
OS process id of the hosting process.
uint64 totalQueueBytes
Bytes currently waiting in the component's queues.
uint64 totalInputCount
Total number of input messages.
uint64 currentCPUTicks
Cumulative CPU ticks consumed so far.
uint64 totalOutputCount
Total number of output messages.
uint64 totalOutputBytes
Total bytes posted as output messages.
uint64 currentMemoryBytes
Current memory footprint in bytes.
uint64 totalInputBytes
Total bytes received as input messages.
uint16 nodeID
Id of the node the component runs on.
uint64 firstRunStartTime
Timestamp (µs) of the very first run.
uint32 totalQueueCount
Messages currently waiting in the component's queues.
One process (space) record in the shared process table.
uint16 cmdQID
Queue index of this process's command queue.
uint16 sigQID
Queue index of this process's signal queue.
uint64 createdTime
Creation timestamp (µs).
uint16 msgQID
Queue index of this process's data-message queue.
uint16 reqQID
Queue index of this process's request queue.
uint32 osID
OS process id (pid).
char name[MAXKEYNAMELEN+1]
Process (space) name.
uint32 id
Process id; corresponds to the index in the table.
AveragePerfStats perfStats
Windowed performance averages.
uint64 lastseen
Last heartbeat/status update (µs); staleness implies a crashed space.
ProcessStats stats
Throughput counters and recent-message rings.
uint8 status
PSYPROC_* lifecycle status.
uint16 nodeID
Id of the node this process belongs to.
uint8 type
0: normal space process, 1: runs inside the node process.
Root header of the process-memory segment: fixed process table plus queue index.
ProcessInfoStruct processes[MAXPROC]
Fixed-size process table.
uint64 qIndex[MAXPROC *4]
Byte offsets of each process's four queues (0 = none).
uint32 cid
Check/magic id validated on attach.
uint64 msgOutBytes
Total bytes sent.
uint64 msgInBytes
Total bytes received.
uint64 time
Timestamp of the last update (µs).
uint64 msgOutCount
Total messages sent.
uint64 msgInCount
Total messages received.
uint64 currentCPUTicks
Cumulative CPU ticks consumed.