12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361 |
- /*##############################################################################
- HPCC SYSTEMS software Copyright (C) 2012 HPCC Systems®.
- Licensed under the Apache License, Version 2.0 (the "License");
- you may not use this file except in compliance with the License.
- You may obtain a copy of the License at
- http://www.apache.org/licenses/LICENSE-2.0
- Unless required by applicable law or agreed to in writing, software
- distributed under the License is distributed on an "AS IS" BASIS,
- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- See the License for the specific language governing permissions and
- limitations under the License.
- ############################################################################## */
- ///* simple test
- #include "udplib.hpp"
- #include "roxiemem.hpp"
- //#include "udptrr.hpp"
- //#include "udptrs.hpp"
- #include "jthread.hpp"
- #include "jsocket.hpp"
- #include "jsem.hpp"
- #include "jdebug.hpp"
- #include <time.h>
- #if defined(_DEBUG) && defined(_WIN32) && !defined(USING_MPATROL)
- #define new new(_NORMAL_BLOCK, __FILE__, __LINE__)
- #endif
- /*=============================================================================================
- Findings:
- - Detect gaps in incoming sequence
- - Implement then test lost packet resending
- - Probably worth special casing self->self comms. (later)
- */
- #if 1
- void usage()
- {
- printf("USAGE: uttest [options] iprange\n");
- printf("Options are:\n");
- printf(
- "--jumboFrames\n"
- "--udpLocalWriteSocketSize nn\n"
- "--udpRetryBusySenders nn\n"
- "--maxPacketsPerSender nn\n"
- "--udpQueueSize nn\n"
- "--udpRTSTimeout nn\n"
- "--udpSnifferEnabled 0|1\n"
- "--udpTraceCategories nn\n"
- "--udpTraceLevel nn\n"
- "--dontSendToSelf\n"
- "--sendSize nnMB\n"
- "--rawSpeedTest\n"
- "--rawBufferSize nn\n"
- );
- ExitModuleObjects();
- releaseAtoms();
- exit(1);
- }
- const char *multicastIP = "239.1.1.1";
- unsigned udpNumQs = 1;
- unsigned numNodes;
- unsigned myIndex;
- unsigned udpQueueSize = 100;
- unsigned maxPacketsPerSender = 0x7fffffff;
- bool sending = true;
- bool receiving = true;
- bool dontSendToSelf = false;
- offset_t sendSize = 0;
- bool doSortSimulator = false;
- bool simpleSequential = true;
- float slowNodeSkew = 1.0;
- unsigned numSortSlaves = 50;
- bool doRawTest = false;
- unsigned rawBufferSize = 1024;
- unsigned rowSize = 100; // MORE - take params
- bool variableRows = true;
- unsigned maxMessageSize=10000;
- bool incrementRowSize = variableRows;
- unsigned maxRowSize=5000;
- unsigned minRowSize=1;
- bool readRows = true;
- struct TestHeader
- {
- unsigned sequence;
- unsigned nodeIndex;
- };
- class SendAsFastAsPossible : public Thread
- {
- ISocket *flowSocket;
- unsigned size;
- static SpinLock ratelock;
- static unsigned lastReport;
- static unsigned totalSent;
- public:
- SendAsFastAsPossible(unsigned port, unsigned sendSize)
- {
- SocketEndpoint ep(port, getNodeAddress(0));
- flowSocket = ISocket::udp_connect(ep);
- size = sendSize;
- }
- virtual int run()
- {
- byte *buffer = new byte[65535];
- if (!lastReport)
- lastReport = msTick();
- for (;;)
- {
- unsigned lim = (1024 * 1024) / size;
- for (unsigned i = 0; i < lim; i++)
- flowSocket->write(buffer, size);
- SpinBlock b(ratelock);
- totalSent += lim * size;
- unsigned now = msTick();
- unsigned elapsed = now - lastReport;
- if (elapsed >= 1000)
- {
- unsigned rate = (((__int64) totalSent) * 1000) / elapsed;
- DBGLOG("%.2f Mbytes/sec", ((float) rate) / (1024*1024));
- totalSent = 0;
- lastReport = now;
- }
- }
- throwUnexpected(); // loop never terminates, but some compilers complain about missing return without this line
- }
- };
- SpinLock SendAsFastAsPossible::ratelock;
- unsigned SendAsFastAsPossible::lastReport = 0;
- unsigned SendAsFastAsPossible::totalSent = 0;
- class Receiver : public Thread
- {
- bool running;
- Semaphore started;
- offset_t allReceived;
- CriticalSection arsect;
- public:
- Receiver() : Thread("Receiver")
- {
- running = false;
- allReceived = 0;
- }
- virtual void start()
- {
- Thread::start();
- started.wait();
- }
- void stop(offset_t torecv)
- {
- CriticalBlock block(arsect);
- while (allReceived<torecv) {
- PROGLOG("Waiting for Receiver (%" I64F "d remaining)",torecv-allReceived);
- CriticalUnblock unblock(arsect);
- Sleep(1000);
- }
- running = false;
- }
- virtual int run()
- {
- Owned<IReceiveManager> rcvMgr = createReceiveManager(7000, 7001, 7002, 7003, multicastIP, udpQueueSize, maxPacketsPerSender, myIndex);
- Owned<roxiemem::IRowManager> rowMgr = roxiemem::createRowManager(0, NULL, queryDummyContextLogger(), NULL);
- Owned<IMessageCollator> collator = rcvMgr->createMessageCollator(rowMgr, 1);
- unsigned lastReport = 0;
- unsigned receivedTotal = 0;
- unsigned *received = new unsigned[numNodes];
- unsigned *lastSequence = new unsigned[numNodes];
- for (unsigned i = 0; i < numNodes; i++)
- {
- received[i] = 0;
- lastSequence[i] = 0;
- }
- running = true;
- started.signal();
- unsigned start = msTick();
- unsigned lastReceived = start;
- while (running)
- {
- bool dummy;
- Owned <IMessageResult> result = collator->getNextResult(2000, dummy);
- if (result)
- {
- if (!lastReport)
- {
- start = msTick(); // get first message free....
- lastReport = msTick();
- }
- // process data here....
- unsigned headerLength;
- const TestHeader *header = (const TestHeader *) result->getMessageHeader(headerLength);
- assertex (headerLength == sizeof(TestHeader) && header->nodeIndex < numNodes);
- if (header->sequence > lastSequence[header->nodeIndex])
- {
- if (header->sequence != lastSequence[header->nodeIndex]+1)
- DBGLOG("Missing messages %u-%u from node %u", lastSequence[header->nodeIndex]+1, header->sequence-1, header->nodeIndex);
- lastSequence[header->nodeIndex] = header->sequence;
- }
- else
- {
- DBGLOG("Out-of-sequence message %u from node %u", header->sequence, header->nodeIndex);
- if (header->sequence+256 < lastSequence[header->nodeIndex])
- {
- DBGLOG("Assuming receiver restart");
- lastSequence[header->nodeIndex] = header->sequence;
- }
- }
- if (readRows)
- {
- Owned<IMessageUnpackCursor> cursor = result->getCursor(rowMgr);
- for (;;)
- {
- if (variableRows)
- {
- RecordLengthType *rowlen = (RecordLengthType *) cursor->getNext(sizeof(RecordLengthType));
- if (rowlen)
- {
- const void *data = cursor->getNext(*rowlen);
- // MORE - check contents
- received[header->nodeIndex] += *rowlen;
- receivedTotal += *rowlen;
- allReceived += *rowlen;
- ReleaseRoxieRow(rowlen);
- ReleaseRoxieRow(data);
- }
- else
- break;
- }
- else
- UNIMPLEMENTED;
- }
- }
- }
- lastReceived = msTick();
- if (lastReport && (lastReceived - lastReport > 10000))
- {
- lastReport = lastReceived;
- DBGLOG("Received %u bytes total, rate = %.2f MB/s", receivedTotal, ((double)receivedTotal)/10485760.0);
- for (unsigned i = 0; i < numNodes; i++)
- {
- DBGLOG(" %u bytes from node %u", received[i], i);
- received[i] = 0;
- }
- receivedTotal = 0;
- }
- }
- {
- CriticalBlock block(arsect);
- double totalRate = (((double)allReceived)/1048576.0)/((lastReceived-start)/1000.0);
- DBGLOG("Node %d All Received %" I64F "d bytes, rate = %.2f MB/s", myIndex, allReceived, totalRate);
- }
- rcvMgr->detachCollator(collator);
- delete [] received;
- delete [] lastSequence;
- return 0;
- }
- };
- void testNxN()
- {
- if (maxPacketsPerSender > udpQueueSize)
- maxPacketsPerSender = udpQueueSize;
- Owned <ISendManager> sendMgr = createSendManager(7000, 7001, 7002, 7003, multicastIP, 100, udpNumQs, NULL, myIndex);
- Receiver receiver;
- IMessagePacker **packers = new IMessagePacker *[numNodes];
- unsigned *sequences = new unsigned[numNodes];
- for (unsigned i = 0; i < numNodes; i++)
- {
- sequences[i] = 1;
- packers[i] = NULL;
- }
- DBGLOG("Ready to start");
- if (receiving)
- {
- receiver.start();
- if (numNodes > 1)
- Sleep(5000);
- }
- offset_t sentTotal = 0;
- if (sending)
- {
- Sleep(5000); // Give receivers a fighting chance
- unsigned dest = 0;
- unsigned start = msTick();
- unsigned last = start;
- if (sendSize)
- {
- unsigned n = dontSendToSelf ? numNodes -1 : numNodes;
- sendSize /= 100*n;
- sendSize *= 100*n;
- }
- for (;;)
- {
- do {
- dest++;
- if (dest == numNodes)
- dest = 0;
- }
- while (dontSendToSelf&&(dest==myIndex));
- if (!packers[dest])
- {
- TestHeader t = {sequences[dest], myIndex};
- packers[dest] = sendMgr->createMessagePacker(1, sequences[dest], &t, sizeof(t), dest, 0);
- }
- void *row = packers[dest]->getBuffer(rowSize, variableRows);
- memset(row, 0xaa, rowSize);
- packers[dest]->putBuffer(row, rowSize, variableRows);
- if (packers[dest]->size() > maxMessageSize)
- {
- unsigned now = msTick();
- if (now-last>10000) {
- DBGLOG("Sent %" I64F "d bytes total, rate = %.2f MB/s", sentTotal, (((double)sentTotal)/1048576.0)/((now-start)/1000.0));
- last = now;
- }
- packers[dest]->flush(true);
- packers[dest]->Release();
- packers[dest] = NULL;
- sequences[dest]++;
- }
- sentTotal += rowSize;
- if (incrementRowSize)
- {
- rowSize++;
- if (rowSize==maxRowSize)
- rowSize = minRowSize;
- }
- if (sendSize && sentTotal>=sendSize)
- break;
- }
- for (unsigned i = 0; i < numNodes; i++)
- {
- if (packers[i])
- packers[i]->flush(true);
- }
- DBGLOG("Node %d All Sent %" I64F "d bytes total, rate = %.2f MB/s", myIndex, sentTotal, (((double)sentTotal)/1048576.0)/((msTick()-start)/1000.0));
- while (!sendMgr->allDone())
- {
- DBGLOG("Node %d waiting for queued data to be flushed", myIndex);
- Sleep(1000);
- }
- DBGLOG("All data sent");
- }
- else if (receiving)
- Sleep(3000000);
- receiver.stop(sentTotal);
- receiver.join();
- Sleep(10*1000); // possible receivers may request retries so should leave senders alive for a bit
- for (unsigned ii = 0; ii < numNodes; ii++)
- {
- if (packers[ii])
- packers[ii]->Release();
- }
- delete [] packers;
- delete [] sequences;
- }
- void rawSendTest()
- {
- unsigned startPort = 7002;
- for (unsigned senders = 0; senders < 10; senders++)
- {
- DBGLOG("Starting sender %d on port %d", senders+1, startPort);
- SendAsFastAsPossible *newSender = new SendAsFastAsPossible(startPort++, rawBufferSize);
- newSender->start();
- Sleep(10000);
- }
- }
- class SortMaster
- {
- unsigned __int64 receivingMask;
- unsigned __int64 sendingMask;
- unsigned __int64 *nodeData;
- int numSlaves;
- CriticalSection masterCrit;
- int *nextNode;
- Semaphore *receiveSem;
- int *receivesCompleted;
- public:
- SortMaster(int _numSlaves)
- {
- receivingMask = 0;
- sendingMask = 0;
- numSlaves = _numSlaves;
- nodeData = new unsigned __int64[numSlaves];
- for (int i = 0; i < numSlaves; i++)
- nodeData[i] = ((unsigned __int64) 1) << i;
- nextNode = NULL;
- receiveSem = NULL;
- receivesCompleted = NULL;
- if (simpleSequential)
- {
- nextNode = new int[numSlaves];
- receiveSem = new Semaphore[numSlaves];
- for (int i = 0; i < numSlaves; i++)
- {
- nextNode[i] = i+1;
- receiveSem[i].signal();
- }
- nextNode[numSlaves-1] = 0;
- }
- else
- {
- receivesCompleted = new int[numSlaves];
- for (int i = 0; i < numSlaves; i++)
- receivesCompleted[i] = 0;
- }
- }
- ~SortMaster()
- {
- delete [] nodeData;
- delete [] nextNode;
- delete [] receiveSem;
- delete [] receivesCompleted;
- }
- int requestToSend(int sendingNode)
- {
- int receivingNode = 0;
- if (simpleSequential)
- {
- {
- CriticalBlock b(masterCrit);
- receivingNode = nextNode[sendingNode];
- nextNode[sendingNode]++;
- if (nextNode[sendingNode] >= numSlaves)
- nextNode[sendingNode] = 0;
- sendingMask |= (((unsigned __int64) 1)<<sendingNode);
- receivingMask |= (((unsigned __int64) 1)<<receivingNode);
- }
- receiveSem[receivingNode].wait();
- }
- else
- {
- // Nigel's algorithm - find a node that this slave hasn't yet sent to, which is idle, which is furthest behind, and (if poss) which is not currently sending
- int bestScore = -1;
- while (bestScore == -1)
- {
- CriticalBlock b(masterCrit);
- for (int i = 0; i < numSlaves; i++)
- {
- if ((nodeData[sendingNode] & (((unsigned __int64) 1) << i)) == 0)
- {
- // I still need to send to this node
- if ((receivingMask & (((unsigned __int64) 1) << i)) == 0)
- {
- // and it is idle...
- int score = 2*receivesCompleted[i];
- if ((sendingMask & (((unsigned __int64) 1) << i)) != 0)
- score++;
- if (score < bestScore || bestScore==-1)
- {
- bestScore = score;
- receivingNode = i;
- }
- }
- }
- }
- if (bestScore == -1)
- {
- CriticalUnblock b(masterCrit);
- Sleep(10); // MORE - should wait until something changes then retry
- }
- else
- {
- sendingMask |= (((unsigned __int64) 1)<<sendingNode);
- receivingMask |= (((unsigned __int64) 1)<<receivingNode);
- nodeData[sendingNode] |= (((unsigned __int64) 1) << receivingNode);
- break;
- }
- }
- }
- return receivingNode;
- };
- void noteTransferStart(int sendingNode, int receivingNode)
- {
- // Nothing here at the moment - we set the masks in requestToSend to ensure atomicity
- };
- void noteTransferEnd(int sendingNode, int receivingNode)
- {
- CriticalBlock b(masterCrit);
- sendingMask &= ~(((unsigned __int64) 1)<<sendingNode);
- receivingMask &= ~(((unsigned __int64) 1)<<receivingNode);
- if (simpleSequential)
- receiveSem[receivingNode].signal();
- else
- receivesCompleted[receivingNode]++;
- };
- inline int queryNumSlaves() const
- {
- return numSlaves;
- }
- };
- class SortSlave : public Thread
- {
- SortMaster *master;
- int myIdx;
- int slavesDone;
- int dataSize(int targetIndex)
- {
- if (targetIndex==0 && slowNodeSkew)
- return (int)(1000 * slowNodeSkew);
- else
- return 1000;
- }
- public:
- SortSlave()
- {
- master = NULL;
- myIdx = -1;
- slavesDone = 0;
- }
- void init(SortMaster *_master, unsigned _myIdx)
- {
- master = _master;
- myIdx = _myIdx;
- slavesDone = 0;
- }
- void sendTo(unsigned datasize, unsigned slaveIdx)
- {
- assert(slaveIdx != myIdx);
- DBGLOG("Node %d sending %d bytes to node %d", myIdx, datasize, slaveIdx);
- master->noteTransferStart(myIdx, slaveIdx);
- Sleep(datasize);
- master->noteTransferEnd(myIdx, slaveIdx);
- }
- virtual int run()
- {
- while (slavesDone < (master->queryNumSlaves() - 1))
- {
- unsigned nextDest = master->requestToSend(myIdx);
- sendTo(dataSize(nextDest), nextDest);
- slavesDone++;
- }
- return 0;
- }
- };
- void sortSimulator()
- {
- // test out various ideas for determining the order in which n nodes should exchange data....
- SortMaster master(numSortSlaves);
- SortSlave *slaves = new SortSlave[numSortSlaves];
- unsigned start = msTick();
- for (unsigned i = 0; i < numSortSlaves; i++)
- {
- slaves[i].init(&master, i);
- slaves[i].start();
- }
- for (unsigned j = 0; j < numSortSlaves; j++)
- {
- slaves[j].join();
- }
- unsigned elapsed = msTick() - start;
- DBGLOG("Complete in %d.%03d seconds", elapsed / 1000, elapsed % 1000);
- DBGLOG("sequential=%d, skewFactor %f", (int) simpleSequential, slowNodeSkew);
- delete[] slaves;
- }
- int main(int argc, char * argv[] )
- {
- InitModuleObjects();
- if (argc < 2)
- usage();
- strdup("Make sure leak checking is working");
- queryStderrLogMsgHandler()->setMessageFields(MSGFIELD_time | MSGFIELD_thread | MSGFIELD_prefix);
- {
- Owned<IComponentLogFileCreator> lf = createComponentLogFileCreator("UDPTRANSPORT");
- lf->setCreateAliasFile(false);
- lf->setRolling(false);
- lf->setAppend(false);
- lf->setMaxDetail(TopDetail);
- lf->setMsgFields(MSGFIELD_STANDARD);
- lf->beginLogging();
- }
- StringBuffer cmdline;
- int c;
- for (c = 0; c < argc; c++) {
- if (c)
- cmdline.append(' ');
- cmdline.append(argv[c]);
- }
- DBGLOG("%s",cmdline.str());
- // queryLogMsgManager()->enterQueueingMode();
- // queryLogMsgManager()->setQueueDroppingLimit(512, 32);
- udpRequestToSendTimeout = 5000;
- for (c = 1; c < argc; c++)
- {
- const char *ip = argv[c];
- const char *dash = strchr(ip, '-');
- if (dash==ip)
- {
- if (strcmp(ip, "--udpQueueSize")==0)
- {
- c++;
- if (c==argc || !isdigit(*argv[c]))
- usage();
- udpQueueSize = atoi(argv[c]);
- }
- if (strcmp(ip, "--udpRTSTimeout")==0)
- {
- c++;
- if (c==argc || !isdigit(*argv[c]))
- usage();
- udpRequestToSendTimeout = atoi(argv[c]);
- }
- else if (strcmp(ip, "--jumboFrames")==0)
- {
- roxiemem::setDataAlignmentSize(0x2000);
- }
- else if (strcmp(ip, "--rawSpeedTest")==0)
- {
- doRawTest = true;
- }
- else if (strcmp(ip, "--udpLocalWriteSocketSize")==0)
- {
- c++;
- if (c==argc)
- usage();
- udpLocalWriteSocketSize = atoi(argv[c]);
- }
- else if (strcmp(ip, "--udpRetryBusySenders")==0)
- {
- c++;
- if (c==argc)
- usage();
- udpRetryBusySenders = atoi(argv[c]);
- }
- else if (strcmp(ip, "--maxPacketsPerSender")==0)
- {
- c++;
- if (c==argc)
- usage();
- maxPacketsPerSender = atoi(argv[c]);
- }
- else if (strcmp(ip, "--udpSnifferEnabled")==0)
- {
- c++;
- if (c==argc)
- usage();
- udpSnifferEnabled = atoi(argv[c]) != 0;
- }
- else if (strcmp(ip, "--udpTraceLevel")==0)
- {
- c++;
- if (c==argc)
- usage();
- udpTraceLevel = atoi(argv[c]);
- }
- else if (strcmp(ip, "--udpTraceCategories")==0)
- {
- c++;
- if (c==argc)
- usage();
- udpTraceCategories = atoi(argv[c]);
- }
- else if (strcmp(ip, "--dontSendToSelf")==0)
- {
- dontSendToSelf = true;
- }
- else if (strcmp(ip, "--sortSimulator")==0)
- {
- doSortSimulator = true;
- }
- else if (strcmp(ip, "--sendSize")==0)
- {
- c++;
- if (c==argc)
- usage();
- sendSize = (offset_t)atoi(argv[c])*(offset_t)0x100000;
- }
- else if (strcmp(ip, "--rawBufferSize")==0)
- {
- c++;
- if (c==argc)
- usage();
- rawBufferSize = atoi(argv[c]);
- }
- else
- usage();
- }
- else if (dash && isdigit(dash[1]) && dash>ip && isdigit(dash[-1]))
- {
- const char *startrange = dash-1;
- while (isdigit(startrange[-1]))
- startrange--;
- char *endptr;
- unsigned firstnum = atoi(startrange);
- unsigned lastnum = strtol(dash+1, &endptr, 10);
- while (firstnum <= lastnum)
- {
- StringBuffer ipstr;
- ipstr.append(startrange - ip, ip).append(firstnum).append(endptr);
- unsigned nodeIdx = addRoxieNode(ipstr.str());
- const IpAddress &nodeIP = getNodeAddress(nodeIdx);
- nodeIP.getIpText(ipstr.clear());
- printf("Node %u is %s\n", nodeIdx, ipstr.str());
- firstnum++;
- }
- }
- else
- {
- StringBuffer ipstr;
- unsigned nodeIdx = addRoxieNode(ip);
- const IpAddress &nodeIP = getNodeAddress(nodeIdx);
- nodeIP.getIpText(ipstr.clear());
- printf("Node %u is %s\n", nodeIdx, ipstr.str());
- }
- }
- if (doRawTest)
- rawSendTest();
- else if (doSortSimulator)
- sortSimulator();
- else
- {
- numNodes = getNumNodes();
- myIndex = addRoxieNode(GetCachedHostName());
- if (myIndex >= numNodes)
- {
- printf("ERROR: my ip does not appear to be in range\n");
- usage();
- }
- roxiemem::setTotalMemoryLimit(false, true, false, 1048576000, 0, NULL, NULL);
- testNxN();
- roxiemem::releaseRoxieHeap();
- }
- ExitModuleObjects();
- releaseAtoms();
- return 0;
- }
- #else
- // Ole's old test - look at sometime!
- #define MAX_PACKERS 10
- #define MAX_PACKETS 20
- struct PackerInfo {
- unsigned numPackets;
- unsigned packetsSizes[MAX_PACKETS];
- };
- char *progName;
- bool noendwait = false;
- unsigned thisTrace = 1;
- unsigned modeType = 0;
- unsigned myIndex = 0;
- unsigned destA = 0;
- unsigned destB = 0;
- char *multiCast = "239.1.1.2";
- unsigned udpNumQs = 3;
- unsigned numPackers = 2;
- unsigned numSizes = 4;
- unsigned numSends = 10;
- unsigned initSize = 100;
- unsigned sizeMulti = 2;
- unsigned delayPackers = 0;
- unsigned getUnpackerTimeout = 10000;
- unsigned packerHdrSize = 32;
- struct PackerInfo packersInfo[MAX_PACKERS]; // list of packers info, if used. each is alist of sizes (msgs).
- unsigned numPackersInfo = 0;
- void usage(char *err = NULL)
- {
- if (err) fprintf(stderr, "Usage Error: %s\n", err);
- fprintf(stderr, "Usage: %s [ -send [-destA IP] [-destB IP] ] [-receive]\n", progName);
- fprintf(stderr, " [-multiCast IP] [-udpTimeout sec] [-udpMaxTimeouts val]\n");
- fprintf(stderr, " [-udpNumQs val] [-udpQsPriority val] [-packerHdrSize val]\n");
- fprintf(stderr, " [-numPackers val] [-numSizes val] [-numSends val]\n");
- fprintf(stderr, " [-initSize val] [-sizeMulti val] [-delayPackers msec]\n");
- fprintf(stderr, " [-udpTrace val] [-thisTrace val] [-noendwait]\n");
- fprintf(stderr, " [-send] : Sets the mode to sender mode (i.e roxie slave like) <default dual mode>\n");
- fprintf(stderr, " [-receive] : Sets the mode to receiver mode (i.e roxie server like) <default dual mode>\n");
- fprintf(stderr, " [-destA IP] : Sets the sender destination ip address to IP (i.e roxie server IP) <default to local host>\n");
- fprintf(stderr, " [-destB IP] : Sets the sender second destination ip address to IP <default no sec dest>\n");
- fprintf(stderr, " [-multiCast IP] : Sets the sniffer multicast ip address to IP <default %s>\n", multiCast);
- fprintf(stderr, " [-udpTimeout msec] : Sets the sender udpRequestToSendTimeout value <default %i>\n", udpRequestToSendTimeout);
- fprintf(stderr, " [-udpMaxTimeouts val]: Sets the sender udpMaxRetryTimedoutReqs value <default %i>\n", udpMaxRetryTimedoutReqs);
- fprintf(stderr, " [-udpNumQs val] : Sets the sender's number of output queues <default %i>\n", udpNumQs);
- fprintf(stderr, " [-udpQsPriority val] : Sets the sender's output queues priority udpQsPriority <default %i>\n", udpOutQsPriority);
- fprintf(stderr, " [-packerHdrSize val] : Sets the packers header size (like RoxieHeader) <default %i>\n", packerHdrSize);
- fprintf(stderr, " [-numPackers val] : Sets the number of packers/unpackers to create/expect <default %i>\n", numPackers);
- fprintf(stderr, " [-packers val vale .]: Sets a packer specific packet sizes, this option can be repeated as many packers as needed\n");
- fprintf(stderr, " [-numSizes val] : Sets the number of packet data sizes to try sending/receiving <default %i>\n", numSizes);
- fprintf(stderr, " [-numSends val]] : Sets the number of msgs per size per packer to send <default %i>\n", numSends);
- fprintf(stderr, " [-initSize val] : Sets the size of the first msg(s) per packer to send <default %i>\n", initSize);
- fprintf(stderr, " [-sizeMulti val] : Sets the multiplier value of the size of subsequent msgs per packer <default %i>\n", sizeMulti);
- fprintf(stderr, " [-delayPackers msec] : Sets the delay value between sent packers (simulate roxie server/slave) <default %i>\n", delayPackers);
- fprintf(stderr, " [-getUnpackerTimeout msec] : Sets the timeout value used when calling getNextUnpacker <default %i>\n", getUnpackerTimeout);
- fprintf(stderr, " [-thisTrace val] : Sets the trace level of this program <default %i>\n", thisTrace);
- fprintf(stderr, " [-udpTrace val] : Sets the udpTraveLevel value <default %i>\n", udpTraceLevel);
- fprintf(stderr,"\n\nEnter q to terminate program : ");
- fflush(stdout);
- char tmpBuf[10]; scanf("%s", tmpBuf);
- exit(1);
- }
- #define SND_MODE_BIT 0x01
- #define RCV_MODE_BIT 0x02
- int main(int argc, char * argv[] )
- {
- InitModuleObjects();
- progName = argv[0];
- destA = myIndex = addRoxieNode(GetCachedHostName());
- udpRequestToSendTimeout = 5000;
- udpMaxRetryTimedoutReqs = 3;
- udpOutQsPriority = 5;
- udpTraceLevel = 1;
- setTotalMemoryLimit(104857600);
- char errBuff[100];
- for (int i = 1; i < argc; i++)
- {
- if (*argv[i] == '-')
- {
- if(stricmp(argv[i]+1,"send")==0)
- modeType |= SND_MODE_BIT;
- else if(stricmp(argv[i]+1,"receive")==0)
- modeType |= RCV_MODE_BIT;
- else if(stricmp(argv[i]+1,"noendwait")==0)
- noendwait = true;
- else if(stricmp(argv[i]+1,"destA")==0)
- {
- if (i+1 < argc)
- {
- destA = addRoxieNode(argv[++i]);
- }
- else
- {
- sprintf(errBuff,"Missing IP address after \"%s\"", argv[i]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"destB")==0)
- {
- if (i+1 < argc)
- {
- destB = addRoxieNode(argv[++i]);
- }
- else
- {
- sprintf(errBuff,"Missing IP address after \"%s\"", argv[i]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"multiCast")==0)
- {
- if (++i < argc)
- {
- multiCast = argv[i];
- }
- else
- {
- sprintf(errBuff,"Missing IP address after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"udpTimeout")==0)
- {
- if (++i < argc)
- {
- udpRequestToSendTimeout = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"udpMaxTimeouts")==0)
- {
- if (++i < argc)
- {
- udpMaxRetryTimedoutReqs = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"udpNumQs")==0)
- {
- if (++i < argc)
- {
- udpNumQs = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"udpQsPriority")==0)
- {
- if (++i < argc)
- {
- udpOutQsPriority = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
-
- else if(stricmp(argv[i]+1,"packerHdrSize")==0)
- {
- if (++i < argc)
- {
- packerHdrSize = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"numPackers")==0)
- {
- if (++i < argc)
- {
- numPackers = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"packer")==0)
- {
- if (numPackersInfo >= MAX_PACKERS)
- {
- sprintf(errBuff,"Too many packers are listed - max=%i", MAX_PACKERS);
- usage(errBuff);
- }
- struct PackerInfo &packerInfo = packersInfo[numPackersInfo];
- packerInfo.numPackets = 0;
- while ((++i < argc) && (*argv[i] != '-'))
- {
- if (packerInfo.numPackets >= MAX_PACKETS)
- {
- sprintf(errBuff,"Too many packets in packer - max=%i", MAX_PACKETS);
- usage(errBuff);
- }
- packerInfo.packetsSizes[packerInfo.numPackets] = atoi(argv[i]);
- packerInfo.numPackets++;
- }
- if (packerInfo.numPackets == 0)
- {
- sprintf(errBuff,"Missing packer packets info");
- usage(errBuff);
- }
- --i;
- numPackersInfo++;
- }
- else if(stricmp(argv[i]+1,"numSizes")==0)
- {
- if (++i < argc)
- {
- numSizes = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"numSends")==0)
- {
- if (++i < argc)
- {
- numSends = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"initSize")==0)
- {
- if (++i < argc)
- {
- initSize = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"sizeMulti")==0)
- {
- if (++i < argc)
- {
- sizeMulti = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"delayPackers")==0)
- {
- if (++i < argc)
- {
- delayPackers = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"getUnpackerTimeout")==0)
- {
- if (++i < argc)
- {
- getUnpackerTimeout = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"thisTrace")==0)
- {
- if (++i < argc)
- {
- thisTrace = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else if(stricmp(argv[i]+1,"udpTrace")==0)
- {
- if (++i < argc)
- {
- udpTraceLevel = atoi(argv[i]);
- }
- else
- {
- sprintf(errBuff,"Missing value after \"%s\"", argv[i-1]);
- usage(errBuff);
- }
- }
- else
- {
- sprintf(errBuff,"Invalid argument option \"%s\"", argv[i]);
- usage(errBuff);
- }
- }
- else
- {
- sprintf(errBuff,"Argument option \"%s\" missing \"-\" ", argv[i]);
- usage(errBuff);
- }
- }
- // default is daul mode (send and receive)
- if (!modeType) modeType = SND_MODE_BIT | RCV_MODE_BIT;
-
- IReceiveManager *rcvMgr = NULL;
- IRowManager *rowMgr = NULL;
- IMessageCollator *msgCollA = NULL;
- IMessageCollator *msgCollB = NULL;
- ISendManager *sendMgr = NULL;
- if (modeType & RCV_MODE_BIT)
- {
- rcvMgr = createReceiveManager(7000, 7001, 7002, 7003, multiCast, 100, 0x7fffffff, myIndex);
- rowMgr = createRowManager(0, NULL, queryDummyContextLogger(), NULL);
- msgCollA = rcvMgr->createMessageCollator(rowMgr, 100);
- if (destB)
- {
- msgCollB = rcvMgr->createMessageCollator(rowMgr, 200);
- }
- Sleep(1000);
- }
- if (modeType & SND_MODE_BIT)
- {
- sendMgr = createSendManager(7000, 7001, 7002, 7003, multiCast, 100, udpNumQs, 100, NULL, myIndex);
- Sleep(5000);
- char locBuff[100000];
- for (unsigned packerNum=0; packerNum < numPackers; packerNum++)
- {
- unsigned totalSize = 0;
- char packAHdr[100];
- char packBHdr[100];
- sprintf(packAHdr,"helloA%i", packerNum);
- if (thisTrace)
- printf("Creating packer - hdrLen=%i header %s\n", packerHdrSize, packAHdr);
- IMessagePacker *msgPackA = sendMgr->createMessagePacker(100, 0, packAHdr, packerHdrSize, destA, 1);
- IMessagePacker *msgPackB = NULL;
- if (destB)
- {
- sprintf(packBHdr,"helloB%i", packerNum);
- if (thisTrace)
- printf("Creating packer - hdrLen=%i header %s\n", packerHdrSize, packBHdr);
- msgPackB = sendMgr->createMessagePacker(200, 0, packBHdr, packerHdrSize, destB, 0);
- }
- unsigned buffSize = initSize;
- int pkIx = packerNum;
- int nmSizes = numSizes;
- if (numPackersInfo)
- {
- if (pkIx >= numPackersInfo) pkIx %= numPackersInfo;
- nmSizes = packersInfo[pkIx].numPackets;
- }
- for (unsigned sizeNum=0; sizeNum < nmSizes; sizeNum++, buffSize *= sizeMulti)
- {
- unsigned nmSends = numSends;
- if (numPackersInfo)
- {
- nmSends = 1;
- buffSize = packersInfo[pkIx].packetsSizes[sizeNum];
- }
- for (unsigned sendNum=0; sendNum < nmSends; sendNum++)
- {
- sprintf(locBuff,"size=%i num=%i multi=%i packer=%i hello world",
- buffSize, sendNum, sizeNum, packerNum);
- if (thisTrace > 1)
- printf("Sending data : %s\n", locBuff);
- char *transBuff = (char*) msgPackA->getBuffer(buffSize, false);
- strncpy(transBuff, locBuff, buffSize);
- msgPackA->putBuffer(transBuff, buffSize, false);
-
- if (msgPackB)
- {
- transBuff = (char*) msgPackB->getBuffer(buffSize, false);
- strncpy(transBuff, locBuff, buffSize);
- msgPackB->putBuffer(transBuff, buffSize, false);
- }
- totalSize += buffSize;
- }
- }
- msgPackA->flush(true);
- msgPackA->Release();
- if (thisTrace)
- printf("Packer %s total data size = %i\n", packAHdr, totalSize);
- if (msgPackB)
- {
- msgPackB->flush(true);
- msgPackB->Release();
- if (thisTrace)
- printf("Packer %s total data size = \n", packBHdr, totalSize);
- }
- if (delayPackers) Sleep(delayPackers);
- }
-
- while(!sendMgr->allDone()) Sleep(50);
- }
- if (modeType & RCV_MODE_BIT)
- {
- for (unsigned unpackerNum=0; unpackerNum < numPackers; unpackerNum++)
- {
- bool anyActivity_a;
- bool anyActivity_b;
- IMessageResult *resultA = msgCollA->getNextResult(getUnpackerTimeout, anyActivity_a);
- if (!resultA)
- {
- printf("timeout waiting on msgCollA->getNextResult(%i,..)\n", getUnpackerTimeout);
- }
- IMessageResult *resultB = NULL;
- if (msgCollB)
- {
- resultB = msgCollB->getNextResult(getUnpackerTimeout, anyActivity_b);
- if (!resultB)
- {
- printf("timeout waiting on msgCollB->getNextResult(%i,..)\n", getUnpackerTimeout);
- }
- }
- unsigned len;
- const void *hdr;
- char locBuff[100000];
- if (resultA)
- {
- hdr = resultA->getMessageHeader(len);
- if (thisTrace)
- printf("Got unpacker - hdrLen=%i header %s\n", len, hdr);
- }
- if (resultB)
- {
- hdr = resultB->getMessageHeader(len);
- if (thisTrace)
- printf("Got unpacker - hdrLen=%i header \"%s\"\n", len, hdr);
- }
-
- if (!resultA && resultB)
- {
- resultA = resultB;
- resultB = NULL;
- }
- if (!resultA) continue;
- Owned<IMessageUnpackCursor> unpackA = resultA->getCursor(rowMgr);
- Owned<IMessageUnpackCursor> unpackB = resultB ? resultB->getCursor(rowMgr) : NULL;
- unsigned totalSize = 0;
- unsigned buffSize = initSize;
- if (unpackerNum)
- {
- int size;
- if (thisTrace)
- printf("Calling getNext() for all data available in packer \"%s\"\n", hdr);
- void * p= unpackA->getNext(0x0ffffffff,&size);
- totalSize += size;
- }
- else
- {
- if (thisTrace)
- printf("Calling getNext() with diff sizes for packer \"%s\"\n", hdr);
- buffSize = initSize;
- int pkIx = unpackerNum;
- int nmSizes = numSizes;
- if (numPackersInfo)
- {
- if (pkIx >= numPackersInfo) pkIx %= numPackersInfo;
- nmSizes = packersInfo[pkIx].numPackets;
- }
- for (unsigned sizeNum=0; sizeNum < nmSizes; sizeNum++, buffSize *= sizeMulti)
- {
- unsigned nmSends = numSends;
- if (numPackersInfo)
- {
- nmSends = 1;
- buffSize = packersInfo[pkIx].packetsSizes[sizeNum];
- }
- for (unsigned sendNum=0; sendNum < nmSends; sendNum++)
- {
- int size;
- void *transBuff= unpackA->getNext(buffSize, &size);
- if (!transBuff)
- {
- if (thisTrace > 1)
- printf("end of data\n");
- }
- else {
- totalSize += size;
- memcpy(locBuff, transBuff, size);
- locBuff[size]=0;
- if (thisTrace > 1)
- printf("Received (for size=%i num=%i multi=%i unpacker=%i) data : %s\n",
- buffSize, sendNum, sizeNum, unpackerNum, locBuff);
- }
- }
- }
- }
-
- if (thisTrace)
- printf("Unpacker %s total data size = %i\n", hdr, totalSize);
- buffSize=initSize;
- if (thisTrace > 1)
- printf("Trying to read more than written\n");
- void *transBuff = unpackA->getNext(buffSize);
- if (!transBuff)
- {
- if (thisTrace > 1)
- printf("OK: Could not read more than written\n");
- }
- else
- {
- memcpy(locBuff, transBuff, buffSize);
- locBuff[buffSize]=0;
- printf("WARNING: read more than written: (%s)\n", locBuff);
- }
- printf("\n\n\n");
-
- unpackA->Release();
- if (unpackB) unpackB->Release();
- }
- }
- if (msgCollA)
- {
- rcvMgr->detachCollator(msgCollA);
- msgCollA->Release();
- }
- if (msgCollB)
- {
- rcvMgr->detachCollator(msgCollB);
- msgCollB->Release();
- }
- if (sendMgr) sendMgr->Release();
- if (rcvMgr) rcvMgr->Release();
- if (!noendwait)
- {
- printf("\n\nEnter q to terminate program : ");
- scanf("%s", errBuff);
- }
- return 0;
- }
- #endif
|