8 #ifndef smSubscribableMessageQueue_H_
9 #define smSubscribableMessageQueue_H_
18 #include <vlutil/vlTime.h>
19 #include <vlutil/vlUtil.h>
20 #include <vlutil/vlInetAddr.h>
23 using namespace MAKRti;
26 #define DtSMQ_PRIO_0 0x00
27 #define DtSMQ_PRIO_1 0x01
28 #define DtSMQ_PRIO_2 0x02
29 #define DtSMQ_PRIO_3 0x03
30 #define DtSMQ_PRIO_4 0x04
31 #define DtSMQ_PRIO_5 0x05
32 #define DtSMQ_PRIO_6 0x06
33 #define DtSMQ_PRIO_7 0x07
34 #define DtSMQ_PRIO_MASK 0x07
35 #define DtSMQ_MSG_RELIABLE 0x10
36 #define DtSMQ_MSG_BESTEFFORT 0x20
37 //#define DtSMQ_MSG_TBD 0x40
38 #define DtSMQ_SKIP_RELIABLE 0x10
39 #define DtSMQ_SKIP_BESTEFFORT 0x20
54 #define DtRtiSubMsgQ_theMaxSubscribers 16384
117 virtual void unsubscribeSelf();
127 virtual void unsubscribe(DtU32
id);
136 virtual bool queueIsReliable()
const;
139 virtual bool isSubscribed(DtU32 subscriberId)
const;
142 virtual int subscriberCount()
const;
149 virtual bool sendMessage(
const void *message,
unsigned int size);
150 virtual bool sendMessage(
const void *message,
unsigned int size,
unsigned int flags,
151 const MAKRti::DtInetAddr& destAddr = MAKRti::DtInetAddr::inaddrAny());
155 virtual void getRcvMsgTransportInfo(MAKRti::DtInetAddr *destAddr, DtU32 *flags);
159 virtual void getRcvMsgSenderInfo(DtU32* senderId, DtU32* msgSn =
NULL);
163 virtual int waitForMessages(MAKRti::DtTime waitTime = 0);
167 virtual int waitUntilFull(MAKRti::DtTime waitTime = 0);
173 virtual void resetQueue(DtU32 numberOfBuckets, DtU32 payloadSize);
183 virtual void initSyncVars();
187 virtual void clearSyncVars();
190 virtual std::ostream &printDataToStream(std::ostream &str)
const;
198 static unsigned int poolSize(DtU32 numberOfBuckets, DtU32 payloadSize);
207 virtual unsigned int maxSubscribers()
const;
213 DtU32 queueHeaderSize()
const;
216 virtual DtU32 numberOfBuckets()
const;
219 virtual DtU32 payloadSize()
const;
222 virtual volatile DtU32 queueState()
const;
226 virtual volatile DtU32 queueFront(
unsigned int sub)
const;
229 virtual volatile DtU32 queueBack()
const;
233 virtual volatile DtU32 queueBackstop()
const;
237 virtual DtU32 updateFront(
unsigned int sub,
unsigned int val);
241 virtual DtU32 updateBack(
unsigned int val);
245 virtual DtU32 updateBackstop(
unsigned int val);
248 virtual void setSubscribed(
unsigned int sub,
unsigned int val);
251 virtual int setSubscriberCount(DtU32 val);
254 virtual int incrementSubscriberCount();
257 virtual int decrementSubscriberCount();
260 virtual void setState(DtU32 val);
263 virtual void lockForRead();
265 virtual void unlockForRead();
268 virtual bool stalledSubscribers(
unsigned int bucketNum,
unsigned int numStalled);
272 virtual bool handleStall(
unsigned int bucketNum,
unsigned int sub);
275 virtual DtU32 bucketHeaderSize()
const;
278 virtual DtU32 bucketSenderId(
unsigned int bucketNum)
const;
279 virtual DtU32 bucketMsgSn(
unsigned int bucketNum)
const;
280 virtual DtU32 bucketOutstandingReadCount(
unsigned int bucketNum)
const;
281 virtual DtInetAddr bucketDestinationAddr(
unsigned int bucketNum)
const;
282 virtual DtU32 bucketPayloadLen(
unsigned int bucketNum)
const;
283 virtual DtU32 bucketPriority(
unsigned int bucketNum)
const;
284 virtual bool bucketIsReliable(
unsigned int bucketNum)
const;
285 virtual bool bucketIsBestEffort(
unsigned int bucketNum)
const;
286 virtual bool bucketIsFirst(
unsigned int bucketNum)
const;
287 virtual bool bucketIsLast(
unsigned int bucketNum)
const;
290 virtual void setSenderId(
unsigned int bucketNum, DtU32 sid);
291 virtual void setMsgSn(
unsigned int bucketNum, DtU32 seqnum);
292 virtual void setOutstandingReadCount(
unsigned int bucketNum, DtU32 val);
293 virtual void decrementOutstandingReadCount(
unsigned int bucketNum);
294 virtual void setDestinationAddr(
unsigned int bucketNum,
const DtInetAddr& dest);
295 virtual void setPayloadLen(
unsigned int bucketNum, DtU32 len);
296 virtual void setPriority(
unsigned int bucketNum, DtU32 priority);
297 virtual void setReliable(
unsigned int bucketNum, DtU32 val);
298 virtual void setBestEffort(
unsigned int bucketNum, DtU32 val);
299 virtual void setFirst(
unsigned int bucketNum, DtU32 val);
300 virtual void setLast(
unsigned int bucketNum, DtU32 val);
306 virtual pthread_rwlock_t * rdWrLock();
307 virtual pthread_mutex_t * orcMutex();
308 virtual pthread_cond_t * newMessageCond();
309 virtual pthread_mutex_t * newMessageCondMutex();
310 virtual pthread_cond_t * queueFullCond();
311 virtual pthread_mutex_t * queueFullCondMutex();
347 DtU32 sizeInBytes : 24;
348 DtU32 msgPriority : 3;
350 DtU32 besteffort : 1;
363 DtU32 payloadSize : 24;
365 DtU32 reserved_1 : 7;
391 bucketIndex(bucket), stalledId(id), senderId(sender), senderSn(seq), posted(tstamp)