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
50 #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();
287 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;
358 #ifdef DtUSE_UTILITIES_NAMESPACE