12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997 |
- #include "cmPrefix.h"
- #include "cmGlobal.h"
-
- #include "cmRpt.h"
- #include "cmErr.h"
- #include "cmMem.h"
- #include "cmMallocDebug.h"
-
- #include "cmThread.h"
-
- #include <pthread.h>
- #include <unistd.h> // usleep
- //#include <atomic_ops.h>
-
- cmThreadH_t cmThreadNullHandle = {NULL};
-
- enum
- {
- kDoExitThFl = 0x01,
- kDoPauseThFl = 0x02,
- kDoRunThFl = 0x04
- };
-
- typedef struct
- {
- cmErr_t err;
- cmThreadFunc_t funcPtr;
- pthread_t pthreadH;
- cmThStateId_t state;
- void* funcParam;
- unsigned doFlags;
- unsigned pauseMicroSecs;
- unsigned waitMicroSecs;
- } cmThThread_t;
-
- cmThRC_t _cmThError( cmErr_t* err, cmThRC_t rc, int sysErr, const char* fmt, ... )
- {
- va_list vl;
- va_start(vl,fmt);
- cmErrVSysMsg(err,rc,sysErr,fmt,vl);
- va_end(vl);
- return rc;
- }
-
- void _cmThThreadCleanUpCallback(void* t)
- {
- ((cmThThread_t*)t)->state = kExitedThId;
- }
-
- void* _cmThThreadCallback(void* param)
- {
- cmThThread_t* t = (cmThThread_t*)param;
-
- // set a clean up handler - this will be called when the
- // thread terminates unexpectedly or pthread_cleanup_pop() is called.
- pthread_cleanup_push(_cmThThreadCleanUpCallback,t);
-
- while( cmIsFlag(t->doFlags,kDoExitThFl) == false )
- {
-
- if( t->state == kPausedThId )
- {
- usleep( t->pauseMicroSecs );
-
- if( cmIsFlag(t->doFlags,kDoRunThFl) )
- {
- t->doFlags = cmClrFlag(t->doFlags,kDoRunThFl);
- t->state = kRunningThId;
- }
- }
- else
- {
-
- if( t->funcPtr(t->funcParam)==false )
- break;
-
- if( cmIsFlag(t->doFlags,kDoPauseThFl) )
- {
- t->doFlags = cmClrFlag(t->doFlags,kDoPauseThFl);
- t->state = kPausedThId;
- }
- }
- }
-
- pthread_cleanup_pop(1);
-
- pthread_exit(NULL);
-
- return t;
- }
-
- cmThThread_t* _cmThThreadFromHandle( cmThreadH_t h )
- {
- cmThThread_t* tp = (cmThThread_t*)h.h;
-
- assert(tp != NULL);
-
- return tp->state==kNotInitThId ? NULL : tp;
- }
-
- cmThRC_t _cmThWaitForState( cmThThread_t* t, unsigned stateId )
- {
- unsigned waitTimeMicroSecs = 0;
-
- while( t->state != stateId && waitTimeMicroSecs < t->waitMicroSecs )
- {
- //usleep( t->waitMicroSecs );
- usleep( 15000 );
- waitTimeMicroSecs += 15000; //t->waitMicroSecs;
- }
-
- return t->state==stateId ? kOkThRC : kTimeOutThRC;
- }
-
- cmThRC_t cmThreadCreate( cmThreadH_t* hPtr, cmThreadFunc_t funcPtr, void* funcParam, cmRpt_t* rpt )
- {
- //pthread_attr_t attr;
- cmThRC_t rc = kOkThRC;
- cmThThread_t* tp = cmMemAllocZ( cmThThread_t, 1 );
- int sysErr;
-
- cmErrSetup(&tp->err,rpt,"Thread");
-
- tp->funcPtr = funcPtr;
- tp->funcParam = funcParam;
- tp->state = kPausedThId;
- tp->doFlags = 0;
- tp->pauseMicroSecs = 50000;
- tp->waitMicroSecs = 1000000;
-
- if((sysErr = pthread_create(&tp->pthreadH,NULL,_cmThThreadCallback, (void*)tp )) != 0 )
- {
- tp->state = kNotInitThId;
- rc = _cmThError(&tp->err,kCreateFailThRC,sysErr,"Thread create failed.");
- }
-
- hPtr->h = tp;
-
- return rc;
- }
-
-
- cmThRC_t cmThreadDestroy( cmThreadH_t* hPtr )
- {
- cmThRC_t rc = kOkThRC;
-
- if( hPtr==NULL || cmThreadIsValid(*hPtr)==false )
- return rc;
-
- cmThThread_t* t = _cmThThreadFromHandle(*hPtr );
-
- if( t == NULL )
- return kInvalidHandleThRC;
-
- // tell the thread to exit
- t->doFlags = cmSetFlag(t->doFlags,kDoExitThFl);
-
- // wait for the thread to exit and then deallocate the thread object
- if((rc = _cmThWaitForState(t,kExitedThId)) == kOkThRC )
- {
- cmMemFree(t);
- hPtr->h = NULL;
- }
- else
- {
- rc = _cmThError(&t->err,rc,0,"Thread timed out waiting for destroy.");
- }
-
- return rc;
- }
-
- cmThRC_t cmThreadPause( cmThreadH_t h, unsigned cmdFlags )
- {
- cmThRC_t rc = kOkThRC;
- bool pauseFl = cmIsFlag(cmdFlags,kPauseThFl);
- bool waitFl = cmIsFlag(cmdFlags,kWaitThFl);
-
- cmThThread_t* t = _cmThThreadFromHandle(h);
- unsigned waitId;
-
- if( t == NULL )
- return kInvalidHandleThRC;
-
- bool isPausedFl = t->state == kPausedThId;
-
- if( isPausedFl == pauseFl )
- return kOkThRC;
-
- if( pauseFl )
- {
- t->doFlags = cmSetFlag(t->doFlags,kDoPauseThFl);
- waitId = kPausedThId;
- }
- else
- {
- t->doFlags = cmSetFlag(t->doFlags,kDoRunThFl);
- waitId = kRunningThId;
- }
-
- if( waitFl )
- rc = _cmThWaitForState(t,waitId);
-
- if( rc != kOkThRC )
- _cmThError(&t->err,rc,0,"Thread timed out waiting for '%s'.", pauseFl ? "pause" : "un-pause");
-
- return rc;
- }
-
- cmThStateId_t cmThreadState( cmThreadH_t h )
- {
- cmThThread_t* tp = _cmThThreadFromHandle(h);
-
- if( tp == NULL )
- return kNotInitThId;
-
- return tp->state;
- }
-
- bool cmThreadIsValid( cmThreadH_t h )
- { return h.h != NULL; }
-
- unsigned cmThreadPauseTimeOutMicros( cmThreadH_t h )
- {
- cmThThread_t* tp = _cmThThreadFromHandle(h);
- return tp->pauseMicroSecs;
- }
-
- void cmThreadSetPauseTimeOutMicros( cmThreadH_t h, unsigned usecs )
- {
- cmThThread_t* tp = _cmThThreadFromHandle(h);
- tp->pauseMicroSecs = usecs;
- }
-
- unsigned cmThreadWaitTimeOutMicros( cmThreadH_t h )
- {
- cmThThread_t* tp = _cmThThreadFromHandle(h);
- return tp->waitMicroSecs;
- }
-
- void cmThreadSetWaitTimeOutMicros( cmThreadH_t h, unsigned usecs )
- {
- cmThThread_t* tp = _cmThThreadFromHandle(h);
- tp->waitMicroSecs = usecs;
- }
-
-
- bool _cmThreadTestCb( void* p )
- {
- unsigned* ip = (unsigned*)p;
- ip[0]++;
- return true;
- }
-
-
- void cmThreadTest(cmRpt_t* rpt)
- {
- cmThreadH_t th0;
- unsigned val = 0;
-
- if( cmThreadCreate(&th0,_cmThreadTestCb,&val,rpt) == kOkThRC )
- {
- if( cmThreadPause(th0,0) != kOkThRC )
- {
- cmRptPrintf(rpt,"Thread start failed.\n");
- return;
- }
-
- char c = 0;
-
- cmRptPrintf(rpt,"o=print p=pause s=state q=quit\n");
-
- while( c != 'q' )
- {
-
- c = (char)fgetc(stdin);
- fflush(stdin);
-
- switch(c)
- {
- case 'o':
- cmRptPrintf(rpt,"val: 0x%x\n",val);
- break;
-
- case 's':
- cmRptPrintf(rpt,"state=%i\n",cmThreadState(th0));
- break;
-
- case 'p':
- {
- cmRC_t rc;
- if( cmThreadState(th0) == kPausedThId )
- rc = cmThreadPause(th0,kWaitThFl);
- else
- rc = cmThreadPause(th0,kPauseThFl|kWaitThFl);
-
- if( rc == kOkThRC )
- cmRptPrintf(rpt,"new state:%i\n", cmThreadState(th0));
- else
- cmRptPrintf(rpt,"cmThreadPause() failed.");
-
- }
- break;
-
- case 'q':
- break;
-
- //default:
- //cmRptPrintf(rpt,"Unknown:%c\n",c);
-
- }
- }
-
- if( cmThreadDestroy(&th0) != kOkThRC )
- cmRptPrintf(rpt,"Thread destroy failed.\n");
- }
-
- }
-
- //-----------------------------------------------------------------------------
- //-----------------------------------------------------------------------------
- //-----------------------------------------------------------------------------
-
- typedef struct
- {
- cmErr_t err;
- pthread_mutex_t mutex;
- pthread_cond_t cvar;
- } cmThreadMutex_t;
-
- cmThreadMutexH_t kThreadMutexNULL = {NULL};
-
-
- cmThreadMutex_t* _cmThreadMutexFromHandle( cmThreadMutexH_t h )
- {
- cmThreadMutex_t* p = (cmThreadMutex_t*)h.h;
- assert(p != NULL);
- return p;
- }
-
- cmThRC_t cmThreadMutexCreate( cmThreadMutexH_t* hPtr, cmRpt_t* rpt )
- {
- int sysErr;
- cmThreadMutex_t* p = cmMemAllocZ( cmThreadMutex_t, 1 );
-
- cmErrSetup(&p->err,rpt,"Thread Mutex");
-
- if((sysErr = pthread_mutex_init(&p->mutex,NULL)) != 0 )
- return _cmThError(&p->err,kCreateFailThRC,sysErr,"Thread mutex create failed.");
-
- if((sysErr = pthread_cond_init(&p->cvar,NULL)) != 0 )
- return _cmThError(&p->err,kCreateFailThRC,sysErr,"Thread Condition var. create failed.");
-
- hPtr->h = p;
- return kOkThRC;
- }
-
- cmThRC_t cmThreadMutexDestroy( cmThreadMutexH_t* hPtr )
- {
- int sysErr;
- cmThreadMutex_t* p = _cmThreadMutexFromHandle(*hPtr);
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- if((sysErr = pthread_cond_destroy(&p->cvar)) != 0)
- return _cmThError(&p->err,kDestroyFailThRC,sysErr,"Thread condition var. destroy failed.");
-
- if((sysErr = pthread_mutex_destroy(&p->mutex)) != 0)
- return _cmThError(&p->err,kDestroyFailThRC,sysErr,"Thread mutex destroy failed.");
-
-
- cmMemFree(p);
- hPtr->h = NULL;
-
- return kOkThRC;
- }
-
- cmThRC_t cmThreadMutexTryLock( cmThreadMutexH_t h, bool* lockFlPtr )
- {
- cmThreadMutex_t* p = _cmThreadMutexFromHandle(h);
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- int sysErr = pthread_mutex_trylock(&p->mutex);
-
- switch(sysErr)
- {
- case EBUSY:
- *lockFlPtr = false;
- break;
-
- case 0:
- *lockFlPtr = true;
- break;
-
- default:
- return _cmThError(&p->err,kLockFailThRC,sysErr,"Thread mutex try-lock failed.");;
- }
-
- return kOkThRC;
- }
-
- cmThRC_t cmThreadMutexLock( cmThreadMutexH_t h )
- {
- cmThreadMutex_t* p = _cmThreadMutexFromHandle(h);
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- int sysErr = pthread_mutex_lock(&p->mutex);
-
- if( sysErr == 0 )
- return kOkThRC;
-
- return _cmThError(&p->err,kLockFailThRC,sysErr,"Thread mutex lock failed.");
- }
-
- cmThRC_t cmThreadMutexUnlock( cmThreadMutexH_t h )
- {
- cmThreadMutex_t* p = _cmThreadMutexFromHandle(h);
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- int sysErr = pthread_mutex_unlock(&p->mutex);
-
- if( sysErr == 0 )
- return kOkThRC;
-
- return _cmThError(&p->err,kUnlockFailThRC,sysErr,"Thread mutex unlock failed.");
- }
-
- bool cmThreadMutexIsValid( cmThreadMutexH_t h )
- { return h.h != NULL; }
-
- cmThRC_t cmThreadMutexWaitOnCondVar( cmThreadMutexH_t h, bool lockFl )
- {
- cmThreadMutex_t* p = _cmThreadMutexFromHandle(h);
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- int sysErr;
-
- if( lockFl )
- if( (sysErr=pthread_mutex_lock(&p->mutex)) != 0 )
- _cmThError(&p->err,kLockFailThRC,sysErr,"Thread lock failed on cond. var. wait.");
-
- if((sysErr = pthread_cond_wait(&p->cvar,&p->mutex)) != 0 )
- _cmThError(&p->err,kCVarWaitFailThRC,sysErr,"Thread cond. var. wait failed.");
-
- return kOkThRC;
- }
-
- cmThRC_t cmThreadMutexSignalCondVar( cmThreadMutexH_t h )
- {
- int sysErr;
- cmThreadMutex_t* p = _cmThreadMutexFromHandle(h);
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- if((sysErr = pthread_cond_signal(&p->cvar)) != 0 )
- return _cmThError(&p->err,kCVarSignalFailThRC,sysErr,"Thread cond. var. signal failed.");
-
- return kOkThRC;
- }
-
-
- //-----------------------------------------------------------------------------
- //-----------------------------------------------------------------------------
- //-----------------------------------------------------------------------------
-
- cmTsQueueH_t cmTsQueueNullHandle = { NULL };
-
- enum { cmTsQueueBufCnt = 2 };
-
-
- typedef struct
- {
- unsigned allocCnt; // count of bytes allocated for the buffer
- unsigned fullCnt; // count of bytes used in the buffer
- char* basePtr; // base of buffer memory
- unsigned* msgPtr; // pointer to first msg
- unsigned msgCnt;
- } cmTsQueueBuf;
-
- typedef struct
- {
- cmThreadMutexH_t mutexH;
- cmTsQueueBuf bufArray[cmTsQueueBufCnt];
- unsigned inBufIdx;
- unsigned outBufIdx;
- char* memPtr;
- cmTsQueueCb_t cbFunc;
- void* userCbPtr;
- } cmTsQueue_t;
-
- cmTsQueue_t* _cmTsQueueFromHandle( cmTsQueueH_t h )
- {
- cmTsQueue_t* p = h.h;
- assert(p != NULL);
- return p;
- }
-
- cmThRC_t _cmTsQueueDestroy( cmTsQueue_t* p )
- {
- cmThRC_t rc;
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- if( p->mutexH.h != NULL )
- if((rc = cmThreadMutexDestroy(&p->mutexH)) != kOkThRC )
- return rc;
-
- if( p->memPtr != NULL )
- cmMemPtrFree(&p->memPtr);
-
-
- cmMemPtrFree(&p);
-
-
- return kOkThRC;
- }
-
-
- cmThRC_t cmTsQueueCreate( cmTsQueueH_t* hPtr, unsigned bufByteCnt, cmTsQueueCb_t cbFunc, void* userCbPtr, cmRpt_t* rpt )
- {
- cmTsQueue_t* p = cmMemAllocZ( cmTsQueue_t, 1 );
- unsigned i;
-
- if( cmThreadMutexCreate(&p->mutexH,rpt) != kOkThRC )
- goto errLabel;
-
- p->memPtr = cmMemAllocZ( char, bufByteCnt*cmTsQueueBufCnt );
- p->outBufIdx = 0;
- p->inBufIdx = 1;
- p->cbFunc = cbFunc;
- p->userCbPtr = userCbPtr;
-
- for(i=0; i<cmTsQueueBufCnt; ++i)
- {
- p->bufArray[i].allocCnt = bufByteCnt;
- p->bufArray[i].fullCnt = 0;
- p->bufArray[i].basePtr = p->memPtr + (i*bufByteCnt);
- p->bufArray[i].msgPtr = NULL;
- p->bufArray[i].msgCnt = 0;
- }
-
- hPtr->h = p;
-
- return kOkThRC;
-
- errLabel:
-
- _cmTsQueueDestroy(p);
-
- return kCreateFailThRC;
- }
-
-
- cmThRC_t cmTsQueueDestroy( cmTsQueueH_t* hPtr )
- {
- cmThRC_t rc = kOkThRC;
-
- if( (hPtr != NULL) && cmTsQueueIsValid(*hPtr))
- if((rc = _cmTsQueueDestroy(_cmTsQueueFromHandle(*hPtr))) == kOkThRC )
- hPtr->h = NULL;
-
- return rc;
- }
-
- cmThRC_t cmTsQueueSetCallback( cmTsQueueH_t h, cmTsQueueCb_t cbFunc, void* cbArg )
- {
- cmTsQueue_t* p = _cmTsQueueFromHandle(h);
- p->cbFunc = cbFunc;
- p->userCbPtr = cbArg;
- return kOkThRC;
- }
-
- unsigned cmTsQueueAllocByteCount( cmTsQueueH_t h )
- {
- cmTsQueue_t* p = _cmTsQueueFromHandle(h);
- unsigned n = 0;
-
- if( cmThreadMutexLock(p->mutexH) == kOkThRC )
- {
- n = p->bufArray[ p->inBufIdx ].allocCnt;
- cmThreadMutexUnlock(p->mutexH);
- }
- return n;
- }
-
- unsigned cmTsQueueAvailByteCount( cmTsQueueH_t h )
- {
- cmTsQueue_t* p = _cmTsQueueFromHandle(h);
- unsigned n = 0;
-
- if(cmThreadMutexLock(p->mutexH) == kOkThRC )
- {
- n = p->bufArray[ p->inBufIdx ].allocCnt - p->bufArray[ p->inBufIdx].fullCnt;
- cmThreadMutexUnlock(p->mutexH);
- }
- return n;
- }
-
- cmThRC_t _cmTsQueueEnqueueMsg( cmTsQueueH_t h, const void* msgPtrArray[], unsigned msgByteCntArray[], unsigned arrayCnt )
- {
- cmThRC_t rc;
- cmTsQueue_t* p = _cmTsQueueFromHandle(h);
-
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- // lock the mutex
- if((rc = cmThreadMutexLock(p->mutexH)) == kOkThRC )
- {
-
- cmTsQueueBuf* b = p->bufArray + p->inBufIdx; // ptr to buf recd
- const char* ep = b->basePtr + b->allocCnt; // end of buf data space
- unsigned *mp = (unsigned*)(b->basePtr + b->fullCnt); // ptr to size of new msg space
- char* dp = (char*)(mp+1); // ptr to data area of new msg space
- unsigned ttlByteCnt = 0; // track size of msg data
- unsigned i = 0;
-
- // get the total size of the msg
- for(i=0; i<arrayCnt; ++i)
- ttlByteCnt += msgByteCntArray[i];
-
- // if the msg is too big for the queue buf
- if( dp + ttlByteCnt > ep )
- rc = kBufFullThRC;
- else
- {
- // for each segment of the incoming msg
- for(i=0; i<arrayCnt; ++i)
- {
- // get the size of the segment
- unsigned n = msgByteCntArray[i];
-
- // copy in the segment
- memcpy(dp,msgPtrArray[i],n);
-
- dp += n; //
- }
-
- assert(dp <= ep );
-
- // write the size ofthe msg into the buffer
- *mp = ttlByteCnt;
-
- // update the pointer to the first msg
- if( b->msgPtr == NULL )
- b->msgPtr = mp;
-
- // track the count of msgs in this buffer
- ++b->msgCnt;
-
- // update fullCnt last since dequeue uses fullCnt to
- // notice that a msg may be waiting
- b->fullCnt += sizeof(unsigned) + ttlByteCnt;
- }
-
- cmThreadMutexUnlock(p->mutexH);
- }
-
- return rc;
- }
-
- cmThRC_t cmTsQueueEnqueueSegMsg( cmTsQueueH_t h, const void* msgPtrArray[], unsigned msgByteCntArray[], unsigned arrayCnt )
- { return _cmTsQueueEnqueueMsg(h,msgPtrArray,msgByteCntArray,arrayCnt); }
-
- cmThRC_t cmTsQueueEnqueueMsg( cmTsQueueH_t h, const void* dataPtr, unsigned byteCnt )
- {
- const void* msgPtrArray[] = { dataPtr };
- unsigned msgByteCntArray[] = { byteCnt };
-
- return _cmTsQueueEnqueueMsg(h,msgPtrArray,msgByteCntArray,1);
- }
-
- cmThRC_t cmTsQueueEnqueueIdMsg( cmTsQueueH_t h, unsigned id, const void* dataPtr, unsigned byteCnt )
- {
- const void* msgPtrArray[] = { &id, dataPtr };
- unsigned msgByteCntArray[] = { sizeof(id), byteCnt };
-
- return _cmTsQueueEnqueueMsg(h,msgPtrArray,msgByteCntArray,2);
- }
-
- cmThRC_t _cmTsQueueDequeueMsg( cmTsQueue_t* p, void* retBuf, unsigned refBufByteCnt )
- {
- cmTsQueueBuf* b = p->bufArray + p->outBufIdx;
-
- // if the output buffer is empty - there is nothing to do
- if( b->fullCnt == 0 )
- return kBufEmptyThRC;
-
- assert( b->msgPtr != NULL );
-
- // get the output msg size and data
- unsigned msgByteCnt = *b->msgPtr;
- char* msgDataPtr = (char*)(b->msgPtr + 1);
-
- // transmit the msg via a callback
- if( retBuf == NULL && p->cbFunc != NULL )
- p->cbFunc(p->userCbPtr,msgByteCnt,msgDataPtr);
- else
- {
- // retBuf may be NULL if the func is being used by cmTsQueueDequeueByteCount()
- if( retBuf == NULL || msgByteCnt > refBufByteCnt )
- return kBufTooSmallThRC;
-
- // copy the msg to a buffer
- if( retBuf != NULL )
- memcpy(retBuf,msgDataPtr,msgByteCnt);
-
- }
-
- // update the buffer
- b->fullCnt -= sizeof(unsigned) + msgByteCnt;
- b->msgPtr = (unsigned*)(msgDataPtr + msgByteCnt);
- --(b->msgCnt);
-
- if( b->fullCnt == 0 )
- {
- assert(b->msgCnt == 0);
- b->msgPtr = NULL;
- }
-
- return kOkThRC;
-
- }
-
- cmThRC_t cmTsQueueDequeueMsg( cmTsQueueH_t h, void* retBuf, unsigned refBufByteCnt )
- {
- cmThRC_t rc;
- cmTsQueue_t* p = _cmTsQueueFromHandle(h);
-
- if( p == NULL )
- return kInvalidHandleThRC;
-
- // dequeue the next msg from the current output buffer
- if((rc =_cmTsQueueDequeueMsg( p, retBuf, refBufByteCnt )) != kBufEmptyThRC )
- return rc;
-
- // the current output buffer was empty
-
- cmTsQueueBuf* b = p->bufArray + p->inBufIdx;
-
- // if the input buffer has msg's ...
- if( b->fullCnt > 0 )
- {
- bool lockFl = false;
-
- // ...attempt to lock the mutex ...
- if( (cmThreadMutexTryLock(p->mutexH,&lockFl) == kOkThRC) && lockFl )
- {
- // ... swap the input and the output buffers ...
- unsigned tmp = p->inBufIdx;
- p->inBufIdx = p->outBufIdx;
- p->outBufIdx = tmp;
-
- // .. unlock the mutex
- cmThreadMutexUnlock(p->mutexH);
-
- // ... and dequeue the first msg from the new output buffer
- rc = _cmTsQueueDequeueMsg( p, retBuf, refBufByteCnt );
- }
-
- }
-
- return rc;
-
- }
-
- bool cmTsQueueMsgWaiting( cmTsQueueH_t h )
- {
- cmTsQueue_t* p = _cmTsQueueFromHandle(h);
-
- if( p == NULL )
- return false;
-
- if( p->bufArray[p->outBufIdx].fullCnt )
- return true;
-
- return p->bufArray[p->inBufIdx].fullCnt > 0;
- }
-
- unsigned cmTsQueueDequeueMsgByteCount( cmTsQueueH_t h )
- {
- cmTsQueue_t* p = _cmTsQueueFromHandle(h);
-
- if( p == NULL )
- return 0;
-
- // if output msgs are available then the msgPtr points to the size of the msg
- if( p->bufArray[p->outBufIdx].fullCnt )
- return *(p->bufArray[p->outBufIdx].msgPtr);
-
- // no msgs are waiting in the output buffer
-
- // force the buffers to swap - returns kBufEmptyThRC if there are
- // still no msgs waiting after the swap (the input buf was also empty)
- if( cmTsQueueDequeueMsg(h,NULL,0) == kBufTooSmallThRC )
- {
- // the buffers swapped so there must be msg waiting
- assert( p->bufArray[p->outBufIdx].fullCnt );
-
- return *(p->bufArray[p->outBufIdx].msgPtr);
- }
-
- return 0;
- }
-
- bool cmTsQueueIsValid( cmTsQueueH_t h )
- { return h.h != NULL; }
-
-
-
- //--------------------------------------------------------------------------------------------------
- //--------------------------------------------------------------------------------------------------
- //--------------------------------------------------------------------------------------------------
- #ifdef NOT_DEF
- enum { kThBufCnt=2 };
-
- typedef struct
- {
- char* buf;
- volatile unsigned ii;
- volatile unsigned oi;
- } cmThBuf_t;
-
- typedef struct
- {
- cmErr_t err;
- cmThBuf_t a[kThBufCnt];
- volatile unsigned ibi;
- unsigned bn;
- cmTsQueueCb_t cbFunc;
- void* cbArg;
- } cmTs1p1c_t;
-
- cmTs1p1c_t* _cmTs1p1cHandleToPtr( cmTs1p1cH_t h )
- {
- cmTs1p1c_t* p = (cmTs1p1c_t*)h.h;
- assert( p != NULL );
- return p;
- }
-
- cmThRC_t _cmTs1p1cDestroy( cmTs1p1c_t* p )
- {
- unsigned i;
- for(i=0; i<kThBufCnt; ++i)
- cmMemFree(p->a[i].buf);
-
- cmMemFree(p);
-
- return kOkThRC;
- }
-
- cmThRC_t cmTs1p1cCreate( cmTs1p1cH_t* hPtr, unsigned bufByteCnt, cmTsQueueCb_t cbFunc, void* cbArg, cmRpt_t* rpt )
- {
- cmThRC_t rc;
- if((rc = cmTs1p1cDestroy(hPtr)) != kOkThRC )
- return rc;
-
- unsigned i;
- cmTs1p1c_t* p = cmMemAllocZ(cmTs1p1c_t,1);
- cmErrSetup(&p->err,rpt,"TS 1p1c Queue");
- for(i=0; i<kThBufCnt; ++i)
- {
- p->a[i].buf = cmMemAllocZ(char,bufByteCnt);
- p->a[i].ii = 0;
- p->a[i].oi = bufByteCnt;
- }
-
- p->ibi = 0;
- p->bn = bufByteCnt;
- p->cbFunc = cbFunc;
- p->cbArg = cbArg;
-
- hPtr->h = p;
-
- return rc;
- }
-
- cmThRC_t cmTs1p1cDestroy( cmTs1p1cH_t* hp )
- {
- cmThRC_t rc = kOkThRC;
-
- if( hp == NULL || cmTs1p1cIsValid(*hp)==false )
- return kOkThRC;
-
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(*hp);
- if(( rc = _cmTs1p1cDestroy(p)) != kOkThRC )
- return rc;
-
- hp->h = NULL;
- return rc;
- }
-
- cmThRC_t cmTs1p1cEnqueueSegMsg( cmTs1p1cH_t h, const void* msgPtrArray[], unsigned msgByteCntArray[], unsigned arrayCnt )
- {
- cmThRC_t rc = kOkThRC;
- unsigned mn = 0;
- unsigned i;
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- cmThBuf_t* ib = p->a + p->ibi;
-
-
- // get the total count of bytes for this msg
- for(i=0; i<arrayCnt; ++i)
- mn += msgByteCntArray[i];
-
- unsigned dn = mn + sizeof(unsigned);
-
- // if the message is too big for even an empty buffer
- if( dn > p->bn )
- return cmErrMsg(&p->err,kBufFullThRC,"A msg containing %i bytes will never be able to fit in a queue with an empty size of %i bytes.",dn,p->bn);
-
- // if the msg won't fit in the current input buffer then try swapping buffers.
- if( ib->ii + dn > p->bn )
- {
- // get the current output buffer
- cmThBuf_t* ob = p->a + (p->ibi==0 ? 1 : 0);
-
- // Empty buffers will be set such that: oi==bn and ii==0.
- //
- // Note that setting ii to 0 in an output buffer is the last operation
- // performed on an empty output buffer. ii==0 is therefore the
- // signal that an output buffer can be reused for input.
-
- // if the output buffer is not empty - then an overflow occurred
- if( ob->ii != 0 )
- return cmErrMsg(&p->err,kBufFullThRC,"The msq queue cannot accept a %i byte msg into %i bytes.",dn, p->bn - ib->ii);
-
- // setup the initial output location of the new output buffer
- ib->oi = 0;
-
- // swap buffers
- p->ibi = (p->ibi + 1) % kThBufCnt;
-
- // get the new input buffer
- ib = ob;
- }
-
- // get a pointer to the base of the write location
- char* dp = ib->buf + ib->ii;
-
- // write the length of the message
- *(unsigned*)dp = mn;
- dp += sizeof(unsigned);
-
- // write the body of the message
- for(i=0; i<arrayCnt; ++i)
- {
- memcpy(dp,msgPtrArray[i],msgByteCntArray[i]);
- dp += msgByteCntArray[i];
- }
-
- // this MUST be executed last - we'll use 'dp' in the calculation
- // (even though ib->ii += dn would be more straight forward way
- // to accomplish the same thing) to prevent the optimizer from
- // moving the assignment prior to the for loop.
- ib->ii += dp - (ib->buf + ib->ii);
-
- return rc;
- }
-
- cmThRC_t cmTs1p1cEnqueueMsg( cmTsQueueH_t h, const void* dataPtr, unsigned byteCnt )
- { return cmTs1p1cEnqueueSegMsg(h,&dataPtr,&byteCnt,1); }
-
- unsigned cmTs1p1cAllocByteCount( cmTs1p1cH_t h )
- {
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- return p->bn;
- }
-
- unsigned cmTs1p1cAvailByteCount( cmTs1p1cH_t h )
- {
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- return p->bn - p->a[ p->ibi ].ii;
- }
- cmThRC_t cmTs1p1cDequeueMsg( cmTs1p1cH_t h, void* dataPtr, unsigned byteCnt )
- {
- cmThRC_t rc = kOkThRC;
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- cmThBuf_t* ob = p->a + (p->ibi == 0 ? 1 : 0);
-
- // empty buffers always are set to: oi==bn && ii==0
- if( ob->oi >= ob->ii )
- return kBufEmptyThRC;
-
- // get the size of the msg
- unsigned mn = *(unsigned*)(ob->buf + ob->oi);
-
- // increment the current output location to the msg body
- ob->oi += sizeof(unsigned);
-
- // copy or send the msg
- if( dataPtr != NULL )
- {
- if( byteCnt < mn )
- return cmErrMsg(&p->err,kBufTooSmallThRC,"The return buffer constains too few bytes (%i) to contain %i bytes.",byteCnt,mn);
-
- memcpy(dataPtr, ob->buf + ob->oi, mn);
- }
- else
- {
- p->cbFunc(p->cbArg, mn, ob->buf + ob->oi );
- }
-
- ob->oi += mn;
-
- // if we are reading correctly ob->oi should land
- // exactly on ob->ii when the buffer is empty
- assert( ob->oi <= ob->ii );
-
- // if the buffer is empty
- if( ob->oi == ob->ii )
- {
- ob->oi = p->bn; // mark the buffer as empty
- ob->ii = 0; //
- }
-
- return rc;
- }
-
-
- unsigned cmTs1p1cDequeueMsgByteCount( cmTsQueueH_t h )
- {
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- cmThBuf_t* ob = p->a + (p->ibi == 0 ? 1 : 0);
-
- // empty buffers always are set to: oi==bn && ii==0
- if( ob->oi >= ob->ii )
- return 0;
-
- // get the size of the msg
- return *(unsigned*)(ob->buf + ob->oi);
- }
-
-
- bool cmTs1p1cMsgWaiting( cmTsQueueH_t h )
- { return cmTs1p1cDequeueMsgByteCount(h) > 0; }
-
- bool cmTs1p1cIsValid( cmTs1p1cH_t h )
- { return h.h != NULL; }
- #endif
-
- //--------------------------------------------------------------------------------------------------
- //--------------------------------------------------------------------------------------------------
- //--------------------------------------------------------------------------------------------------
-
- typedef struct
- {
- volatile unsigned ii;
- cmErr_t err;
- char* buf;
- unsigned bn;
- cmTsQueueCb_t cbFunc;
- void* cbArg;
- volatile unsigned oi;
- } cmTs1p1c_t;
-
- cmTs1p1c_t* _cmTs1p1cHandleToPtr( cmTs1p1cH_t h )
- {
- cmTs1p1c_t* p = (cmTs1p1c_t*)h.h;
- assert( p != NULL );
- return p;
- }
-
- cmThRC_t cmTs1p1cCreate( cmTs1p1cH_t* hPtr, unsigned bufByteCnt, cmTsQueueCb_t cbFunc, void* cbArg, cmRpt_t* rpt )
- {
- cmThRC_t rc;
- if((rc = cmTs1p1cDestroy(hPtr)) != kOkThRC )
- return rc;
-
- cmTs1p1c_t* p = cmMemAllocZ(cmTs1p1c_t,1);
- cmErrSetup(&p->err,rpt,"1p1c Queue");
- p->buf = cmMemAllocZ(char,bufByteCnt+sizeof(unsigned));
- p->ii = 0;
- p->oi = 0;
- p->bn = bufByteCnt;
- p->cbFunc = cbFunc;
- p->cbArg = cbArg;
- hPtr->h = p;
-
- return rc;
- }
-
- cmThRC_t cmTs1p1cDestroy( cmTs1p1cH_t* hp )
- {
- cmThRC_t rc = kOkThRC;
-
- if( hp == NULL || cmTs1p1cIsValid(*hp)==false )
- return kOkThRC;
-
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(*hp);
-
- cmMemFree(p->buf);
- cmMemFree(p);
-
- hp->h = NULL;
- return rc;
- }
-
- cmThRC_t cmTs1p1cEnqueueSegMsg( cmTs1p1cH_t h, const void* msgPtrArray[], unsigned msgByteCntArray[], unsigned arrayCnt )
- {
- cmThRC_t rc = kOkThRC;
- unsigned mn = 0;
- unsigned i;
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
-
- // get the total count of bytes for this msg
- for(i=0; i<arrayCnt; ++i)
- mn += msgByteCntArray[i];
-
- int dn = mn + sizeof(unsigned);
- int oi = p->oi;
- int bi = p->ii; // 'bi' is the idx of the leftmost cell which can be written
- int en = p->bn; // 'en' is the idx of the cell just to the right of the rightmost cell that can be written
-
- // note: If 'oi' marks the rightmost location then 'en' must be set
- // one cell to the left of 'oi', because 'ii' can never be allowed to
- // advance onto 'oi' - because 'oi'=='ii' marks an empty (NOT a full)
- // queue.
- //
- // If 'bn' marks the rightmost location then 'ii' can advance onto 'bn'
- // beause the true queue length is bn+1.
-
- // if we need to wrap
- if( en-bi < dn && oi<=bi )
- {
- bi = 0;
- en = oi - 1; // note if oi==0 then en is negative - see note above re: oi==ii
- assert( p->ii>=0 && p->ii <= p->bn );
- *(unsigned*)(p->buf + p->ii) = cmInvalidIdx; // mark the wrap location
- }
-
- // if oi is between ii and bn
- if( oi > bi )
- en = oi - 1; // never allow ii to advance onto oi - see note above
-
- // if the msg won't fit
- if( en - bi < dn )
- return cmErrMsg(&p->err,kBufFullThRC,"%i consecutive bytes is not available in the queue.",dn);
-
- // set the msg byte count - the msg byte cnt precedes the msg body
- char* dp = p->buf + bi;
- *(unsigned*)dp = dn - sizeof(unsigned);
- dp += sizeof(unsigned);
-
- // copy the msg into the buffer
- for(i=0,dn=0; i<arrayCnt; ++i)
- {
- memcpy(dp,msgPtrArray[i],msgByteCntArray[i]);
- dp += msgByteCntArray[i];
- dn += msgByteCntArray[i];
- }
-
- // incrementing p->ii must occur last - the unnecessary accumulation
- // of dn in the above loop is intended to prevent this line from
- // begin moved before the copy loop.
- p->ii = bi + dn + sizeof(unsigned);
-
- assert( p->ii >= 0 && p->ii <= p->bn);
-
- return rc;
- }
-
- cmThRC_t cmTs1p1cEnqueueMsg( cmTs1p1cH_t h, const void* dataPtr, unsigned byteCnt )
- { return cmTs1p1cEnqueueSegMsg(h,&dataPtr,&byteCnt,1); }
-
- unsigned cmTs1p1cAllocByteCount( cmTs1p1cH_t h )
- {
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- return p->bn;
- }
-
- unsigned cmTs1p1cAvailByteCount( cmTs1p1cH_t h )
- {
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- unsigned oi = p->oi;
- unsigned ii = p->ii;
- return oi < ii ? p->bn - ii + oi : oi - ii;
- }
-
- unsigned _cmTs1p1cDequeueMsgByteCount( cmTs1p1c_t* p )
- {
- // if the buffer is empty
- if( p->ii == p->oi )
- return 0;
-
- // get the length of the next msg
- unsigned mn = *(unsigned*)(p->buf + p->oi);
-
- // if the msg length is cmInvalidIdx ...
- if( mn == cmInvalidIdx )
- {
- p->oi = 0; // ... wrap to buf begin and try again
-
- return _cmTs1p1cDequeueMsgByteCount(p);
- }
-
- return mn;
- }
-
- cmThRC_t cmTs1p1cDequeueMsg( cmTs1p1cH_t h, void* dataPtr, unsigned byteCnt )
- {
- cmThRC_t rc = kOkThRC;
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
-
- unsigned mn;
-
- if((mn = _cmTs1p1cDequeueMsgByteCount(p)) == 0 )
- return kBufEmptyThRC;
-
- void* mp = p->buf + p->oi + sizeof(unsigned);
-
- if( dataPtr != NULL )
- {
- if( byteCnt < mn )
- return cmErrMsg(&p->err,kBufTooSmallThRC,"The return buffer constains too few bytes (%i) to contain %i bytes.",byteCnt,mn);
-
- memcpy(dataPtr,mp,mn);
- }
- else
- {
- p->cbFunc(p->cbArg,mn,mp);
- }
-
- p->oi += mn + sizeof(unsigned);
-
- return rc;
- }
-
-
- unsigned cmTs1p1cDequeueMsgByteCount( cmTs1p1cH_t h )
- {
- cmTs1p1c_t* p = _cmTs1p1cHandleToPtr(h);
- return _cmTs1p1cDequeueMsgByteCount(p);
- }
-
- bool cmTs1p1cMsgWaiting( cmTs1p1cH_t h )
- { return cmTs1p1cDequeueMsgByteCount(h) > 0; }
-
- bool cmTs1p1cIsValid( cmTs1p1cH_t h )
- { return h.h != NULL; }
-
- //============================================================================================================================
-
-
- bool cmThIntCAS( int* addr, int old, int new )
- { return __sync_bool_compare_and_swap(addr,old,new); }
-
- bool cmThUIntCAS( unsigned* addr, unsigned old, unsigned new )
- { return __sync_bool_compare_and_swap(addr,old,new); }
-
- bool cmThFloatCAS( float* addr, float old, float new )
- { return __sync_bool_compare_and_swap((unsigned*)addr, *(unsigned*)(&old),*(unsigned*)(&new)); }
-
- bool cmThPtrCAS( void* addr, void* old, void* neww )
- {
- #ifdef OS_64
- return __sync_bool_compare_and_swap((long long*)addr, (long long)old, (long long)neww);
- #else
- return __sync_bool_compare_and_swap((int*)addr,(int)old,(int)neww);
- #endif
- }
-
-
-
- void cmThIntIncr( int* addr, int incr )
- {
- // ... could also use __sync_add_and_fetch() ...
- __sync_fetch_and_add(addr,incr);
- }
-
- void cmThUIntIncr( unsigned* addr, unsigned incr )
- {
- __sync_fetch_and_add(addr,incr);
- }
-
- void cmThFloatIncr(float* addr, float incr )
- {
-
- float old,new;
- do
- {
- old = *addr;
- new = old + incr;
- }while( cmThFloatCAS(addr,old,new)==0 );
- }
-
-
- void cmThIntDecr( int* addr, int decr )
- {
- __sync_fetch_and_sub(addr,decr);
- }
-
- void cmThUIntDecr( unsigned* addr, unsigned decr )
- {
- __sync_fetch_and_sub(addr,decr);
- }
-
- void cmThFloatDecr(float* addr, float decr )
- {
-
- float old,new;
- do
- {
- old = *addr;
- new = old - decr;
- }while( cmThFloatCAS(addr,old,new)==0 );
-
- }
-
-
-
- //============================================================================================================================
- //
- //
- typedef pthread_t cmThreadId_t;
-
- typedef struct
- {
- cmThreadId_t id; // id of this thread as returned by pthread_self()
- char* buf; // buf[bn]
- int ii; // input index
- int oi; // output index (oi==ii == empty buffer)
- } cmTsBuf_t;
-
- // msg header - which is actually written AFTER the msg it is associated with
- typedef struct cmTsHdr_str
- {
- int mn; // length of the msg
- int ai; // buffer index
- struct cmTsHdr_str* link; // pointer to next msg
- } cmTsHdr_t;
-
- typedef struct
- {
- cmErr_t err;
- int bn; // bytes per thread buffer
- cmTsBuf_t* a; // a[an] buffer array
- unsigned an; // length of a[] - one buffer per thread
- cmTsQueueCb_t cbFunc;
- void* cbArg;
- cmTsHdr_t* ilp; // prev msg hdr record
- cmTsHdr_t* olp; // prev msg hdr record (wait for olp->link to be set to go to next record)
- } cmTsMp1c_t;
-
- cmTsMp1cH_t cmTsMp1cNullHandle = cmSTATIC_NULL_HANDLE;
-
- void _cmTsMp1cPrint( cmTsMp1c_t* p )
- {
- unsigned i;
- for(i=0; i<p->an; ++i)
- printf("%2i ii:%3i oi:%3i\n",i,p->a[i].ii,p->a[i].oi);
- }
-
- cmTsMp1c_t* _cmTsMp1cHandleToPtr( cmTsMp1cH_t h )
- {
- cmTsMp1c_t* p = (cmTsMp1c_t*)h.h;
- assert(p != NULL);
- return p;
- }
-
- unsigned _cmTsMp1cBufIndex( cmTsMp1c_t* p, cmThreadId_t id )
- {
- unsigned i;
- for(i=0; i<p->an; ++i)
- if( p->a[i].id == id )
- return i;
-
- p->an = i+1;
- p->a = cmMemResizePZ(cmTsBuf_t,p->a,p->an);
- p->a[i].buf = cmMemAllocZ(char,p->bn);
- p->a[i].id = id;
-
- return i;
- }
-
- cmThRC_t cmTsMp1cDestroy( cmTsMp1cH_t* hp )
- {
- if( hp == NULL || cmTsMp1cIsValid(*hp) == false )
- return kOkThRC;
-
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(*hp);
- unsigned i;
- for(i=0; i<p->an; ++i)
- cmMemFree(p->a[i].buf);
-
- cmMemPtrFree(&p->a);
- cmMemFree(p);
-
- hp->h = NULL;
-
- return kOkThRC;
- }
-
- cmThRC_t cmTsMp1cCreate( cmTsMp1cH_t* hp, unsigned bufByteCnt, cmTsQueueCb_t cbFunc, void* cbArg, cmRpt_t* rpt )
- {
- cmThRC_t rc;
- if((rc = cmTsMp1cDestroy(hp)) != kOkThRC )
- return rc;
-
- cmTsMp1c_t* p = cmMemAllocZ(cmTsMp1c_t,1);
-
- cmErrSetup(&p->err,rpt,"TsMp1c Queue");
-
- p->a = NULL;
- p->an = 0;
- p->bn = bufByteCnt;
- p->cbFunc = cbFunc;
- p->cbArg = cbArg;
- p->ilp = NULL;
- p->olp = NULL;
-
- hp->h = p;
- return rc;
- }
-
- void cmTsMp1cSetCbFunc( cmTsMp1cH_t h, cmTsQueueCb_t cbFunc, void* cbArg )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- p->cbFunc = cbFunc;
- p->cbArg = cbArg;
- }
-
- cmTsQueueCb_t cmTsMp1cCbFunc( cmTsMp1cH_t h )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- return p->cbFunc;
- }
-
- void* cmTsMp1cCbArg( cmTsMp1cH_t h )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- return p->cbArg;
- }
-
- //#define CAS(addr,old,new) __sync_bool_compare_and_swap(addr,old,new)
- //#define CAS(addr,old,neww) cmThPtrCAS(addr,old,neww)
-
- cmThRC_t cmTsMp1cEnqueueSegMsg( cmTsMp1cH_t h, const void* msgPtrArray[], unsigned msgByteCntArray[], unsigned arrayCnt )
- {
- cmThRC_t rc = kOkThRC;
- unsigned mn = 0;
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- unsigned ai = _cmTsMp1cBufIndex( p, pthread_self() );
- cmTsBuf_t* b = p->a + ai;
- int i,bi,ei;
- cmTsHdr_t hdr;
-
- // Use a stored oi for the duration of this function.
- // b->oi may be changed by the dequeue thread but storing it here
- // at least prevents it from changing during the course of the this function.
- // Note: b->oi is only used to check for buffer full. Even if it changes
- // it would only mean that more bytes were available than calculated based
- // on the stored value. A low estimate of the actual bytes available is
- // never unsafe.
- volatile int oi = b->oi;
-
- // get the total count of bytes for this msg
- for(i=0; i<arrayCnt; ++i)
- mn += msgByteCntArray[i];
-
- // dn = count of msg bytes + count of header bytes
- int dn = mn + sizeof(hdr);
-
- // if oi is ahead of ii in the queue then we must write
- // in the area between ii and oi
- if( oi > b->ii )
- {
- ei = oi-1; // (never allow ii to equal oi (that's the empty condition))
- bi = b->ii;
- }
- else // otherwise oi is same or before ii in the queue and we have the option to wrap
- {
- // if the new msg will not fit at the end of the queue ....
- if( b->ii + dn > p->bn )
- {
- bi = 0; // ... then wrap to the beginning
- ei = oi-1; // (never allow ii to equal oi (that's the empty condition))
- }
- else
- {
- ei = p->bn; // otherwise write at the current location
- bi = b->ii;
- }
- }
-
- if( bi + dn > ei )
- return cmErrMsg(&p->err,kBufFullThRC,"%i consecutive bytes is not available in the queue.",dn);
-
- char* dp = b->buf + bi;
-
- // write the msg
- for(i=0; i<arrayCnt; ++i)
- {
- memcpy(dp,msgPtrArray[i],msgByteCntArray[i]);
- dp += msgByteCntArray[i];
- }
-
- // setup the msg header
- hdr.ai = ai;
- hdr.mn = mn;
- hdr.link = NULL;
-
- // write the msg header (following the msg body in memory)
- cmTsHdr_t* hp = (cmTsHdr_t*)dp;
- memcpy(hp,&hdr,sizeof(hdr));
-
- // increment the buffers input index
- b->ii = bi + dn;
-
- // update the link list head to point to this msg hdr
- cmTsHdr_t* old_hp, *new_hp;
- do
- {
- old_hp = p->ilp;
- new_hp = hp;
- }while(!cmThPtrCAS(&p->ilp,old_hp,new_hp));
-
- // link the prev recd to this recd
- if( old_hp != NULL )
- old_hp->link = hp;
-
- // if this is the first record written by this queue then prime the output list
- do
- {
- old_hp = p->olp;
- new_hp = hp;
-
- if( old_hp != NULL )
- break;
-
- }while(!cmThPtrCAS(&p->olp,old_hp,new_hp));
-
- return rc;
- }
-
- cmThRC_t cmTsMp1cEnqueueMsg( cmTsMp1cH_t h, const void* dataPtr, unsigned byteCnt )
- { return cmTsMp1cEnqueueSegMsg(h,&dataPtr,&byteCnt,1); }
-
- unsigned cmTsMp1cAllocByteCount( cmTsMp1cH_t h )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- return p->bn;
- }
-
- unsigned cmTsMp1cAvailByteCount( cmTsMp1cH_t h )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- unsigned ai = _cmTsMp1cBufIndex(p,pthread_self());
- const cmTsBuf_t* b = p->a + ai;
-
- if( b->oi > b->ii )
- return b->oi - b->ii - 1;
-
- return (p->bn - b->ii) + b->oi - 1;
- }
-
- unsigned _cmTsMp1cNextMsgByteCnt( cmTsMp1c_t* p )
- {
- if( p->olp == NULL )
- return 0;
-
- // if the current msg has not yet been read
- if( p->olp->mn != 0 )
- return p->olp->mn;
-
- // if the current msg has been read but a new next msg has been linked
- if( p->olp->mn == 0 && p->olp->link != NULL )
- {
- // advance the buffer output msg past the prev msg header
- char* hp = (char*)(p->olp + 1);
- p->a[p->olp->ai].oi = hp - p->a[p->olp->ai].buf;
-
- // advance msg pointer to point to the new msg header
- p->olp = p->olp->link;
-
- // return the size of the new msg
- return p->olp->mn;
- }
-
- return 0;
- }
-
- cmThRC_t cmTsMp1cDequeueMsg( cmTsMp1cH_t h, void* dataPtr, unsigned byteCnt )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
-
- // if there are no messages waiting
- if( _cmTsMp1cNextMsgByteCnt(p) == 0 )
- return kBufEmptyThRC;
-
- char* hp = (char*)p->olp;
- char* dp = hp - p->olp->mn; // the msg body is before the msg hdr
-
- if( dataPtr == NULL )
- {
- p->cbFunc(p->cbArg,p->olp->mn,dp);
- }
- else
- {
- if( p->olp->mn > byteCnt )
- return cmErrMsg(&p->err,kBufTooSmallThRC,"The return buffer constains too few bytes (%i) to contain %i bytes.",byteCnt,p->olp->mn);
-
- memcpy(dataPtr,dp,p->olp->mn);
- }
-
- // advance the buffers output index past the msg body
- p->a[p->olp->ai].oi = hp - p->a[p->olp->ai].buf;
-
- // mark the msg as read
- p->olp->mn = 0;
-
- return kOkThRC;
- }
-
-
- bool cmTsMp1cMsgWaiting( cmTsMp1cH_t h )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- return _cmTsMp1cNextMsgByteCnt(p) != 0;
- }
-
- unsigned cmTsMp1cDequeueMsgByteCount( cmTsMp1cH_t h )
- {
- cmTsMp1c_t* p = _cmTsMp1cHandleToPtr(h);
- return _cmTsMp1cNextMsgByteCnt(p);
- }
-
- bool cmTsMp1cIsValid( cmTsMp1cH_t h )
- { return h.h != NULL; }
-
-
-
- //============================================================================================================================/
- //
- // cmTsQueueTest()
- //
-
- // param recd for use by _cmTsQueueCb0() and
- // the msg record passed between the sender
- // threads and the receiver thread
- typedef struct
- {
- unsigned id;
- cmTsQueueH_t qH;
- int val;
- } _cmTsQCbParam_t;
-
-
- // Generate a random number and put it in a TS queue
- bool _cmTsQueueCb0(void* param)
- {
- _cmTsQCbParam_t* p = (_cmTsQCbParam_t*)param;
-
- p->val = rand(); // generate a random number
-
- // send the msg
- if( cmTsQueueEnqueueMsg( p->qH, p, sizeof(_cmTsQCbParam_t)) == kOkThRC )
- printf("in:%i %i\n",p->id,p->val);
- else
- printf("in error %i\n",p->id);
-
-
- usleep(100*1000);
-
- return true;
- }
-
- // Monitor a TS queue for incoming messages from _cmTsQueueCb1()
- bool _cmTsQueueCb1(void* param)
- {
- // the thread param is a ptr to the TS queue to monitor.
- cmTsQueueH_t* qp = (cmTsQueueH_t*)param;
- cmThRC_t rc;
- _cmTsQCbParam_t msg;
-
- // dequeue any waiting messages
- if((rc = cmTsQueueDequeueMsg( *qp, &msg, sizeof(msg))) == kOkThRC )
- printf("out:%i %i\n",msg.id,msg.val);
- else
- {
- if( rc != kBufEmptyThRC )
- printf("out error:%i\n", rc);
- }
-
- return true;
- }
-
- // Test the TS queue by starting sender threads (threads 0 & 1)
- // and a receiver thread (thread 2) and sending messages
- // from the sender to the receiver.
- void cmTsQueueTest( cmRpt_t* rpt )
- {
- cmThreadH_t th0=cmThreadNullHandle,th1=cmThreadNullHandle,th2=cmThreadNullHandle;
- cmTsQueueH_t q=cmTsQueueNullHandle;
- _cmTsQCbParam_t param0, param1;
-
- // create a TS Queue
- if( cmTsQueueCreate(&q,100,NULL,NULL,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 0
- param0.id = 0;
- param0.qH = q;
- if( cmThreadCreate(&th0,_cmTsQueueCb0,¶m0,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 1
- param1.id = 1;
- param1.qH = q;
- if( cmThreadCreate(&th1,_cmTsQueueCb0,¶m1,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 2
- if( cmThreadCreate(&th2,_cmTsQueueCb1,&q,rpt) != kOkThRC )
- goto errLabel;
-
- // start thread 0
- if( cmThreadPause(th0,0) != kOkThRC )
- goto errLabel;
-
- // start thread 1
- if( cmThreadPause(th1,0) != kOkThRC )
- goto errLabel;
-
- // start thread 2
- if( cmThreadPause(th2,0) != kOkThRC )
- goto errLabel;
-
- printf("any key to quit.");
- getchar();
-
-
- errLabel:
-
- if( cmThreadIsValid(th0) )
- if( cmThreadDestroy(&th0) != kOkThRC )
- printf("Error destroying thread 0\n");
-
- if( cmThreadIsValid(th1) )
- if( cmThreadDestroy(&th1) != kOkThRC )
- printf("Error destroying thread 1\n");
-
- if( cmThreadIsValid(th2) )
- if( cmThreadDestroy(&th2) != kOkThRC )
- printf("Error destroying thread 1\n");
-
- if( cmTsQueueIsValid(q) )
- if( cmTsQueueDestroy(&q) != kOkThRC )
- printf("Error destroying queue\n");
-
- }
-
-
- //============================================================================================================================/
- //
- // cmTs1p1cTest()
- //
-
- // param recd for use by _cmTsQueueCb0() and
- // the msg record passed between the sender
- // threads and the receiver thread
- typedef struct
- {
- unsigned id;
- cmTs1p1cH_t qH;
- int val;
- } _cmTs1p1cCbParam_t;
-
- cmTs1p1cH_t cmTs1p1cNullHandle = cmSTATIC_NULL_HANDLE;
-
- // Generate a random number and put it in a TS queue
- bool _cmTs1p1cCb0(void* param)
- {
- _cmTs1p1cCbParam_t* p = (_cmTs1p1cCbParam_t*)param;
-
- p->val = rand(); // generate a random number
-
- // send the msg
- if( cmTs1p1cEnqueueMsg( p->qH, p, sizeof(_cmTs1p1cCbParam_t)) == kOkThRC )
- printf("in:%i %i\n",p->id,p->val);
- else
- printf("in error %i\n",p->id);
-
- ++p->id;
-
- usleep(100*1000);
-
- return true;
- }
-
- // Monitor a TS queue for incoming messages from _cmTs1p1cCb1()
- bool _cmTs1p1cCb1(void* param)
- {
- // the thread param is a ptr to the TS queue to monitor.
- cmTs1p1cH_t* qp = (cmTs1p1cH_t*)param;
- cmThRC_t rc;
- _cmTs1p1cCbParam_t msg;
-
- // dequeue any waiting messages
- if((rc = cmTs1p1cDequeueMsg( *qp, &msg, sizeof(msg))) == kOkThRC )
- printf("out:%i %i\n",msg.id,msg.val);
- else
- {
- if( rc != kBufEmptyThRC )
- printf("out error:%i\n", rc);
- }
-
- return true;
- }
-
- // Test the TS queue by starting sender threads (threads 0 & 1)
- // and a receiver thread (thread 2) and sending messages
- // from the sender to the receiver.
- void cmTs1p1cTest( cmRpt_t* rpt )
- {
- cmThreadH_t th0=cmThreadNullHandle,th1=cmThreadNullHandle,th2=cmThreadNullHandle;
- cmTs1p1cH_t q=cmTs1p1cNullHandle;
- _cmTs1p1cCbParam_t param1;
-
- // create a TS Queue
- if( cmTs1p1cCreate(&q,28*2,NULL,NULL,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 1
- param1.id = 0;
- param1.qH = q;
- if( cmThreadCreate(&th1,_cmTs1p1cCb0,¶m1,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 2
- if( cmThreadCreate(&th2,_cmTs1p1cCb1,&q,rpt) != kOkThRC )
- goto errLabel;
-
- // start thread 1
- if( cmThreadPause(th1,0) != kOkThRC )
- goto errLabel;
-
- // start thread 2
- if( cmThreadPause(th2,0) != kOkThRC )
- goto errLabel;
-
- printf("any key to quit.");
- getchar();
-
-
- errLabel:
-
- if( cmThreadIsValid(th0) )
- if( cmThreadDestroy(&th0) != kOkThRC )
- printf("Error destroying thread 0\n");
-
- if( cmThreadIsValid(th1) )
- if( cmThreadDestroy(&th1) != kOkThRC )
- printf("Error destroying thread 1\n");
-
- if( cmThreadIsValid(th2) )
- if( cmThreadDestroy(&th2) != kOkThRC )
- printf("Error destroying thread 1\n");
-
- if( cmTs1p1cIsValid(q) )
- if( cmTs1p1cDestroy(&q) != kOkThRC )
- printf("Error destroying queue\n");
-
- }
-
-
- //============================================================================================================================/
- //
- // cmTsMp1cTest()
- //
- // param recd for use by _cmTsQueueCb0() and
- // the msg record passed between the sender
- // threads and the receiver thread
- typedef struct
- {
- unsigned id;
- cmTsMp1cH_t qH;
- int val;
- } _cmTsMp1cCbParam_t;
-
-
- unsigned _cmTsMp1cVal = 0;
-
- // Incr the global value _cmTsMp1cVal and put it in a TS queue
- bool _cmTsMp1cCb0(void* param)
- {
- _cmTsMp1cCbParam_t* p = (_cmTsMp1cCbParam_t*)param;
-
- p->val = __sync_fetch_and_add(&_cmTsMp1cVal,1);
-
- // send the msg
- if( cmTsMp1cEnqueueMsg( p->qH, p, sizeof(_cmTsMp1cCbParam_t)) == kOkThRC )
- printf("in:%i %i\n",p->id,p->val);
- else
- printf("in error %i\n",p->id);
-
- usleep(100*1000);
-
- return true;
- }
-
- // Monitor a TS queue for incoming messages from _cmTsMp1cCb1()
- bool _cmTsMp1cCb1(void* param)
- {
- // the thread param is a ptr to the TS queue to monitor.
- cmTsMp1cH_t* qp = (cmTsMp1cH_t*)param;
- cmThRC_t rc;
- _cmTsMp1cCbParam_t msg;
-
- // dequeue any waiting messages
- if((rc = cmTsMp1cDequeueMsg( *qp, &msg, sizeof(msg))) == kOkThRC )
- printf("out - cons id:%i val:%i\n",msg.id,msg.val);
- else
- {
- if( rc != kBufEmptyThRC )
- printf("out error:%i\n", rc);
- }
-
- return true;
- }
-
- // Test the TS queue by starting sender threads (threads 0 & 1)
- // and a receiver thread (thread 2) and sending messages
- // from the sender to the receiver.
- void cmTsMp1cTest( cmRpt_t* rpt )
- {
- cmThreadH_t th0=cmThreadNullHandle,th1=cmThreadNullHandle,th2=cmThreadNullHandle;
- cmTsMp1cH_t q=cmTsMp1cNullHandle;
- _cmTsMp1cCbParam_t param0, param1;
-
- // create a TS Queue
- if( cmTsMp1cCreate(&q,1000,NULL,NULL,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 0 - producer 0
- param0.id = 0;
- param0.qH = q;
- if( cmThreadCreate(&th0,_cmTsMp1cCb0,¶m0,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 1 - producer 1
- param1.id = 1;
- param1.qH = q;
- if( cmThreadCreate(&th1,_cmTsMp1cCb0,¶m1,rpt) != kOkThRC )
- goto errLabel;
-
- // create thread 2 - consumer 0
- if( cmThreadCreate(&th2,_cmTsMp1cCb1,&q,rpt) != kOkThRC )
- goto errLabel;
-
- // start thread 0
- if( cmThreadPause(th0,0) != kOkThRC )
- goto errLabel;
-
- // start thread 1
- if( cmThreadPause(th1,0) != kOkThRC )
- goto errLabel;
-
- // start thread 2
- if( cmThreadPause(th2,0) != kOkThRC )
- goto errLabel;
-
- printf("any key to quit.");
- getchar();
-
-
- errLabel:
-
- if( cmThreadIsValid(th0) )
- if( cmThreadDestroy(&th0) != kOkThRC )
- printf("Error destroying thread 0\n");
-
- if( cmThreadIsValid(th1) )
- if( cmThreadDestroy(&th1) != kOkThRC )
- printf("Error destroying thread 1\n");
-
- if( cmThreadIsValid(th2) )
- if( cmThreadDestroy(&th2) != kOkThRC )
- printf("Error destroying thread 1\n");
-
- if( cmTsMp1cIsValid(q) )
- if( cmTsMp1cDestroy(&q) != kOkThRC )
- printf("Error destroying queue\n");
-
- }
|