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);
303 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;
373 pthread_rwlock_t rwlock;
391 bucketIndex(bucket), stalledId(id), senderId(sender), senderSn(seq), posted(tstamp)
unsigned int theMaxStallTime
The maximum time to permit a subscriber to be stalled, in seconds.
Definition: smSubscribableMessageQueue.h:334
#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
DtIpcMutex__rti__ * myQueueFullMutex
Definition: smSubscribableMessageQueue.h:321
#define DtSM_INVALID_SUBSCRIBER_ID
Definition: vlShmSubscribableMessageQueue__rti__.h:37
DtSmqEnqueueMode
Definition: smSubscribableMessageQueue.h:41
DtU32 senderId
ID of subscriber which hasn't read message.
Definition: smSubscribableMessageQueue.h:396
DtIpcCondVar__rti__ * myQueueFullCond
Definition: smSubscribableMessageQueue.h:320
DO NOT - REPEAT - DO NOT use critical sections in lieu of named mutexes in Win32 implementations.
Definition: vlIpcMutex__rti__.h:43
This must always be the last item.
Definition: smSubscribableMessageQueue.h:46
MAKRti::DtClock myTimestampClock
Definition: smSubscribableMessageQueue.h:317
#define DtRtiSubMsgQ_theMaxSubscribers
TO DO – DAA workaround for MSVC6 bug ID #241569 A bug in MSVC++6 prevents initializing a static const...
Definition: smSubscribableMessageQueue.h:54
~DtStalledSubscriber()
Definition: smSubscribableMessageQueue.h:393
DtIpcCondVar__rti__ * myNewMessageCond
Definition: smSubscribableMessageQueue.h:318
Definition: smSubscribableMessageQueue.h:44
unsigned int theMaxSubscribers
The maximum number of subscribers allowed in this queue.
Definition: smSubscribableMessageQueue.h:330
Definition: smSubscribableMessageQueue.h:43
DtStalledSubscriber keeps track of which queue subscribers have not read a bucket as well as informat...
Definition: smSubscribableMessageQueue.h:387
DtStalledSubscriberList myStalledSubs
!WIN32
Definition: smSubscribableMessageQueue.h:316
Implements a message queue for multiple subscribers.
Definition: smSubscribableMessageQueue.h:80
#define DT_DLL_LRC
Definition: rtiMsConfig.h:126
MAKRti::DtTime posted
Sender-assigned seq. no. of unread message.
Definition: smSubscribableMessageQueue.h:398
Platform independent condition variable.
Definition: vlIpcCondVar__rti__.h:39
DtStalledSubscriber(DtU32 bucket, DtU32 id, DtU32 sender, DtU32 seq, MAKRti::DtTime tstamp)
Definition: smSubscribableMessageQueue.h:390
std::list< DtStalledSubscriber * > DtStalledSubscriberList
Definition: smSubscribableMessageQueue.h:203
DtU32 stalledId
Definition: smSubscribableMessageQueue.h:395
DtU32 senderSn
ID of sender whose message is unread.
Definition: smSubscribableMessageQueue.h:397
DtIpcMutex__rti__ * myNewMessageMutex
Definition: smSubscribableMessageQueue.h:319