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 // ⚠️ If you ever need to identify a queued message from inside this function (as the
909 // BakeLive Anomaly investigation did), note two traps:
910 // - `msg` here is usually the CTRL_TRIGGER WRAPPER from
911 // Node::deliverTriggerMessage. Its serial was NEVER SET before that function was
912 // changed to stamp it, so any trace keyed on msg->getSerial() reads 0 for every
913 // write and cannot tell one message from another.
914 // - reading the attached PAYLOAD's serial here is not cheap: getData() only returns
915 // entries whose cid is CHARDATAID, but setAttachedMessage stores the payload as
916 // DATAMESSAGEID (setRawData), so getData("Message") returns NULL.
917 // getAttachedMessageCopy() works but ALLOCATES - unacceptable on this path.
918 // The wrapper is now stamped with its payload's serial at construction, so reading
919 // msg->getSerial() is the cheap correct answer. @see BakeLiveBug.md
920 uint32 msgSize = msg->getSize();
921 uint64 spaceToEnd, inUse = 0, inUse2 = 0;
922 char* qData = ((char*)qHeader) + sizeof(MessageQueueHeader);
923 uint64 qDataSize = qHeader->size - sizeof(MessageQueueHeader);
924
925 //if (qHeader->endPos >= qHeader->startPos)
926 // inUse = qHeader->endPos - qHeader->startPos;
927 //else
928 // inUse = qDataSize - qHeader->startPos - qHeader->padding + qHeader->endPos;
929
930 if (qHeader->endPos == qHeader->startPos) {
931 spaceToEnd = qDataSize;
932 if (spaceToEnd >= msgSize) {
933 // add msg at buffer start
934 memcpy(qData, msg->data, msgSize);
935 qHeader->startPos = 0;
936 qHeader->endPos = msgSize;
937 qHeader->count++;
938 qHeader->padding = 0;
939 }
940 else {
941 // we need to resize queue
942 if (!resizeQ(qHeader->id, qHeader->size*2)) {
943 // ⚠️ Full AND cannot grow: the CALLER will drop this message, and the
944 // only trace is a level-0 log there which the test framework suppresses.
945 // A message lost this way is invisible to every instrument. @see
946 // BakeLiveBug.md - this was one of the sites traced during that hunt.
947 return false;
948 }
949 // call this function again...
950 return addToQ(procID, qType, msg);
951 }
952 }
953 else if (qHeader->endPos > qHeader->startPos) {
954 spaceToEnd = qDataSize - qHeader->endPos;
955 if (spaceToEnd >= msgSize) {
956 // add msg here
957 memcpy(qData+qHeader->endPos, msg->data, msgSize);
958 qHeader->endPos += msgSize;
959 qHeader->count++;
960 qHeader->padding = 0;
961 }
962 else if (qHeader->startPos > msgSize) {
963 // fill end rest of space with 0
964 memset(qData+qHeader->endPos, 0, (size_t)spaceToEnd);
965 qHeader->padding = spaceToEnd;
966 // add msg at buffer start
967 memcpy(qData, msg->data, msgSize);
968 qHeader->endPos = msgSize;
969 qHeader->count++;
970 }
971 else {
972 // we need to resize queue
973 if (!resizeQ(qHeader->id, qHeader->size*2)) {
974 // ⚠️ Full AND cannot grow: the CALLER will drop this message, and the
975 // only trace is a level-0 log there which the test framework suppresses.
976 // A message lost this way is invisible to every instrument. @see
977 // BakeLiveBug.md - this was one of the sites traced during that hunt.
978 return false;
979 }
980 // call this function again...
981 return addToQ(procID, qType, msg);
982 }
983 }
984 else {
985 spaceToEnd = qHeader->startPos - qHeader->endPos;
986 if (spaceToEnd > msgSize) {
987 // add msg here
988 memcpy(qData+qHeader->endPos, msg->data, msgSize);
989 qHeader->endPos += msgSize;
990 qHeader->count++;
991 // NOTE: this branch, unlike the two above, does NOT reset padding. If a
992 // wrap earlier set padding non-zero, it stays set while endPos now advances
993 // in the wrapped region - and waitForQ's end-of-buffer test subtracts
994 // padding. Recording padding here is the point of the trace.
995 }
996 else {
997 // we need to resize queue
998 if (!resizeQ(qHeader->id, qHeader->size*2)) {
999 // ⚠️ Full AND cannot grow: the CALLER will drop this message, and the
1000 // only trace is a level-0 log there which the test framework suppresses.
1001 // A message lost this way is invisible to every instrument. @see
1002 // BakeLiveBug.md - this was one of the sites traced during that hunt.
1003 return false;
1004 }
1005 // call this function again...
1006 return addToQ(procID, qType, msg);
1007 }
1008 }
1009
1010 //if (qHeader->endPos >= qHeader->startPos)
1011 // inUse2 = qHeader->endPos - qHeader->startPos;
1012 //else
1013 // inUse2 = qDataSize - qHeader->startPos - qHeader->padding + qHeader->endPos;
1014 //if (inUse2 < (qHeader->count * msgSize))
1015 // int ee = 0;
1016
1017 utils::Semaphore** sem = &qSemaphores[qHeader->id];
1018 if (!*sem)
1019 if (!(*sem = utils::GetSemaphore(qHeader->name)))
1020 return false;
1021 (*sem)->signal();
1022
1023// utils::SignalSemaphore(qHeader->name);
1024// printf("Signal %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(GetTimeNow(), true).c_str());
1025 return true;
1026}
1027
1028DataMessage* ProcessMemory::waitForQ(uint16 procID, uint8 qType, uint32 timeout) {
1029 //char* name = NULL;
1030 MessageQueueHeader* qHeader = getQHeader(procID, qType);
1031 if (!qHeader)
1032 return NULL;
1033
1034 uint64 start = GetTimeNow();
1035 int32 remain;
1036 while (!qHeader->count) {
1037 utils::Semaphore** sem = &qSemaphores[qHeader->id];
1038 if (!*sem) {
1039 if (!(*sem = utils::GetSemaphore(qHeader->name))) {
1040 return NULL;
1041 }
1042 }
1043
1044 //if (!name) {
1045 // name = new char[strlen(qHeader->name)+1];
1046 // utils::strcpyavail(name, qHeader->name, (uint32)strlen(qHeader->name)+1, true);
1047 //}
1048 if ((remain = timeout - GetTimeAgeMS(start)) <= 0) {
1049 // printf("StartWait1 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(start, true).c_str());
1050 // printf("Timeout1 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(GetTimeNow(), true).c_str());
1051 // delete [] name;
1052 return NULL;
1053 }
1054 mutex->leave();
1055 bool result = (*sem)->wait(remain);
1056 if (!mutex || !mutex->enter(5000, __FUNCTION__))
1057 return NULL;
1059
1060 if (!result) {
1061 // printf("StartWait2 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(start, true).c_str());
1062 // printf("Timeout2 %s at %s\n", qHeader->name, PrintTimeSortableMicrosecString(GetTimeNow(), true).c_str());
1063 // delete [] name;
1064 // printf("Waited %.3fms (%ums)...\n", GetTimeAge(start)/1000.0, remain);
1065 return NULL;
1066 }
1067
1068 // header might be changed
1069 qHeader = getQHeader(procID, qType);
1070 if (!qHeader) {
1071 // mutex->leave();
1072 // delete [] name;
1073 return NULL;
1074 }
1075 }
1076 //delete [] name;
1077
1078 // ⚠️ LATENT (not GAP4): qHeader is computed at the top of this function, and the
1079 // `while (!qHeader->count)` loop above is SKIPPED ENTIRELY when the queue is
1080 // non-empty on entry - so the re-fetch inside that loop ("header might be
1081 // changed") never runs on the fast path. Meanwhile ProcessMemory::resize() does
1082 // not grow the mapping in place: it creates a NEW segment, memcpy's, unmaps the
1083 // old one and repoints `header`, so anything derived from the old header dangles.
1084 //
1085 // Adding `CHECKPROCESSMEMORYSERIAL` + a re-fetch here was TRIED as the GAP4 fix
1086 // and MEASURABLY DID NOT FIX IT: with the guard in place bakelive still failed
1087 // with byte-identical figures (published=64 distinct=58 missing=6 duped=8, same
1088 // six serials). So this is NOT the GAP4 mechanism, and the guard was reverted
1089 // rather than left in place looking like a fix. It is still a real latent hazard
1090 // worth closing on its own merits, with its own test - see Psyclone.md.
1091 DataMessage* msg = NULL;
1092 if (qHeader->count) {
1093 // ⚠️ READ-POINTER CONTINUITY was measured here during the BakeLive Anomaly hunt:
1094 // does every read take its message from the offset the PREVIOUS read of this queue
1095 // left the pointer at? Result: CLEAN, including inside genuine failing runs - the
1096 // pointer never rewound and never advanced by anything other than msg->getSize().
1097 // The instrumentation was removed with the rest of that investigation.
1098 // If you re-add it: record in memory only. A file append on this path once took
1099 // bakelive from failing 1-in-26 to 76 consecutive passes - the probe hid the bug.
1100 // @see BakeLiveBug.md
1101 msg = new DataMessage(((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos, true);
1102 qHeader->startPos += msg->getSize();
1103 if (qHeader->startPos >= qHeader->size - sizeof(MessageQueueHeader) - qHeader->padding) {
1104 qHeader->startPos = 0;
1105 qHeader->padding = 0;
1106 }
1107 qHeader->count--;
1108//if (qHeader->count && (GetObjID(((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos) != DATAMESSAGEID)) {
1109// printQ(qHeader);
1110// return false;
1111//}
1112 }
1113
1114 return msg;
1115}
1116
1117bool ProcessMemory::printQ(MessageQueueHeader* qHeader) {
1118 printf("[%u] Data Size: %llu, count: %u startpos: %llu, endpos: %llu, padding: %llu\n",
1119 qHeader->id, (uint64)(qHeader->size - sizeof(MessageQueueHeader)), qHeader->count, qHeader->startPos, qHeader->endPos, qHeader->padding);
1120 char* data = ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos;
1121 for (uint32 n=0; n<qHeader->count; n++) {
1122 if (GetObjID(data) == DATAMESSAGEID)
1123 printf(" [%u: %llu - %llu] Msg size: %u Time: %s Type: %u.%u.%u\n", n,
1124 (uint64)(data - ((char*)qHeader) - sizeof(MessageQueueHeader)),
1125 (uint64)(data - ((char*)qHeader) - sizeof(MessageQueueHeader) + ((DataMessageHeader*)data)->size),
1126 ((DataMessageHeader*)data)->size, PrintTimeString(((DataMessageHeader*)data)->time).c_str(),
1127 (((DataMessageHeader*)data)->type).levels[0],
1128 (((DataMessageHeader*)data)->type).levels[1],
1129 (((DataMessageHeader*)data)->type).levels[2]);
1130 else {
1131 printf(" [%u: %llu] Msg corrupted\n",
1132 n, (uint64)(data - ((char*)qHeader) - sizeof(MessageQueueHeader)));
1133 return false;
1134 }
1135 data += ((DataMessageHeader*)data)->size;
1136 if (data == ((char*)qHeader) + qHeader->size - qHeader->padding)
1137 data = ((char*)qHeader) + sizeof(MessageQueueHeader);
1138 }
1139 return true;
1140}
1141
1142bool ProcessMemory::resize(uint64 newMemorySize) {
1143
1144 if (newMemorySize > 100000000) {
1145 //LogPrint(0,LOG_MEMORY,0,"Memory Error: Process Memory requested size %s denied...", utils::BytifySize((double)newMemorySize).c_str());
1146 printf("Memory Error: Process Memory requested size %s denied...\n", utils::BytifySize((double)newMemorySize).c_str());
1147 fflush(stdout);
1148 return false;
1149 }
1150
1151 serial = master->incrementProcessShmemSerial();
1152 //LogPrint(0,LOG_MEMORY,0," --------- Resizing ProcessMemory section %u -> %u serial %u --------", memorySize, newMemorySize, serial);
1153
1154 if (!mutex || !mutex->enter(5000, __FUNCTION__))
1155 return false;
1156
1157 //ProcessMemoryStruct* newHeader = (ProcessMemoryStruct*) utils::OpenSharedMemorySegment(utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), newMemorySize);
1158 //if (newHeader) {
1159 // utils::CloseSharedMemorySegment((char*)newHeader, newMemorySize);
1160 // //LogPrint(0,LOG_MEMORY,2,"ProcessMemory removing stale shared memory (%u/%u)...", port, serial);
1161 //}
1162
1163 ProcessMemoryStruct* newHeader = (ProcessMemoryStruct*) utils::CreateSharedMemorySegment(utils::StringFormat("PsycloneProcessMemory_%u_%u", port, serial).c_str(), newMemorySize, true);
1164 if (!newHeader) {
1165 mutex->leave();
1166 return false;
1167 }
1168 memcpy(newHeader, header, (size_t)header->size);
1169 newHeader->size = newMemorySize;
1170
1171 utils::CloseSharedMemorySegment((char*)header, memorySize);
1172 header = newHeader;
1173 memorySize = newMemorySize;
1174 master->setProcessShmemSize(memorySize);
1175 mutex->leave();
1176
1177 LogPrint(0,LOG_MEMORY,2,"Process Memory resized to %s...", utils::BytifySize((double)newMemorySize).c_str());
1178 if (newMemorySize > 20000000) {
1179 LogPrint(0,LOG_MEMORY,2,"Memory Warning: Process Memory now %s...", utils::BytifySize((double)newMemorySize).c_str());
1180 }
1181
1182 return true;
1183}
1184
1185bool ProcessMemory::resizeQ(uint16 id, uint64 newSize) {
1186 if (!header->qIndex[id])
1187 return false;
1188 MessageQueueHeader* qHeader = (MessageQueueHeader*) (((char*)header) + header->qIndex[id]);
1189 if (qHeader->id != id)
1190 return false;
1191
1192 uint64 difSize = newSize - qHeader->size;
1193 if (header->size - header->usage < difSize) {
1194 if (!resize(memorySize * 2))
1195 return false;
1196 return resizeQ(id, newSize);
1197 }
1198
1199 uint16 nextID = 0;
1200 uint32 n;
1201 for (n=id+1; n<MAXPROC4; n++) {
1202 if (header->qIndex[n]) {
1203 nextID = n;
1204 break;
1205 }
1206 }
1207
1208 // shift qID > id down by difsize
1209 if (nextID) {
1210 // Push every subsequent queue by size
1211 char* data = ((char*)header) + header->qIndex[nextID];
1212 uint64 sizeMove = header->usage - header->qIndex[nextID];
1213 memmove(data+difSize, data, (size_t)sizeMove);
1214 // Adjust indexes
1215 for (n=nextID; n<MAXPROC4; n++) {
1216 if (header->qIndex[n])
1217 header->qIndex[n] += difSize;
1218 }
1219 }
1220
1221 // Now rearrange messages in queue
1222 int64 size = qHeader->endPos - qHeader->startPos;
1223 uint64 firstSize = 0;
1224 uint64 secondSize = 0;
1225 if (size == 0)
1226 qHeader->endPos = qHeader->startPos = 0;
1227 else if ((size > 0) && qHeader->startPos) {
1228 memmove(
1229 ((char*)qHeader) + sizeof(MessageQueueHeader),
1230 ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos,
1231 (size_t)size);
1232 qHeader->endPos = (uint64)size;
1233 qHeader->startPos = 0;
1234 }
1235 else if (size < 0) {
1236 firstSize = qHeader->size - sizeof(MessageQueueHeader) - qHeader->startPos - qHeader->padding;
1237 secondSize = qHeader->endPos;
1238
1239 // is there enough space to move things around without a safety copy involving a malloc?
1240 if (difSize > secondSize) {
1241 // move first part of memory to the newly available space at the end
1242 memmove(
1243 ((char*)qHeader) + qHeader->size,
1244 ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos,
1245 (size_t)firstSize);
1246 // move last part of memory to the right place
1247 memmove(
1248 ((char*)qHeader) + sizeof(MessageQueueHeader) + firstSize,
1249 ((char*)qHeader) + sizeof(MessageQueueHeader),
1250 (size_t)secondSize);
1251 // copy first part back
1252 memmove(
1253 ((char*)qHeader) + sizeof(MessageQueueHeader),
1254 ((char*)qHeader) + qHeader->size,
1255 (size_t)firstSize);
1256 }
1257 else {
1258 // make safety copy of first part
1259 char* copy = (char*)malloc((size_t)firstSize);
1260 memcpy(copy, ((char*)qHeader) + sizeof(MessageQueueHeader) + qHeader->startPos, (size_t)firstSize);
1261 // move last part of memory to the right place
1262 memmove(
1263 ((char*)qHeader) + sizeof(MessageQueueHeader) + firstSize,
1264 ((char*)qHeader) + sizeof(MessageQueueHeader),
1265 (size_t)secondSize);
1266 // copy first part back from the safety copy
1267 memcpy(((char*)qHeader) + sizeof(MessageQueueHeader), copy, (size_t)firstSize);
1268 free(copy);
1269 }
1270 qHeader->endPos = firstSize + secondSize;
1271 qHeader->startPos = 0;
1272 }
1273
1274 qHeader->size = newSize;
1275 qHeader->padding = 0;
1276 header->usage += difSize;
1277
1278//if (qHeader->count) {
1279// printQ(qHeader);
1280// return false;
1281//}
1282
1283 return true;
1284}
1285
1286
1288 unittest::progress(0, "init memory manager");
1289
1290 // First create and initialise the MemoryManager
1291 MemoryManager* manager = new MemoryManager();
1292 if (!manager->create(0)) {
1293 unittest::fail("MemoryManager create(0) failed");
1295 delete(manager);
1296 return false;
1297 }
1298
1299 uint32 count = 10000;
1300 uint32 writeCount = 0;
1301 uint32 readCount = 0;
1302 uint32 n;
1303 DataMessage* msg;
1304 char str[128];
1305 const char* cstr;
1306
1307 uint32 queueTestThreadID;
1308 g_processTestStop = false;
1309 if (!ThreadManager::CreateThread(QueueTest, manager, queueTestThreadID)) {
1310 unittest::fail("Could not create Queue Test thread");
1312 delete(manager);
1313 return false;
1314 }
1315 g_processTestThreadID = queueTestThreadID;
1316
1317 uint64 time = GetTimeNow();
1318
1319 // empty-queue wait must actually block for ~1s
1320 unittest::progress(10, "empty queue wait");
1321 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) != NULL) {
1322 unittest::fail("Got unexpected Message from empty Queue");
1324 delete(manager);
1325 return false;
1326 }
1327
1328 uint64 t = 0;
1329 if ( (t = GetTimeAge(time)) < 950000) {
1330 unittest::fail("Message Queue wait 1000ms didn't actually wait more than %.3fms", (t/1000.0));
1332 delete(manager);
1333 return false;
1334 }
1335
1336 // Single-threaded round-trip latency (same thread adds + reads sig queue)
1337 unittest::progress(25, "single-thread latency");
1338 const uint32 latIters = 20000;
1339 uint64 singleTotal = 0;
1340 t = 0;
1341 for (n=0; n<latIters; n++) {
1342 msg = new DataMessage(CTRL_TEST, 0);
1343 if (!manager->processMemory->addToSigQ(0, msg)) {
1344 unittest::fail("Could not add Message %u to SigQueue", n);
1345 delete(manager);
1346 return false;
1347 }
1348 delete(msg);
1349 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) == NULL) {
1350 unittest::fail("Could not get Message [%u] from SigQueue", n);
1351 delete(manager);
1352 return false;
1353 }
1354 t += GetTimeAge(msg->getCreatedTime());
1355 delete(msg);
1356 if (n && n % 10000 == 0) {
1357 unittest::detail("Single threaded message delay[%u]: %.3fus", n, ((double)t)/(2*10000.0));
1358 singleTotal += t;
1359 t = 0;
1360 }
1361 }
1362 singleTotal += t;
1363 // each iteration measures the created-time age twice (add + read)
1364 double singleLatency = (double)singleTotal / (2.0 * latIters);
1365 unittest::metric("single_thread_latency", singleLatency, "us", false);
1366
1367 // warm up the queue (cross-thread path via cmd queue -> slave -> sig queue)
1368 msg = new DataMessage(CTRL_TEST, 0);
1369 if (!manager->processMemory->addToCmdQ(0, msg)) {
1370 unittest::fail("Could not add startup Message to Queue");
1372 delete(manager);
1373 return false;
1374 }
1375 delete(msg);
1376 utils::Sleep(10);
1377 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) == NULL) {
1378 unittest::fail("Could not get startup Message from Queue");
1380 delete(manager);
1381 return false;
1382 }
1383 delete(msg);
1384
1385 // Multi-threaded round-trip latency (slave thread relays cmd -> sig)
1386 unittest::progress(45, "multi-thread latency");
1387 uint64 multiTotal = 0;
1388 t = 0;
1389 for (n=0; n<latIters; n++) {
1390 msg = new DataMessage(CTRL_TEST, 0);
1391 if (!manager->processMemory->addToCmdQ(0, msg)) {
1392 unittest::fail("Could not add Message %u to CmdQueue", n);
1393 delete(manager);
1394 return false;
1395 }
1396 delete(msg);
1397 if ( (msg = manager->processMemory->waitForSigQ(0, 1000)) == NULL) {
1398 unittest::fail("Could not get Message [%u] from SigQueue", n);
1399 delete(manager);
1400 return false;
1401 }
1402 t += GetTimeAge(msg->getCreatedTime());
1403 delete(msg);
1404 if (n && n % 10000 == 0) {
1405 unittest::detail("Multi threaded message delay[%u]: %.3fus", n, ((double)t)/(2*10000.0));
1406 multiTotal += t;
1407 t = 0;
1408 }
1409 }
1410 multiTotal += t;
1411 double multiLatency = (double)multiTotal / (2.0 * latIters);
1412 unittest::metric("multi_thread_latency", multiLatency, "us", false);
1413
1414
1415
1416
1417
1418
1419 // Message queue correctness + write throughput
1420 unittest::progress(65, "msg queue fill");
1421 time = GetTimeNow();
1422 uint64 writeStart = GetTimeNow();
1423 for (n=0; n<count/4; n++) {
1424 msg = new DataMessage(CTRL_TEST, writeCount);
1425 snprintf(str, 128, "Test%2u", writeCount);
1426 msg->setString("TestEntry", str);
1427 msg->setTime("TestTime", time);
1428 msg->setString("TestEntry2", str);
1429 msg->setTime("TestTime2", time);
1430 if (!manager->processMemory->addToMsgQ(0, msg)) {
1431 unittest::fail("Could not add Message %u to MsgQueue", n);
1432 delete(manager);
1433 return false;
1434 }
1435 // if (n == 1000)
1436 // utils::PrintBinary(msg->data, msg->data->size, false, "Msg1000");
1437 delete(msg);
1438 writeCount++;
1439 }
1440
1441 if ( (n = manager->processMemory->getMsgQCount(0)) != count/4) {
1442 unittest::fail("Queue contains %u messages instead of %u", n, count/4);
1444 delete(manager);
1445 return false;
1446 }
1447
1448 unittest::progress(75, "msg queue drain");
1449 for (n=0; n<count/8; n++) {
1450 if ( (msg = manager->processMemory->waitForMsgQ(0, 1000)) == NULL) {
1451 unittest::fail("Could not get Message %u from Queue", n);
1452 delete(manager);
1453 return false;
1454 }
1455 snprintf(str, 128, "Test%2u", readCount);
1456 if (strcmp(str, msg->getString("TestEntry")) != 0) {
1457 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, msg->getString("TestEntry"));
1458 delete(msg);
1459 delete(manager);
1460 return false;
1461 }
1462 if (msg->getTime("TestTime") != time) {
1463 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime"));
1464 delete(msg);
1465 delete(manager);
1466 return false;
1467 }
1468 cstr = msg->getString("TestEntry2");
1469 if (!cstr || strcmp(str, cstr) != 0) {
1470 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, cstr);
1471 delete(msg);
1472 delete(manager);
1473 return false;
1474 }
1475 if (msg->getTime("TestTime2") != time) {
1476 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime2"));
1477 delete(msg);
1478 delete(manager);
1479 return false;
1480 }
1481 readCount++;
1482 delete(msg);
1483 }
1484
1485 if ( (n = manager->processMemory->getMsgQCount(0)) != count/8) {
1486 unittest::fail("Queue contains %u messages instead of %u", n, count/8);
1488 delete(manager);
1489 return false;
1490 }
1491
1492 unittest::progress(85, "msg queue refill");
1493 for (n=0; n<3*count/4; n++) {
1494 msg = new DataMessage(CTRL_TEST, writeCount);
1495 snprintf(str, 128, "Test%2u", writeCount);
1496 msg->setString("TestEntry", str);
1497 msg->setTime("TestTime", time);
1498 msg->setString("TestEntry2", str);
1499 msg->setTime("TestTime2", time);
1500 if (!manager->processMemory->addToMsgQ(0, msg)) {
1501 unittest::fail("Could not add Message %u to MsgQueue", n);
1502 delete(manager);
1503 return false;
1504 }
1505 delete(msg);
1506 writeCount++;
1507 }
1508
1509 if ( (n = manager->processMemory->getMsgQCount(0)) != 7*count/8) {
1510 unittest::fail("Queue contains %u messages instead of %u", n, 7*count/8);
1512 delete(manager);
1513 return false;
1514 }
1515
1516 unittest::progress(92, "msg queue verify");
1517 for (n=0; n<7*count/8; n++) {
1518 if ( (msg = manager->processMemory->waitForMsgQ(0, 1000)) == NULL) {
1519 unittest::fail("Could not get Message %u from Queue", n);
1520 delete(manager);
1521 return false;
1522 }
1523 snprintf(str, 128, "Test%2u", readCount);
1524 if (strcmp(str, msg->getString("TestEntry")) != 0) {
1525 unittest::fail("Message %u from Queue corrupted string '%s' (queue has %u)", readCount, msg->getString("TestEntry"), manager->processMemory->getMsgQCount(0));
1526 delete(msg);
1527 delete(manager);
1528 return false;
1529 }
1530 if (msg->getTime("TestTime") != time) {
1531 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime"));
1532 delete(msg);
1533 delete(manager);
1534 return false;
1535 }
1536 cstr = msg->getString("TestEntry2");
1537 if (!cstr || strcmp(str, cstr) != 0) {
1538 unittest::fail("Message %u from Queue corrupted string '%s'", readCount, cstr);
1539 delete(msg);
1540 delete(manager);
1541 return false;
1542 }
1543 if (msg->getTime("TestTime2") != time) {
1544 unittest::fail("Message %u from Queue corrupted time '%llu'", readCount, msg->getTime("TestTime2"));
1545 delete(msg);
1546 delete(manager);
1547 return false;
1548 }
1549 readCount++;
1550 delete(msg);
1551 }
1552
1553 if ( (n = manager->processMemory->getMsgQCount(0)) != 0) {
1554 unittest::fail("Queue contains %u messages instead of 0", n);
1556 delete(manager);
1557 return false;
1558 }
1559
1560 // Aggregate write throughput over all enqueued messages (count/4 + 3*count/4)
1561 double writeUs = (double)(GetTimeNow() - writeStart);
1562 if (writeUs > 0.0)
1563 unittest::metric("msg_queue_throughput", (double)count / writeUs * 1e6, "msg/s", true);
1564
1565
1566// printf("Testing Request Queue...\n\n");
1567//
1568// // First create two components
1569// uint16 procID1, procID2;
1570// uint32 compID1 = 1, compID2 = 2;
1571// uint32 qID;
1572// //if (!MemoryMaps::CreateNewNode(1, "", 1)) {
1573// // printf("Error: Testing Request Queue - CreateNewNode\n");
1574// // return false;
1575// //}
1576// if (!MemoryMaps::CreateNewProcess("Proc1", procID1)) {
1577// printf("Error: Testing Request Queue - CreateNewProcess\n");
1578// return false;
1579// }
1580// if (!MemoryMaps::CreateNewProcess("Proc2", procID2)) {
1581// printf("Error: Testing Request Queue - CreateNewProcess 2\n");
1582// return false;
1583// }
1584// if (!MemoryQueues::CreateMessageQueue(qID)) {
1585// printf("Error: Testing Request Queue - CreateMessageQueue\n");
1586// return false;
1587// }
1588// if (!MemoryMaps::SetProcessQueueID(procID1, MSGQ_ID, qID)) {
1589// printf("Error: Testing Request Queue - SetProcessQueueID\n");
1590// return false;
1591// }
1592// if (!MemoryQueues::CreateMessageQueue(qID)) {
1593// printf("Error: Testing Request Queue - CreateMessageQueue\n");
1594// return false;
1595// }
1596// if (!MemoryMaps::SetProcessQueueID(procID2, REQQ_ID, qID)) {
1597// printf("Error: Testing Request Queue - SetProcessQueueID\n");
1598// return false;
1599// }
1600//
1601// ComponentData* compData1 = ComponentData::CreateComponent(compID1, "Comp1", 10*1024, 1, procID1);
1602// if (compData1 == NULL) {
1603// printf("Error: Testing Request Queue - CreateComponent1\n");
1604// return false;
1605// }
1606// ComponentData* compData2 = ComponentData::CreateComponent(compID2, "Comp2", 10*1024, 1, procID2);
1607// if (compData2 == NULL) {
1608// printf("Error: Testing Request Queue - CreateComponent2\n");
1609// return false;
1610// }
1611//
1612// DataMessage* reqMsg = new DataMessage(CTRL_TEST, compID1, compID2);
1613// uint32 reqID = 0;
1614// if (!MemoryQueues::AddRequest(reqMsg, reqID)) {
1615// printf("Error: Testing Request Queue - AddRequest\n");
1616// return false;
1617// }
1620// delete(reqMsg);
1621//
1622// if (!MemoryMaps::GetProcessQueueID(procID2, REQQ_ID, qID)) {
1623// printf("Error: Testing Request Queue - GetProcessQueueID\n");
1624// return false;
1625// }
1626// reqMsg = MemoryQueues::WaitForMessageQueue(qID, 1000);
1627// if (reqMsg == NULL) {
1628// printf("Error: Testing Request Queue - WaitForMessageQueue\n");
1629// return false;
1630// }
1631//
1632// uint32 reqID2 = (uint32)reqMsg->getReference();
1633// if (reqID != reqID2) {
1634// printf("Error: Testing Request Queue - getReference mismatch\n");
1635// return false;
1636// }
1637//
1638// DataMessage* replyMsg = new DataMessage(CTRL_TEST, compID2, 0, 1000000000);
1639// if (!MemoryQueues::AddReply(reqID, true, replyMsg)) {
1640// printf("Error: Testing Request Queue - AddReply\n");
1641// return false;
1642// }
1643// delete(reqMsg);
1644// delete(replyMsg);
1645// replyMsg = NULL;
1646//
1647// uint8 status;
1648// if (!MemoryQueues::WaitForReply(reqID2, 1000, status, &replyMsg)) {
1649// printf("Error: Testing Request Queue - WaitForReply\n");
1650// return false;
1651// }
1652//
1653// if (status != REQ_SUCCESS_DATA) {
1654// printf("Error: Testing Request Queue - Status mismatch\n");
1655// return false;
1656// }
1657//
1658// if (replyMsg == NULL) {
1659// printf("Error: Testing Request Queue - ReplyMsg NULL\n");
1660// return false;
1661// }
1662
1663 //delete(replyMsg);
1664 //delete(compData1);
1665 //delete(compData2);
1666
1667 unittest::progress(100, "done");
1669 delete(manager);
1670 return true;
1671}
1672
1673
1674//bool ProcessMemory::UnitTestQueues() {
1675// printf("Testing Multithreaded Memory Queues...\n\n");
1676//
1677// // First create and initialise the MemoryManager
1678// MemoryManager* manager = new MemoryManager();
1679// if (!manager->create(0, 30)) {
1680// fprintf(stderr, "MemoryManager init() failed...\n");
1681// delete(manager);
1682// return false;
1683// }
1684//
1685// uint64 t1, t2, t3;
1686// DataMessage* msg;
1687// int64 c = 0;
1688// uint32 q1, q2;
1689//
1690// if (!MemoryQueues::CreateMessageQueue(q1, "Q1")) {
1691// LogPrint(0, 0, 0, "[1] Could not create Test Queue 1...");
1692// delete(manager);
1693// return false;
1694// }
1695// if (!MemoryQueues::CreateMessageQueue(q2, "Q2")) {
1696// LogPrint(0, 0, 0, "[1] Could not create Test Queue 2...");
1697// delete(manager);
1698// return false;
1699// }
1700//
1701// uint32 queueTestThreadID;
1702// if (!ThreadManager::CreateThread(QueueTest, NULL, queueTestThreadID)) {
1703// LogPrint(0, 0, 0, "[1] Could not create Queue Test thread...");
1704// delete(manager);
1705// return false;
1706// }
1707//
1708// msg = new DataMessage();
1709// msg->setInt("Counter", 0);
1710//
1711// if (!MemoryQueues::AddMessageToQueue(q1, msg)) {
1712// LogPrint(0, 0, 0, "[1] Could not add initial message to Q1...");
1713// delete(manager);
1714// return false;
1715// }
1716//
1717// uint64 start = 0;
1718//
1719// while (true) {
1720// t1 = GetTimeNow();
1721// if (msg = MemoryQueues::WaitForMessageQueue(q2, 50)) {
1722// t2 = GetTimeNow();
1723// msg->getInt("Counter", c);
1724// if (c && (c % 99999 == 0)) {
1725// if (start)
1726// LogPrint(0,0,0,"[1] Average path time: %.3fus", (double)GetTimeAge(start)/99999.0);
1727// start = GetTimeNow();
1728// }
1729// msg->setInt("Counter", c+1);
1730// msg->setSendTime(GetTimeNow());
1731// if (!MemoryQueues::AddMessageToQueue(q1, msg)) {
1732// LogPrint(0, 0, 0, "[1] Could not add message %lld to Q1...", c+1);
1733// delete(msg);
1734// thread_ret_val(0);
1735// }
1736// t3 = GetTimeNow();
1737// if (t3-t1 > 10000)
1738// 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());
1739// delete(msg);
1740// }
1741// else {
1742// t3 = GetTimeNow();
1743// LogPrint(0,0,0,"[1 - %lld] Timeout: %lld\n", c+1, t3-t1);
1744// }
1745// }
1746//
1747// return true;
1748//}
1749//
1751
1752 MemoryManager* manager = (MemoryManager*) arg;
1753
1754 uint64 t = 0;
1755 uint32 c = 0;
1756 DataMessage* msg;
1757
1758 while (!g_processTestStop) {
1759
1760 if (msg = manager->processMemory->waitForCmdQ(0, 1000)) {
1761 c++;
1762 if (!manager->processMemory->addToSigQ(0, msg)) {
1763 fprintf(stderr, "Test Slave could not add Message %u to Queue...\n", c);
1764 delete(msg);
1765 thread_ret_val(0);
1766 }
1767 delete(msg);
1768 }
1769 }
1770 thread_ret_val(0);
1771
1772 //uint64 t1, t2, t3;
1773 //int64 c = 0;
1774 //uint32 q1, q2;
1775
1776 //if (!MemoryQueues::GetMessageQueueByName(q1, "Q1")) {
1777 // LogPrint(0, 0, 0, "[2] Could not Get Test Queue 1...");
1778 // thread_ret_val(0);
1779 //}
1780
1781 //if (!MemoryQueues::GetMessageQueueByName(q2, "Q2")) {
1782 // LogPrint(0, 0, 0, "[2] Could not Get Test Queue 2...");
1783 // thread_ret_val(0);
1784 //}
1785
1786 //while (true) {
1787 // t1 = GetTimeNow();
1788 // if (msg = MemoryQueues::WaitForMessageQueue(q1, 50)) {
1789 // t2 = GetTimeNow();
1790 // msg->getInt("Counter", c);
1791 // msg->setInt("Counter", c+1);
1792 // msg->setSendTime(GetTimeNow());
1793 // if (!MemoryQueues::AddMessageToQueue(q2, msg)) {
1794 // LogPrint(0, 0, 0, "[2] Could not add message %lld to Q2...", c+1);
1795 // delete(msg);
1796 // thread_ret_val(0);
1797 // }
1798 // t3 = GetTimeNow();
1799 // if (t3-t1 > 10000)
1800 // 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());
1801 // delete(msg);
1802 // }
1803 // else {
1804 // t3 = GetTimeNow();
1805 // LogPrint(0,0,0,"[2 - %lld] Timeout: %lld\n", c+1, t3-t1);
1806 // }
1807 //}
1808
1809 thread_ret_val(0);
1810}
1811
1813
1814 MemoryManager* manager = (MemoryManager*) arg;
1815
1816 uint64 t = 0;
1817 uint32 c = 0;
1818 DataMessage* msg;
1819
1820 while (!g_processTestStop) {
1821 if (msg = manager->processMemory->waitForCmdQ(0, 1000)) {
1822 c++;
1823 if (!manager->processMemory->addToMsgQ(0, msg)) {
1824 fprintf(stderr, "Test Slave could not add Message %u to Queue...\n", c);
1825 thread_ret_val(-1);
1826 }
1827 delete(msg);
1828 }
1829 }
1830 thread_ret_val(0);
1831}
1832
1834 unittest::progress(0, "init memory manager");
1835
1836 // First create and initialise the MemoryManager
1837 MemoryManager* manager = new MemoryManager();
1838 uint32 maxPageCount = 100;
1839 if (!manager->create(0)) {
1840 unittest::fail("MemoryManager create(0) failed");
1842 delete(manager);
1843 return false;
1844 }
1845
1846 uint32 count = 10000;
1847 uint32 writeCount = 0;
1848 uint32 readCount = 0;
1849 uint32 n;
1850 DataMessage* msg;
1851
1852 uint32 queueTestThreadID;
1853 g_processTestStop = false;
1854 if (!ThreadManager::CreateThread(ProcessMemoryPerfTest, manager, queueTestThreadID)) {
1855 unittest::fail("Could not create Queue Test thread");
1857 delete(manager);
1858 return false;
1859 }
1860 g_processTestThreadID = queueTestThreadID;
1861
1862 unittest::progress(10, "dual-thread round trips");
1863 uint64 perfStart = GetTimeNow();
1864 const uint32 perfIters = 20000;
1865 uint64 dualTotal = 0;
1866 uint64 t = 0;
1867 for (n=0; n<perfIters; n++) {
1868 msg = new DataMessage(CTRL_TEST, 0);
1869 if (!manager->processMemory->addToCmdQ(0, msg)) {
1870 unittest::fail("Could not add Message %u to CmdQueue", n);
1871 delete(manager);
1872 return false;
1873 }
1874 delete(msg);
1875 if ( (msg = manager->processMemory->waitForMsgQ(0, 10000)) == NULL) {
1876 unittest::fail("Could not get Message [%u] from MsgQueue", n);
1877 delete(manager);
1878 return false;
1879 }
1880 t += GetTimeAge(msg->getCreatedTime());
1881 delete(msg);
1882 if (n && n % 10000 == 0) {
1883 unittest::detail("Dual threaded message delay[%u]: %.3fus", n, ((double)t)/(2*10000.0));
1884 dualTotal += t;
1885 t = 0;
1886 unittest::progress(10 + (int)(80ULL * n / perfIters), "dual-thread round trips");
1887 }
1888 }
1889 dualTotal += t;
1890
1891 double elapsedUs = (double)(GetTimeNow() - perfStart);
1892 // each iteration measures created-time age twice (enqueue + dequeue)
1893 double dualLatency = (double)dualTotal / (2.0 * perfIters);
1894 unittest::metric("dual_thread_latency", dualLatency, "us", false);
1895 if (elapsedUs > 0.0)
1896 unittest::metric("dual_thread_throughput", (double)perfIters / elapsedUs * 1e6, "msg/s", true);
1897
1898 unittest::progress(100, "done");
1900 delete(manager);
1901 return true;
1902}
1903
1906 "Process memory queues: signal/command/message queue round-trips and integrity", "memory");
1908 "Process memory dual-threaded queue latency and throughput", "memory");
1910 "Process memory circular-queue seeded oracle fuzz (single process)", "memory", false);
1911 // Phase 2: OFF by default like the Phase 1 fuzz (last arg false) - it spawns threads
1912 // and runs for tens of seconds, so it must be named explicitly, never dragged into
1913 // test=cmsdk.
1915 "Process memory ring: concurrent multi-reader / burst-drain / resize-under-read / "
1916 "blocking-path fuzz (set-equality oracle + read-pointer continuity)", "memory", false);
1917}
1918
1919} // namespace cmlabs
1920
1921
1922
1923
1924
1925
1926
1927
1928
1929
1930
1931
1932
#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.
static bool FuzzTest()
Seeded, oracle-checked fuzz of the circular queue arithmetic (single process).
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).
static bool ConcFuzzTest()
Phase 2: CONCURRENT + cross-process fuzz of the ring (ProcessMemoryConcFuzz.cpp).
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:519
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
bool Sleep(uint32 ms)
Suspend the calling thread.
Definition Utils.cpp:3121
char * OpenSharedMemorySegment(const char *name, uint64 size)
Open and map an existing named shared memory segment.
Definition Utils.cpp:2489
Semaphore * GetSemaphore(const char *name, bool autocreate=true)
Look up (and optionally create) a named semaphore in the global registry.
Definition Utils.cpp:748
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:2368
bool CloseSharedMemorySegment(char *data, uint64 size)
Unmap a segment previously created/opened here.
Definition Utils.cpp:2604
uint32 strcpyavail(char *dst, const char *src, uint32 maxlen, bool copyAvailable)
Bounded strcpy that always NUL-terminates.
Definition Utils.cpp:7497
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.
uint64 GetProcessMemoryUsage()
Current resident memory usage of this process.
Definition Utils.cpp:5983
std::string BytifySize(double val)
Format a byte count with binary units, e.g.
Definition Utils.cpp:9165
std::string StringFormat(const char *format,...)
printf into a std::string.
Definition Utils.cpp:8067
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:83
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.