added named pipe

This commit is contained in:
randogoth 2024-03-02 19:01:51 +02:00
parent 5389cf9df1
commit 4b66e957cc
2 changed files with 366 additions and 186 deletions

1
.gitignore vendored
View file

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

View file

@ -1,8 +1,16 @@
#define _DEFAULT_SOURCE
#ifdef _WIN32
#include <Windows.h>
#include <conio.h>
#else
#include <unistd.h>
#include <fcntl.h>
#include <sys/stat.h>
#include <sys/select.h>
#endif
#include <sys/time.h>
#include <iostream>
#include <iomanip>
#include <stdint.h>
@ -11,12 +19,15 @@
#include <omp.h>
#include <queue>
#include <deque>
#include <vector>
#include <string>
#include <memory>
#include <cmath>
#include <mutex>
#include <condition_variable>
#include <deque>
#include "BiasAndAC.h"
#include "Monkey.h"
#include "Serial.h"
@ -41,8 +52,45 @@ using namespace std;
}
#endif
#define _BSD_SOURCE
#include <sys/time.h>
template<typename T>
class ThreadSafeQueue {
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()
{
@ -90,7 +138,7 @@ int main()
time(&timePrev);
timePrev -= 8;
queue<shared_ptr<vector<uint32_t> > > dataQueue;
//queue<shared_ptr<vector<uint32_t> > > dataQueue;
double bitsThroughputCount = 0;
double prevBitsThroughputCount = 0;
@ -119,82 +167,135 @@ 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<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
{
fd_set readfds;
int fifoFd = open ("/tmp/myfifo", O_RDONLY | O_NONBLOCK);
if (fifoFd == -1)
{
perror ("Failed to open FIFO");
exit (EXIT_FAILURE);
}
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 (FD_ISSET (fifoFd, &readfds))
{
bytesRead = read (fifoFd, buffer, sizeof (buffer));
if (bytesRead > 0)
{
dataQueue.push (std::string (buffer, bytesRead));
}
}
}
close (fifoFd);
}
#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();
// Convert read bytes to uint32_t and store in a vector
shared_ptr<vector<uint32_t>> newBuffer(new vector<uint32_t>());
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);
}
}
if (!newBuffer->empty()) {
#pragma omp critical
dataQueue.push(newBuffer);
}
}
doExit = true; // Exit if stdin closes or reaches EOF
}
#pragma omp section
std::string buffer; // This will hold the raw binary string data from the queue
while (!doExit)
{
// This section processes the enqueued data
while (!doExit || !dataQueue.empty()) {
shared_ptr<vector<uint32_t>> testBuffer = nullptr;
#pragma omp critical
if (dataQueue.try_pop (buffer))
{
if (!dataQueue.empty()) {
testBuffer = dataQueue.front();
dataQueue.pop();
}
// Assuming each uint32_t is stored in 4 bytes within the string
// Convert the string buffer into a vector<uint32_t> for processing
std::vector < uint32_t > values;
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);
}
if (testBuffer != nullptr) {
// Parallel processing of the testBuffer
// 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 : *testBuffer) biasAndAc.InsertWord32(value); }
{
for (uint32_t value:values)
biasAndAc.InsertWord32 (value);
}
#pragma omp section
{ for (uint32_t value : *testBuffer) oqso.InsertWord32(value); }
{
for (uint32_t value:values)
oqso.InsertWord32 (value);
}
#pragma omp section
{ for (uint32_t value : *testBuffer) serial.InsertWord32(value); }
{
for (uint32_t value:values)
serial.InsertWord32 (value);
}
#pragma omp section
{ for (uint32_t value : *testBuffer) entropy.InsertWord32(value); }
{
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;
QueryPerformanceCounter (&nowCount);
double newTimeInterval =
((double) nowCount.QuadPart -
prevCount.QuadPart) / countFreq.QuadPart;
prevCount = nowCount;
#elif __linux
gettimeofday(&stop, NULL);
gettimeofday (&stop, NULL);
newTimeInterval = (stop.tv_sec - start.tv_sec); // sec
newTimeInterval += (stop.tv_usec - start.tv_usec) /1000000.0; // us to sec
newTimeInterval += (stop.tv_usec - start.tv_usec) / 1000000.0; // us to sec
start = stop;
#elif MACOSX
#endif
@ -203,7 +304,8 @@ int main()
double newRate;
#pragma omp critical
{
newBitInterval = (bitsThroughputCount - prevBitsThroughputCount);
newBitInterval =
(bitsThroughputCount - prevBitsThroughputCount);
prevBitsThroughputCount = bitsThroughputCount;
}
@ -211,69 +313,74 @@ int main()
if (throughput == 0)
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;
// 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; i<meterZs.size(); i++)
for (int i = 0; i < meterZs.size (); i++)
{
// freeze condition
if (fabs(meterZs[i])>4.264897 || (metaPs[i]<0.00001 || metaPs[i]>0.99999))
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))
if (fabs (meterZs[i]) < 2.326348
&& (metaPs[i] > 0.01 && metaPs[i] < 0.99))
meterFlags[i] = 0;
}
@ -283,24 +390,56 @@ int main()
// 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 <<
" 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 <<
" | | |"
<< endl;
for (int i=1; i<=32; i++)
for (int i = 1; i <= 32; i++)
{
switch(i)
switch (i)
{
case 2:
cout << " | Start Time |";
@ -312,51 +451,83 @@ int main()
cout << " | Total Bits Tested |";
break;
case 6:
cout << " | " << scientific << setw(9) << setprecision(2) << bitsTestedCount << fixed << " |";
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 |";
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) << "% |";
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 << " |";
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 << " |";
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)? "+" : " ") << " |";
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 << " " << 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);
SetConsoleCursorPosition (GetStdHandle (STD_OUTPUT_HANDLE),
coord);
#elif __linux
// clear screen
// printf("\E[H");
@ -364,16 +535,17 @@ int main()
#endif
timePrev = timeNow;
cout.flush();
cout.flush ();
buffer.clear ();
}
// End on an 'x' keypress
#ifdef _WIN32
if (kbhit())
if (kbhit ())
{
char c = getch_();
char c = getch_ ();
if ( tolower(c) == 'x' )
if (tolower (c) == 'x')
{
doExit = true;
break;
@ -387,10 +559,17 @@ int main()
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;
}