From 4b66e957cccdd16e5df61fa9b4bf0097f3cf0581 Mon Sep 17 00:00:00 2001 From: randogoth Date: Sat, 2 Mar 2024 19:01:51 +0200 Subject: [PATCH] added named pipe --- .gitignore | 3 +- QNGmeter.cpp | 549 ++++++++++++++++++++++++++++++++++----------------- 2 files changed, 366 insertions(+), 186 deletions(-) diff --git a/.gitignore b/.gitignore index 59a4f91..eefeeb2 100644 --- a/.gitignore +++ b/.gitignore @@ -3,4 +3,5 @@ CMakeFiles/ CMakeCache.txt Makefile *.bin -qngmeter \ No newline at end of file +qngmeter +*.json \ No newline at end of file diff --git a/QNGmeter.cpp b/QNGmeter.cpp index 7feae1e..c937d11 100755 --- a/QNGmeter.cpp +++ b/QNGmeter.cpp @@ -1,8 +1,16 @@ +#define _DEFAULT_SOURCE + #ifdef _WIN32 #include #include +#else +#include +#include +#include +#include #endif +#include #include #include #include @@ -11,12 +19,15 @@ #include #include -#include #include #include #include #include +#include +#include +#include + #include "BiasAndAC.h" #include "Monkey.h" #include "Serial.h" @@ -41,9 +52,46 @@ using namespace std; } #endif -#define _BSD_SOURCE -#include - +template +class ThreadSafeQueue { +private: + mutable std::mutex mtx; + std::deque dataQueue; + std::condition_variable dataCond; + +public: + ThreadSafeQueue() {} + + void push(T newData) { + std::lock_guard lk(mtx); + dataQueue.push_back(std::move(newData)); + dataCond.notify_one(); + } + + bool try_pop(T& value) { + std::lock_guard lk(mtx); + if (dataQueue.empty()) { + return false; + } + value = std::move(dataQueue.front()); + dataQueue.pop_front(); + return true; + } + + std::unique_lock wait_and_pop(T& value) { + std::unique_lock lk(mtx); + dataCond.wait(lk, [this] { return !dataQueue.empty(); }); + value = std::move(dataQueue.front()); + dataQueue.pop_front(); + return lk; + } + + bool empty() const { + std::lock_guard lk(mtx); + return dataQueue.empty(); + } +}; + int main() { // try this to resize terminal @@ -90,7 +138,7 @@ int main() time(&timePrev); timePrev -= 8; - queue > > dataQueue; + //queue > > dataQueue; double bitsThroughputCount = 0; double prevBitsThroughputCount = 0; @@ -119,278 +167,409 @@ int main() double meterScore = 0; bool meterFreeze = false; - omp_set_nested(true); + const char* fifoPath = "/tmp/QNGmeter"; +#ifdef __linux + mkfifo(fifoPath, 0666); // Create FIFO if it doesn't exist + int fifoFd = open(fifoPath, O_RDONLY | O_NONBLOCK); + if (fifoFd == -1) { + std::cerr << "Failed to open named pipe for reading." << std::endl; + return 1; + } +#endif - #pragma omp parallel sections num_threads(2) + fd_set readfds; + int maxfd = max(STDIN_FILENO, fifoFd) + 1; + + ThreadSafeQueue dataQueue; // Use a thread-safe queue for shared data + + omp_set_nested(true); // Ensure nested parallelism is enabled + + #pragma omp parallel sections { - #pragma omp section { - // This section reads from stdin as binary data and enqueues it for testing - const size_t blockSize = 2048; // Number of uint32_t values in a block - const size_t bytesPerValue = sizeof(uint32_t); - const size_t bufferSize = blockSize * bytesPerValue; // Total bytes per block - char buffer[bufferSize]; // Temporary buffer to store bytes - while (!cin.eof() && !doExit) { - cin.read(buffer, bufferSize); - size_t bytesRead = cin.gcount(); + fd_set readfds; + int fifoFd = open ("/tmp/myfifo", O_RDONLY | O_NONBLOCK); + if (fifoFd == -1) + { + perror ("Failed to open FIFO"); + exit (EXIT_FAILURE); + } - // Convert read bytes to uint32_t and store in a vector - shared_ptr> newBuffer(new vector()); - for (size_t i = 0; i < bytesRead; i += bytesPerValue) { - if (i + bytesPerValue <= bytesRead) { - // Ensure we have a full 4 bytes to read - uint32_t value = 0; - memcpy(&value, buffer + i, bytesPerValue); - newBuffer->push_back(value); + int maxfd = max (STDIN_FILENO, fifoFd) + 1; + ThreadSafeQueue < std::string > dataQueue; + + while (!doExit) + { + FD_ZERO (&readfds); + FD_SET (STDIN_FILENO, &readfds); + FD_SET (fifoFd, &readfds); + + if (select (maxfd, &readfds, NULL, NULL, NULL) == -1) + { + perror ("select"); + exit (EXIT_FAILURE); + } + + char buffer[4096]; + ssize_t bytesRead; + + if (FD_ISSET (STDIN_FILENO, &readfds)) + { + bytesRead = read (STDIN_FILENO, buffer, sizeof (buffer)); + if (bytesRead > 0) + { + dataQueue.push (std::string (buffer, bytesRead)); } } - if (!newBuffer->empty()) { - #pragma omp critical - dataQueue.push(newBuffer); + if (FD_ISSET (fifoFd, &readfds)) + { + bytesRead = read (fifoFd, buffer, sizeof (buffer)); + if (bytesRead > 0) + { + dataQueue.push (std::string (buffer, bytesRead)); + } } } - doExit = true; // Exit if stdin closes or reaches EOF - } + close (fifoFd); + } #pragma omp section { - // This section processes the enqueued data - while (!doExit || !dataQueue.empty()) { - shared_ptr> testBuffer = nullptr; - #pragma omp critical + std::string buffer; // This will hold the raw binary string data from the queue + while (!doExit) + { + if (dataQueue.try_pop (buffer)) { - if (!dataQueue.empty()) { - testBuffer = dataQueue.front(); - dataQueue.pop(); - } - } - - if (testBuffer != nullptr) { - // Parallel processing of the testBuffer - #pragma omp parallel sections num_threads(4) + // Assuming each uint32_t is stored in 4 bytes within the string + // Convert the string buffer into a vector for processing + std::vector < uint32_t > values; + const size_t numValues = buffer.size () / sizeof (uint32_t); + for (size_t i = 0; i < numValues; ++i) { - #pragma omp section - { for (uint32_t value : *testBuffer) biasAndAc.InsertWord32(value); } - #pragma omp section - { for (uint32_t value : *testBuffer) oqso.InsertWord32(value); } - #pragma omp section - { for (uint32_t value : *testBuffer) serial.InsertWord32(value); } - #pragma omp section - { for (uint32_t value : *testBuffer) entropy.InsertWord32(value); } + uint32_t value; + memcpy (&value, buffer.data () + i * sizeof (uint32_t), + sizeof (uint32_t)); + values.push_back (value); } - + + // Now, values contains the uint32_t data extracted from buffer + // Process the values as needed + #pragma omp parallel sections num_threads(4) + { + #pragma omp section + { + for (uint32_t value:values) + biasAndAc.InsertWord32 (value); + } + #pragma omp section + { + for (uint32_t value:values) + oqso.InsertWord32 (value); + } + #pragma omp section + { + for (uint32_t value:values) + serial.InsertWord32 (value); + } + #pragma omp section + { + for (uint32_t value:values) + entropy.InsertWord32 (value); + } + } + // Display results - time(&timeNow); - if (difftime(timeNow, timePrev) >= 10) + time (&timeNow); + if (difftime (timeNow, timePrev) >= 10) { // calc rates - #ifdef _WIN32 - QueryPerformanceCounter(&nowCount); - double newTimeInterval = ((double)nowCount.QuadPart - prevCount.QuadPart) / countFreq.QuadPart; - prevCount = nowCount; - #elif __linux - gettimeofday(&stop, NULL); - newTimeInterval = (stop.tv_sec - start.tv_sec); // sec - newTimeInterval += (stop.tv_usec - start.tv_usec) /1000000.0; // us to sec - start = stop; - #elif MACOSX - #endif + #ifdef _WIN32 + QueryPerformanceCounter (&nowCount); + double newTimeInterval = + ((double) nowCount.QuadPart - + prevCount.QuadPart) / countFreq.QuadPart; + prevCount = nowCount; + #elif __linux + gettimeofday (&stop, NULL); + newTimeInterval = (stop.tv_sec - start.tv_sec); // sec + newTimeInterval += (stop.tv_usec - start.tv_usec) / 1000000.0; // us to sec + start = stop; + #elif MACOSX + #endif double newBitInterval; double newRate; - #pragma omp critical + #pragma omp critical { - newBitInterval = (bitsThroughputCount - prevBitsThroughputCount); - prevBitsThroughputCount = bitsThroughputCount; + newBitInterval = + (bitsThroughputCount - prevBitsThroughputCount); + prevBitsThroughputCount = bitsThroughputCount; } newRate = newBitInterval / newTimeInterval; if (throughput == 0) - throughput = newRate; + throughput = newRate; else - throughput = (2*throughput + newRate) / 3; + throughput = (2 * throughput + newRate) / 3; - double newBitsTestedRatio = (bitsTestedCount-prevBitsTestedCount) / newBitInterval; + double newBitsTestedRatio = + (bitsTestedCount - prevBitsTestedCount) / newBitInterval; prevBitsTestedCount = bitsTestedCount; - bitsTestedRatio = (2*bitsTestedRatio + newBitsTestedRatio) / 3; + bitsTestedRatio = + (2 * bitsTestedRatio + newBitsTestedRatio) / 3; if (bitsTestedRatio > 1) - bitsTestedRatio = 1; + bitsTestedRatio = 1; // meta test and meter - metaPs.clear(); - meterZs.clear(); + metaPs.clear (); + meterZs.clear (); if (bitsTestedCount >= 65536) { // Autocorrelation KS test - for (int i=0; i<32; i++) + for (int i = 0; i < 32; i++) { - metaPs.push_back(biasAndAc.AC.P_Chi2[i]); - meterZs.push_back(biasAndAc.AC.cumulativeACZScore[i]); + metaPs.push_back (biasAndAc.AC.P_Chi2[i]); + meterZs.push_back (biasAndAc.AC. + cumulativeACZScore[i]); } double AcKSP; double AcKSN; - ks.KSUP(&AcKSP, &AcKSN, &metaPs[0], metaPs.size()); + ks.KSUP (&AcKSP, &AcKSN, &metaPs[0], metaPs.size ()); // Combined KS test - metaPs.push_back(AcKSP); + metaPs.push_back (AcKSP); - metaPs.push_back(biasAndAc.Bias.P_Chi2); - meterZs.push_back(biasAndAc.Bias.cumulativeBiasZScore); + metaPs.push_back (biasAndAc.Bias.P_Chi2); + meterZs.push_back (biasAndAc.Bias.cumulativeBiasZScore); } if (bitsTestedCount >= 4194304) { - metaPs.push_back(serial.P_Chi2); - serialP = gamma.Gamma(128., serial.cumulativeSerialChi2); - serialZ = ks.PtoZ(serialP); - meterZs.push_back(serialZ); + metaPs.push_back (serial.P_Chi2); + serialP = gamma.Gamma (128., serial.cumulativeSerialChi2); + serialZ = ks.PtoZ (serialP); + meterZs.push_back (serialZ); - metaPs.push_back(entropy.P_Chi2); - meterZs.push_back(entropy.cumulativeZScore); + metaPs.push_back (entropy.P_Chi2); + meterZs.push_back (entropy.cumulativeZScore); } if (bitsTestedCount >= 10485775) { - metaPs.push_back(oqso.P_Chi2); - meterZs.push_back(oqso.cumulativeZScore); + metaPs.push_back (oqso.P_Chi2); + meterZs.push_back (oqso.cumulativeZScore); } if (bitsTestedCount >= 65536) { // This KS is combined AC KSP plus with other tests - ks.KSUP(&KSP, &KSN, &metaPs[32], metaPs.size()-32); + ks.KSUP (&KSP, &KSN, &metaPs[32], metaPs.size () - 32); meterFreeze = false; - for (int i=0; i4.264897 || (metaPs[i]<0.00001 || metaPs[i]>0.99999)) - meterFlags[i] = -1; + if (fabs (meterZs[i]) > 4.264897 + || (metaPs[i] < 0.00001 || metaPs[i] > 0.99999)) + meterFlags[i] = -1; // unfreeze condition if (meterFlags[i] == -1) { - if (fabs(meterZs[i])<2.326348 && (metaPs[i]>0.01 && metaPs[i]<0.99)) - meterFlags[i] = 0; + if (fabs (meterZs[i]) < 2.326348 + && (metaPs[i] > 0.01 && metaPs[i] < 0.99)) + meterFlags[i] = 0; } if (meterFlags[i] == -1) - meterFreeze = true; + meterFreeze = true; } // meter calc if (meterFreeze == false) - meterScore = log(bitsTestedCount)/log(2.); + meterScore = log (bitsTestedCount) / log (2.); } cout << endl; - cout << " QNGmeter Console 1.0 Test Type z-score p[z<=x] p[chi2<=x] " << endl; - cout << " +---------------------------+------------------------------------------------+" << endl; - cout << " | | 1/0 Balance " << setiosflags(ios::fixed) << setprecision(3) << showpos << biasAndAc.Bias.cumulativeBiasZScore << " " << setprecision(4) << noshowpos << CStat::ZtoP(biasAndAc.Bias.cumulativeBiasZScore) << " " << biasAndAc.Bias.P_Chi2 << " |" << endl; - cout << " | | Serial Test " << setiosflags(ios::fixed) << setprecision(3) << showpos << serialZ << " " << setprecision(4) << noshowpos << serialP << " " << serial.P_Chi2 << " |" << endl; - cout << " | | OQSO Test " << setiosflags(ios::fixed) << setprecision(3) << showpos << oqso.cumulativeZScore << " " << setprecision(4) << noshowpos << CStat::ZtoP(oqso.cumulativeZScore) << " " << oqso.P_Chi2 << " |" << endl; - cout << " | | Entropy Test " << setiosflags(ios::fixed) << setprecision(3) << showpos << entropy.cumulativeZScore << " " << setprecision(4) << noshowpos << CStat::ZtoP(entropy.cumulativeZScore) << " " << serial.P_Chi2 << " |" << endl; - cout << " | | H: " << setiosflags(ios::fixed) << setprecision(9) << entropy.E << " |" << endl; - - cout << " | | |" << endl; + cout << + " QNGmeter Console 1.0 Test Type z-score p[z<=x] p[chi2<=x] " + << endl; + cout << + " +---------------------------+------------------------------------------------+" + << endl; + cout << " | | 1/0 Balance " + << setiosflags (ios:: + fixed) << setprecision (3) << showpos << + biasAndAc.Bias. + cumulativeBiasZScore << " " << setprecision (4) << + noshowpos << CStat::ZtoP (biasAndAc.Bias. + cumulativeBiasZScore) << " " << + biasAndAc.Bias.P_Chi2 << " |" << endl; + cout << " | | Serial Test " + << setiosflags (ios:: + fixed) << setprecision (3) << showpos << + serialZ << " " << setprecision (4) << noshowpos << + serialP << " " << serial.P_Chi2 << " |" << endl; + cout << " | | OQSO Test " + << setiosflags (ios:: + fixed) << setprecision (3) << showpos << + oqso. + cumulativeZScore << " " << setprecision (4) << noshowpos + << CStat::ZtoP (oqso.cumulativeZScore) << " " << oqso. + P_Chi2 << " |" << endl; + cout << " | | Entropy Test " + << setiosflags (ios:: + fixed) << setprecision (3) << showpos << + entropy. + cumulativeZScore << " " << setprecision (4) << noshowpos + << CStat::ZtoP (entropy. + cumulativeZScore) << " " << serial. + P_Chi2 << " |" << endl; + cout << " | | H: " << + setiosflags (ios::fixed) << setprecision (9) << entropy. + E << " |" << endl; - for (int i=1; i<=32; i++) - { - switch(i) - { - case 2: - cout << " | Start Time |"; - break; - case 3: - cout << " | " << sStartTime << " |"; - break; - case 5: - cout << " | Total Bits Tested |"; - break; - case 6: - cout << " | " << scientific << setw(9) << setprecision(2) << bitsTestedCount << fixed << " |"; - break; - case 8: - cout << " | Throughput |"; - break; - case 9: - cout << " | " << setiosflags(ios::fixed) << setw(4) << setprecision(1) << (double)(throughput/1000000.0) << " Mbps |"; - break; - case 11: - cout << " | Bits Tested Percent |"; - break; - case 12: - cout << " | " << setiosflags(ios::fixed) << setw(5) << setprecision(1) << (100*bitsTestedRatio) << "% |"; - break; - case 18: - cout << " | Meta KS+ Test |"; - break; - case 19: - cout << " | " << setiosflags(ios::fixed) << setw(5) << setprecision(3) << KSP << " |"; - break; - case 22: - cout << " | Meta KS- Test |"; - break; - case 23: - cout << " | " << setiosflags(ios::fixed) << setw(5) << setprecision(3) << KSN << " |"; - break; - case 29: - cout << " | QNGmeter Score |"; - break; - case 30: - cout << " | " << setiosflags(ios::fixed) << setw(4) << setprecision(1) << abs(meterScore) << ((meterScore<0)? "-" : (meterFreeze==false)? "+" : " ") << " |"; - break; - default: - cout << " | |"; - } - cout << " " << setw(2) << i << "st AutoCorr " << setiosflags(ios::fixed) << setprecision(3) << showpos << biasAndAc.AC.cumulativeACZScore[i-1] << " " << setprecision(4) << noshowpos << CStat::ZtoP(biasAndAc.AC.cumulativeACZScore[i-1]) << " " << biasAndAc.AC.P_Chi2[i-1] << " |" << endl; - } - cout << " +---------------------------+------------------------------------------------+" << endl; + cout << + " | | |" + << endl; - #ifdef _WIN32 - // put cursor in top corner - COORD coord; - coord.X = 0; - coord.Y = 0; - SetConsoleCursorPosition(GetStdHandle(STD_OUTPUT_HANDLE), coord); - #elif __linux - // clear screen - // printf("\E[H"); - #elif MACOSX - #endif + for (int i = 1; i <= 32; i++) + { + switch (i) + { + case 2: + cout << " | Start Time |"; + break; + case 3: + cout << " | " << sStartTime << " |"; + break; + case 5: + cout << " | Total Bits Tested |"; + break; + case 6: + cout << " | " << scientific << setw (9) << + setprecision (2) << bitsTestedCount << fixed << + " |"; + break; + case 8: + cout << " | Throughput |"; + break; + case 9: + cout << " | " << setiosflags (ios:: + fixed) << setw (4) + << setprecision (1) << (double) (throughput / + 1000000.0) << + " Mbps |"; + break; + case 11: + cout << " | Bits Tested Percent |"; + break; + case 12: + cout << " | " << setiosflags (ios:: + fixed) << setw (5) + << setprecision (1) << (100 * + bitsTestedRatio) << + "% |"; + break; + case 18: + cout << " | Meta KS+ Test |"; + break; + case 19: + cout << " | " << setiosflags (ios:: + fixed) << setw (5) + << setprecision (3) << KSP << " |"; + break; + case 22: + cout << " | Meta KS- Test |"; + break; + case 23: + cout << " | " << setiosflags (ios:: + fixed) << setw (5) + << setprecision (3) << KSN << " |"; + break; + case 29: + cout << " | QNGmeter Score |"; + break; + case 30: + cout << " | " << setiosflags (ios:: + fixed) << setw (4) + << setprecision (1) << abs (meterScore) << + ((meterScore < 0) ? "-" : (meterFreeze == + false) ? "+" : " ") << + " |"; + break; + default: + cout << " | |"; + } + cout << " " << setw (2) << i << "st AutoCorr " << + setiosflags (ios:: + fixed) << setprecision (3) << showpos << + biasAndAc.AC.cumulativeACZScore[i - + 1] << " " << + setprecision (4) << noshowpos << CStat::ZtoP (biasAndAc. + AC. + cumulativeACZScore + [i - + 1]) << + " " << biasAndAc.AC.P_Chi2[i - 1] << " |" << endl; + } + cout << + " +---------------------------+------------------------------------------------+" + << endl; + + #ifdef _WIN32 + // put cursor in top corner + COORD coord; + coord.X = 0; + coord.Y = 0; + SetConsoleCursorPosition (GetStdHandle (STD_OUTPUT_HANDLE), + coord); + #elif __linux + // clear screen + // printf("\E[H"); + #elif MACOSX + #endif timePrev = timeNow; - cout.flush(); + cout.flush (); + buffer.clear (); } - + // End on an 'x' keypress - #ifdef _WIN32 - if (kbhit()) + #ifdef _WIN32 + if (kbhit ()) + { + char c = getch_ (); + + if (tolower (c) == 'x') { - char c = getch_(); - - if ( tolower(c) == 'x' ) - { - doExit = true; - break; - } + doExit = true; + break; } - #elif __linux - // CRTL-C - #elif MACOSX - #endif + } + #elif __linux + // CRTL-C + #elif MACOSX + #endif } else { // give this thread a break from tight loop - waiting for data - usleep(1000); + usleep (1000); } } } } + + +#ifdef __linux + close(fifoFd); // Close the FIFO file descriptor when done +#endif + + return 0; }