Revert "added named pipe"

This reverts commit 4b66e957cc.
This commit is contained in:
randogoth 2024-03-02 20:44:27 +02:00
parent 4b66e957cc
commit 74ea72064c
2 changed files with 182 additions and 362 deletions

3
.gitignore vendored
View file

@ -3,5 +3,4 @@ CMakeFiles/
CMakeCache.txt CMakeCache.txt
Makefile Makefile
*.bin *.bin
qngmeter qngmeter
*.json

View file

@ -1,16 +1,8 @@
#define _DEFAULT_SOURCE
#ifdef _WIN32 #ifdef _WIN32
#include <Windows.h> #include <Windows.h>
#include <conio.h> #include <conio.h>
#else
#include <unistd.h>
#include <fcntl.h>
#include <sys/stat.h>
#include <sys/select.h>
#endif #endif
#include <sys/time.h>
#include <iostream> #include <iostream>
#include <iomanip> #include <iomanip>
#include <stdint.h> #include <stdint.h>
@ -19,15 +11,12 @@
#include <omp.h> #include <omp.h>
#include <queue> #include <queue>
#include <deque>
#include <vector> #include <vector>
#include <string> #include <string>
#include <memory> #include <memory>
#include <cmath> #include <cmath>
#include <mutex>
#include <condition_variable>
#include <deque>
#include "BiasAndAC.h" #include "BiasAndAC.h"
#include "Monkey.h" #include "Monkey.h"
#include "Serial.h" #include "Serial.h"
@ -52,46 +41,9 @@ using namespace std;
} }
#endif #endif
template<typename T> #define _BSD_SOURCE
class ThreadSafeQueue { #include <sys/time.h>
private:
mutable std::mutex mtx;
std::deque<T> dataQueue;
std::condition_variable dataCond;
public:
ThreadSafeQueue() {}
void push(T newData) {
std::lock_guard<std::mutex> lk(mtx);
dataQueue.push_back(std::move(newData));
dataCond.notify_one();
}
bool try_pop(T& value) {
std::lock_guard<std::mutex> lk(mtx);
if (dataQueue.empty()) {
return false;
}
value = std::move(dataQueue.front());
dataQueue.pop_front();
return true;
}
std::unique_lock<std::mutex> wait_and_pop(T& value) {
std::unique_lock<std::mutex> 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<std::mutex> lk(mtx);
return dataQueue.empty();
}
};
int main() int main()
{ {
// try this to resize terminal // try this to resize terminal
@ -138,7 +90,7 @@ int main()
time(&timePrev); time(&timePrev);
timePrev -= 8; timePrev -= 8;
//queue<shared_ptr<vector<uint32_t> > > dataQueue; queue<shared_ptr<vector<uint32_t> > > dataQueue;
double bitsThroughputCount = 0; double bitsThroughputCount = 0;
double prevBitsThroughputCount = 0; double prevBitsThroughputCount = 0;
@ -167,409 +119,278 @@ int main()
double meterScore = 0; double meterScore = 0;
bool meterFreeze = false; bool meterFreeze = false;
const char* fifoPath = "/tmp/QNGmeter"; omp_set_nested(true);
#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
fd_set readfds; #pragma omp parallel sections num_threads(2)
int maxfd = max(STDIN_FILENO, fifoFd) + 1;
ThreadSafeQueue<std::string> 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 #pragma omp section
{ {
fd_set readfds; // This section reads from stdin as binary data and enqueues it for testing
int fifoFd = open ("/tmp/myfifo", O_RDONLY | O_NONBLOCK); const size_t blockSize = 2048; // Number of uint32_t values in a block
if (fifoFd == -1) const size_t bytesPerValue = sizeof(uint32_t);
{ const size_t bufferSize = blockSize * bytesPerValue; // Total bytes per block
perror ("Failed to open FIFO"); char buffer[bufferSize]; // Temporary buffer to store bytes
exit (EXIT_FAILURE); while (!cin.eof() && !doExit) {
} cin.read(buffer, bufferSize);
size_t bytesRead = cin.gcount();
int maxfd = max (STDIN_FILENO, fifoFd) + 1; // Convert read bytes to uint32_t and store in a vector
ThreadSafeQueue < std::string > dataQueue; shared_ptr<vector<uint32_t>> newBuffer(new vector<uint32_t>());
for (size_t i = 0; i < bytesRead; i += bytesPerValue) {
while (!doExit) if (i + bytesPerValue <= bytesRead) {
{ // Ensure we have a full 4 bytes to read
FD_ZERO (&readfds); uint32_t value = 0;
FD_SET (STDIN_FILENO, &readfds); memcpy(&value, buffer + i, bytesPerValue);
FD_SET (fifoFd, &readfds); newBuffer->push_back(value);
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 (FD_ISSET (fifoFd, &readfds)) if (!newBuffer->empty()) {
{ #pragma omp critical
bytesRead = read (fifoFd, buffer, sizeof (buffer)); dataQueue.push(newBuffer);
if (bytesRead > 0)
{
dataQueue.push (std::string (buffer, bytesRead));
}
} }
} }
doExit = true; // Exit if stdin closes or reaches EOF
close (fifoFd);
} }
#pragma omp section #pragma omp section
{ {
std::string buffer; // This will hold the raw binary string data from the queue // This section processes the enqueued data
while (!doExit) while (!doExit || !dataQueue.empty()) {
{ shared_ptr<vector<uint32_t>> testBuffer = nullptr;
if (dataQueue.try_pop (buffer)) #pragma omp critical
{ {
// Assuming each uint32_t is stored in 4 bytes within the string if (!dataQueue.empty()) {
// Convert the string buffer into a vector<uint32_t> for processing testBuffer = dataQueue.front();
std::vector < uint32_t > values; dataQueue.pop();
const size_t numValues = buffer.size () / sizeof (uint32_t);
for (size_t i = 0; i < numValues; ++i)
{
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 if (testBuffer != nullptr) {
// Process the values as needed // Parallel processing of the testBuffer
#pragma omp parallel sections num_threads(4) #pragma omp parallel sections num_threads(4)
{ {
#pragma omp section #pragma omp section
{ { for (uint32_t value : *testBuffer) biasAndAc.InsertWord32(value); }
for (uint32_t value:values) #pragma omp section
biasAndAc.InsertWord32 (value); { 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); }
} }
#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 // Display results
time (&timeNow); time(&timeNow);
if (difftime (timeNow, timePrev) >= 10) if (difftime(timeNow, timePrev) >= 10)
{ {
// calc rates // calc rates
#ifdef _WIN32 #ifdef _WIN32
QueryPerformanceCounter (&nowCount); QueryPerformanceCounter(&nowCount);
double newTimeInterval = double newTimeInterval = ((double)nowCount.QuadPart - prevCount.QuadPart) / countFreq.QuadPart;
((double) nowCount.QuadPart - prevCount = nowCount;
prevCount.QuadPart) / countFreq.QuadPart; #elif __linux
prevCount = nowCount; gettimeofday(&stop, NULL);
#elif __linux newTimeInterval = (stop.tv_sec - start.tv_sec); // sec
gettimeofday (&stop, NULL); newTimeInterval += (stop.tv_usec - start.tv_usec) /1000000.0; // us to sec
newTimeInterval = (stop.tv_sec - start.tv_sec); // sec start = stop;
newTimeInterval += (stop.tv_usec - start.tv_usec) / 1000000.0; // us to sec #elif MACOSX
start = stop; #endif
#elif MACOSX
#endif
double newBitInterval; double newBitInterval;
double newRate; double newRate;
#pragma omp critical #pragma omp critical
{ {
newBitInterval = newBitInterval = (bitsThroughputCount - prevBitsThroughputCount);
(bitsThroughputCount - prevBitsThroughputCount); prevBitsThroughputCount = bitsThroughputCount;
prevBitsThroughputCount = bitsThroughputCount;
} }
newRate = newBitInterval / newTimeInterval; newRate = newBitInterval / newTimeInterval;
if (throughput == 0) if (throughput == 0)
throughput = newRate; throughput = newRate;
else else
throughput = (2 * throughput + newRate) / 3; throughput = (2*throughput + newRate) / 3;
double newBitsTestedRatio = double newBitsTestedRatio = (bitsTestedCount-prevBitsTestedCount) / newBitInterval;
(bitsTestedCount - prevBitsTestedCount) / newBitInterval;
prevBitsTestedCount = bitsTestedCount; prevBitsTestedCount = bitsTestedCount;
bitsTestedRatio = bitsTestedRatio = (2*bitsTestedRatio + newBitsTestedRatio) / 3;
(2 * bitsTestedRatio + newBitsTestedRatio) / 3;
if (bitsTestedRatio > 1) if (bitsTestedRatio > 1)
bitsTestedRatio = 1; bitsTestedRatio = 1;
// meta test and meter // meta test and meter
metaPs.clear (); metaPs.clear();
meterZs.clear (); meterZs.clear();
if (bitsTestedCount >= 65536) if (bitsTestedCount >= 65536)
{ {
// Autocorrelation KS test // 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]); metaPs.push_back(biasAndAc.AC.P_Chi2[i]);
meterZs.push_back (biasAndAc.AC. meterZs.push_back(biasAndAc.AC.cumulativeACZScore[i]);
cumulativeACZScore[i]);
} }
double AcKSP; double AcKSP;
double AcKSN; double AcKSN;
ks.KSUP (&AcKSP, &AcKSN, &metaPs[0], metaPs.size ()); ks.KSUP(&AcKSP, &AcKSN, &metaPs[0], metaPs.size());
// Combined KS test // Combined KS test
metaPs.push_back (AcKSP); metaPs.push_back(AcKSP);
metaPs.push_back (biasAndAc.Bias.P_Chi2); metaPs.push_back(biasAndAc.Bias.P_Chi2);
meterZs.push_back (biasAndAc.Bias.cumulativeBiasZScore); meterZs.push_back(biasAndAc.Bias.cumulativeBiasZScore);
} }
if (bitsTestedCount >= 4194304) if (bitsTestedCount >= 4194304)
{ {
metaPs.push_back (serial.P_Chi2); metaPs.push_back(serial.P_Chi2);
serialP = gamma.Gamma (128., serial.cumulativeSerialChi2); serialP = gamma.Gamma(128., serial.cumulativeSerialChi2);
serialZ = ks.PtoZ (serialP); serialZ = ks.PtoZ(serialP);
meterZs.push_back (serialZ); meterZs.push_back(serialZ);
metaPs.push_back (entropy.P_Chi2); metaPs.push_back(entropy.P_Chi2);
meterZs.push_back (entropy.cumulativeZScore); meterZs.push_back(entropy.cumulativeZScore);
} }
if (bitsTestedCount >= 10485775) if (bitsTestedCount >= 10485775)
{ {
metaPs.push_back (oqso.P_Chi2); metaPs.push_back(oqso.P_Chi2);
meterZs.push_back (oqso.cumulativeZScore); meterZs.push_back(oqso.cumulativeZScore);
} }
if (bitsTestedCount >= 65536) if (bitsTestedCount >= 65536)
{ {
// This KS is combined AC KSP plus with other tests // 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; meterFreeze = false;
for (int i = 0; i < meterZs.size (); i++) for (int i=0; i<meterZs.size(); i++)
{ {
// freeze condition // freeze condition
if (fabs (meterZs[i]) > 4.264897 if (fabs(meterZs[i])>4.264897 || (metaPs[i]<0.00001 || metaPs[i]>0.99999))
|| (metaPs[i] < 0.00001 || metaPs[i] > 0.99999)) meterFlags[i] = -1;
meterFlags[i] = -1;
// unfreeze condition // unfreeze condition
if (meterFlags[i] == -1) if (meterFlags[i] == -1)
{ {
if (fabs (meterZs[i]) < 2.326348 if (fabs(meterZs[i])<2.326348 && (metaPs[i]>0.01 && metaPs[i]<0.99))
&& (metaPs[i] > 0.01 && metaPs[i] < 0.99)) meterFlags[i] = 0;
meterFlags[i] = 0;
} }
if (meterFlags[i] == -1) if (meterFlags[i] == -1)
meterFreeze = true; meterFreeze = true;
} }
// meter calc // meter calc
if (meterFreeze == false) if (meterFreeze == false)
meterScore = log (bitsTestedCount) / log (2.); meterScore = log(bitsTestedCount)/log(2.);
} }
cout << endl; cout << endl;
cout << cout << " QNGmeter Console 1.0 Test Type z-score p[z<=x] p[chi2<=x] " << endl;
" QNGmeter Console 1.0 Test Type z-score p[z<=x] p[chi2<=x] " cout << " +---------------------------+------------------------------------------------+" << endl;
<< 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 << 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;
<< endl; cout << " | | Entropy Test " << setiosflags(ios::fixed) << setprecision(3) << showpos << entropy.cumulativeZScore << " " << setprecision(4) << noshowpos << CStat::ZtoP(entropy.cumulativeZScore) << " " << serial.P_Chi2 << " |" << endl;
cout << " | | 1/0 Balance " cout << " | | H: " << setiosflags(ios::fixed) << setprecision(9) << entropy.E << " |" << endl;
<< setiosflags (ios::
fixed) << setprecision (3) << showpos << cout << " | | |" << endl;
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 << for (int i=1; i<=32; i++)
" | | |"
<< endl;
for (int i = 1; i <= 32; i++)
{ {
switch (i) switch(i)
{ {
case 2: case 2:
cout << " | Start Time |"; cout << " | Start Time |";
break; break;
case 3: case 3:
cout << " | " << sStartTime << " |"; cout << " | " << sStartTime << " |";
break; break;
case 5: case 5:
cout << " | Total Bits Tested |"; cout << " | Total Bits Tested |";
break; break;
case 6: case 6:
cout << " | " << scientific << setw (9) << cout << " | " << scientific << setw(9) << setprecision(2) << bitsTestedCount << fixed << " |";
setprecision (2) << bitsTestedCount << fixed << break;
" |"; case 8:
break; cout << " | Throughput |";
case 8: break;
cout << " | Throughput |"; case 9:
break; cout << " | " << setiosflags(ios::fixed) << setw(4) << setprecision(1) << (double)(throughput/1000000.0) << " Mbps |";
case 9: break;
cout << " | " << setiosflags (ios:: case 11:
fixed) << setw (4) cout << " | Bits Tested Percent |";
<< setprecision (1) << (double) (throughput / break;
1000000.0) << case 12:
" Mbps |"; cout << " | " << setiosflags(ios::fixed) << setw(5) << setprecision(1) << (100*bitsTestedRatio) << "% |";
break; break;
case 11: case 18:
cout << " | Bits Tested Percent |"; cout << " | Meta KS+ Test |";
break; break;
case 12: case 19:
cout << " | " << setiosflags (ios:: cout << " | " << setiosflags(ios::fixed) << setw(5) << setprecision(3) << KSP << " |";
fixed) << setw (5) break;
<< setprecision (1) << (100 * case 22:
bitsTestedRatio) << cout << " | Meta KS- Test |";
"% |"; break;
break; case 23:
case 18: cout << " | " << setiosflags(ios::fixed) << setw(5) << setprecision(3) << KSN << " |";
cout << " | Meta KS+ Test |"; break;
break; case 29:
case 19: cout << " | QNGmeter Score |";
cout << " | " << setiosflags (ios:: break;
fixed) << setw (5) case 30:
<< setprecision (3) << KSP << " |"; cout << " | " << setiosflags(ios::fixed) << setw(4) << setprecision(1) << abs(meterScore) << ((meterScore<0)? "-" : (meterFreeze==false)? "+" : " ") << " |";
break; break;
case 22: default:
cout << " | Meta KS- Test |"; cout << " | |";
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 << 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;
" +---------------------------+------------------------------------------------+" }
<< endl; cout << " +---------------------------+------------------------------------------------+" << endl;
#ifdef _WIN32 #ifdef _WIN32
// put cursor in top corner // put cursor in top corner
COORD coord; COORD coord;
coord.X = 0; coord.X = 0;
coord.Y = 0; coord.Y = 0;
SetConsoleCursorPosition (GetStdHandle (STD_OUTPUT_HANDLE), SetConsoleCursorPosition(GetStdHandle(STD_OUTPUT_HANDLE), coord);
coord); #elif __linux
#elif __linux // clear screen
// clear screen // printf("\E[H");
// printf("\E[H"); #elif MACOSX
#elif MACOSX #endif
#endif
timePrev = timeNow; timePrev = timeNow;
cout.flush (); cout.flush();
buffer.clear ();
} }
// End on an 'x' keypress // End on an 'x' keypress
#ifdef _WIN32 #ifdef _WIN32
if (kbhit ()) if (kbhit())
{
char c = getch_ ();
if (tolower (c) == 'x')
{ {
doExit = true; char c = getch_();
break;
if ( tolower(c) == 'x' )
{
doExit = true;
break;
}
} }
} #elif __linux
#elif __linux // CRTL-C
// CRTL-C #elif MACOSX
#elif MACOSX #endif
#endif
} }
else else
{ {
// give this thread a break from tight loop - waiting for data // 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;
} }