8 #ifndef vlShmSubscribableMessageQueue__rti___H_
9 #define vlShmSubscribableMessageQueue__rti___H_
18 #include <vlutil/vlTime.h>
21 #include <vlutil/vlNetTypes.h>
29 #ifdef DtUSE_UTILITIES_NAMESPACE
30 namespace DtUSE_UTILITIES_NAMESPACE
34 #define DtSm_QUEUE_STATE_UNINITIALIZED 0x0
35 #define DtSm_QUEUE_STATE_INITIALIZED 0xFAD2FADE
36 #define DtSm_QUEUE_STATE_SHUTDOWN 0xDEADBEEF
37 #define DtSM_INVALID_SUBSCRIBER_ID UINT_MAX
41 #define DtSMQ_FETCH_DEFAULT 0x00
42 #define DtSMQ_SKIP_MINE 0x40
43 #define DtSMQ_FETCH_ALL 0x80
45 #define DtSmSubMsgQ_theMaxSubscribers 255
101 virtual void unsubscribe(MAKRti::DtU32
id);
104 virtual bool isSubscribed(MAKRti::DtU32 subscriberId)
const;
107 virtual MAKRti::DtU32 subscriberId()
const;
110 virtual int subscriberCount()
const;
118 virtual bool isEmpty();
124 virtual void setDefaultReadBehavior(MAKRti::DtU32 flagsMask);
130 virtual bool sendMessage(
const void *message,
unsigned int size);
137 virtual void* message(
unsigned int *outSize = 0);
138 virtual void* message(
void *recvBuffer,
unsigned int *outSize);
144 virtual void resetQueue(MAKRti::DtU32 numberOfBuckets, MAKRti::DtU32 payloadSize);
147 virtual bool sm_shutdown(
bool ruthless =
false );
150 virtual volatile MAKRti::DtU32 queueState()
const;
156 virtual void initSyncVars();
160 virtual void clearSyncVars();
163 virtual std::ostream &printDataToStream(std::ostream &str)
const;
171 static unsigned int poolSize(MAKRti::DtU32 numberOfBuckets, MAKRti::DtU32 payloadSize);
181 virtual int assignSubscriberId();
184 virtual unsigned int maxSubscribers()
const;
190 virtual MAKRti::DtU32 queueHeaderSize()
const;
193 virtual MAKRti::DtU32 numberOfBuckets()
const;
196 virtual MAKRti::DtU32 payloadSize()
const;
200 virtual volatile MAKRti::DtU32 queueFront(
unsigned int sub)
const;
203 virtual volatile MAKRti::DtU32 queueBack()
const;
207 virtual volatile MAKRti::DtU32 queueBackstop()
const;
211 virtual MAKRti::DtU32 updateFront(
unsigned int sub,
unsigned int val);
215 virtual MAKRti::DtU32 updateBack(
unsigned int val);
219 virtual MAKRti::DtU32 updateBackstop(
unsigned int val);
222 virtual void setSubscribed(
unsigned int sub,
unsigned int val);
225 virtual int setSubscriberCount(MAKRti::DtU32 val);
228 virtual int incrementSubscriberCount();
231 virtual int decrementSubscriberCount();
234 virtual void setState(MAKRti::DtU32 val);
238 virtual volatile char* payloadInBucket(
unsigned int bucketNum)
const;
241 virtual volatile void * bucketHeader(
unsigned int bucketNum)
const;
244 virtual MAKRti::DtU32 bucketHeaderSize()
const;
247 virtual MAKRti::DtU32 bucketSenderId(
unsigned int bucketNum)
const;
248 virtual MAKRti::DtU32 bucketMsgSn(
unsigned int bucketNum)
const;
249 virtual MAKRti::DtU32 bucketOutstandingReadCount(
unsigned int bucketNum)
const;
250 virtual MAKRti::DtU32 bucketPayloadLen(
unsigned int bucketNum)
const;
251 virtual bool bucketIsFirst(
unsigned int bucketNum)
const;
252 virtual bool bucketIsLast(
unsigned int bucketNum)
const;
255 virtual void setSenderId(
unsigned int bucketNum, MAKRti::DtU32 sid);
256 virtual void setMsgSn(
unsigned int bucketNum, MAKRti::DtU32 seqnum);
257 virtual void setOutstandingReadCount(
unsigned int bucketNum, MAKRti::DtU32 val);
258 virtual void decrementOutstandingReadCount(
unsigned int bucketNum);
259 virtual void setPayloadLen(
unsigned int bucketNum, MAKRti::DtU32 len);
260 virtual void setFirst(
unsigned int bucketNum, MAKRti::DtU32 val);
261 virtual void setLast(
unsigned int bucketNum, MAKRti::DtU32 val);
264 virtual void lockForRead();
266 virtual void unlockForRead();
269 virtual void lockForWrite();
272 virtual void unlockForWrite();
276 virtual bool stalledSubscribers(
unsigned int bucketNum,
277 unsigned int numStalled);
281 virtual bool checkForStalls();
284 virtual pthread_rwlock_t * rdWrLock();
288 virtual pthread_mutex_t * orcMutex();
326 MAKRti::DtU32 sizeInBytes : 24;
327 MAKRti::DtU32 reserved_1 : 6;
328 MAKRti::DtU32 last : 1;
329 MAKRti::DtU32 first : 1;
350 pthread_rwlock_t rwlock;
358 #ifdef DtUSE_UTILITIES_NAMESPACE
#define DtSMQ_FETCH_DEFAULT
The following constants are intended as bit-masks to be combined with other values (some as-yet-unspe...
Definition: vlShmSubscribableMessageQueue__rti__.h:41
Implements a generic shared memory message queue for multiple subscribers.
Definition: vlShmSubscribableMessageQueue__rti__.h:66
MAKRti::DtString myShmName
Definition: vlShmSubscribableMessageQueue__rti__.h:294
MAKRti::DtU32 mySubscriberId
Definition: vlShmSubscribableMessageQueue__rti__.h:295
DtIpcRdWrLock__rti__ * myRdWrLock
Definition: vlShmSubscribableMessageQueue__rti__.h:297
#define DtSM_INVALID_SUBSCRIBER_ID
Definition: vlShmSubscribableMessageQueue__rti__.h:37
DO NOT - REPEAT - DO NOT use critical sections in lieu of named mutexes in Win32 implementations.
Definition: vlIpcMutex__rti__.h:43
Contains classes which manage and view a shared memory pool in a platform independant way...
This is an abstract base class presenting an interface for a simple message queue.
Definition: vlShmMessageQueue__rti__.h:23
unsigned int theMaxSubscribers
The maximum number of subscribers allowed in this queue.
Definition: vlShmSubscribableMessageQueue__rti__.h:309
Platform independent read-lock / write-lock.
Definition: vlIpcRdWrLock__rti__.h:38
DtSharedMemoryPoolClient__rti__ * myShmPool
!WIN32
Definition: vlShmSubscribableMessageQueue__rti__.h:293
MAKRti::DtU32 myMessagesSentCnt
Definition: vlShmSubscribableMessageQueue__rti__.h:296
std::vector< char > myReadBuffer
Definition: vlShmSubscribableMessageQueue__rti__.h:299
DtIpcMutex__rti__ * myOrcMutex
Definition: vlShmSubscribableMessageQueue__rti__.h:298
#define DtSmSubMsgQ_theMaxSubscribers
TO DO – DAA workaround for MSVC6 bug ID #241569 A bug in MSVC++6 prevents initializing a static const...
Definition: vlShmSubscribableMessageQueue__rti__.h:50
MAKRti::DtU32 myDefaultReadMask
Definition: vlShmSubscribableMessageQueue__rti__.h:300
#define DT_DLL_RTIUTIL
Definition: rtiMsConfig.h:127
DtSharedMemoryPoolClient__rti__ is used to attach to and access a previously created memory pool...
Definition: vlShmPool__rti__.h:107