8 #ifndef smSubscribableMessageQueue_H_
9 #define smSubscribableMessageQueue_H_
18 #include <vlutil/vlTime.h>
19 #include <vlutil/vlUtil.h>
20 #include <vlutil/vlInetAddr.h>
24 #define DtSMQ_PRIO_0 0x00
25 #define DtSMQ_PRIO_1 0x01
26 #define DtSMQ_PRIO_2 0x02
27 #define DtSMQ_PRIO_3 0x03
28 #define DtSMQ_PRIO_4 0x04
29 #define DtSMQ_PRIO_5 0x05
30 #define DtSMQ_PRIO_6 0x06
31 #define DtSMQ_PRIO_7 0x07
32 #define DtSMQ_PRIO_MASK 0x07
33 #define DtSMQ_MSG_RELIABLE 0x10
34 #define DtSMQ_MSG_BESTEFFORT 0x20
35 //#define DtSMQ_MSG_TBD 0x40
36 #define DtSMQ_SKIP_RELIABLE 0x10
37 #define DtSMQ_SKIP_BESTEFFORT 0x20
52 #define DtRtiSubMsgQ_theMaxSubscribers 16384
115 virtual void unsubscribeSelf();
125 virtual void unsubscribe(MAKRti::DtU32
id);
134 virtual bool queueIsReliable()
const;
137 virtual bool isSubscribed(MAKRti::DtU32 subscriberId)
const;
140 virtual int subscriberCount()
const;
147 virtual bool sendMessage(
const void *message,
unsigned int size);
148 virtual bool sendMessage(
const void *message,
unsigned int size,
unsigned int flags,
149 const MAKRti::DtInetAddr& destAddr = MAKRti::DtInetAddr::inaddrAny());
153 virtual void getRcvMsgTransportInfo(MAKRti::DtInetAddr *destAddr, MAKRti::DtU32 *flags);
157 virtual void getRcvMsgSenderInfo(MAKRti::DtU32* senderId, MAKRti::DtU32* msgSn =
NULL);
161 virtual int waitForMessages(MAKRti::DtTime waitTime = 0);
165 virtual int waitUntilFull(MAKRti::DtTime waitTime = 0);
171 virtual void resetQueue(MAKRti::DtU32 numberOfBuckets, MAKRti::DtU32 payloadSize);
181 virtual void initSyncVars();
185 virtual void clearSyncVars();
188 virtual std::ostream &printDataToStream(std::ostream &str)
const;
196 static unsigned int poolSize(MAKRti::DtU32 numberOfBuckets, MAKRti::DtU32 payloadSize);
205 virtual unsigned int maxSubscribers()
const;
211 MAKRti::DtU32 queueHeaderSize()
const;
214 virtual MAKRti::DtU32 numberOfBuckets()
const;
217 virtual MAKRti::DtU32 payloadSize()
const;
220 virtual volatile MAKRti::DtU32 queueState()
const;
224 virtual volatile MAKRti::DtU32 queueFront(
unsigned int sub)
const;
227 virtual volatile MAKRti::DtU32 queueBack()
const;
231 virtual volatile MAKRti::DtU32 queueBackstop()
const;
235 virtual MAKRti::DtU32 updateFront(
unsigned int sub,
unsigned int val);
239 virtual MAKRti::DtU32 updateBack(
unsigned int val);
243 virtual MAKRti::DtU32 updateBackstop(
unsigned int val);
246 virtual void setSubscribed(
unsigned int sub,
unsigned int val);
249 virtual int setSubscriberCount(MAKRti::DtU32 val);
252 virtual int incrementSubscriberCount();
255 virtual int decrementSubscriberCount();
258 virtual void setState(MAKRti::DtU32 val);
261 virtual void lockForRead();
263 virtual void unlockForRead();
266 virtual bool stalledSubscribers(
unsigned int bucketNum,
unsigned int numStalled);
270 virtual bool handleStall(
unsigned int bucketNum,
unsigned int sub);
273 virtual MAKRti::DtU32 bucketHeaderSize()
const;
276 virtual MAKRti::DtU32 bucketSenderId(
unsigned int bucketNum)
const;
277 virtual MAKRti::DtU32 bucketMsgSn(
unsigned int bucketNum)
const;
278 virtual MAKRti::DtU32 bucketOutstandingReadCount(
unsigned int bucketNum)
const;
279 virtual MAKRti::DtInetAddr bucketDestinationAddr(
unsigned int bucketNum)
const;
280 virtual MAKRti::DtU32 bucketPayloadLen(
unsigned int bucketNum)
const;
281 virtual MAKRti::DtU32 bucketPriority(
unsigned int bucketNum)
const;
282 virtual bool bucketIsReliable(
unsigned int bucketNum)
const;
283 virtual bool bucketIsBestEffort(
unsigned int bucketNum)
const;
284 virtual bool bucketIsFirst(
unsigned int bucketNum)
const;
285 virtual bool bucketIsLast(
unsigned int bucketNum)
const;
288 virtual void setSenderId(
unsigned int bucketNum, MAKRti::DtU32 sid);
289 virtual void setMsgSn(
unsigned int bucketNum, MAKRti::DtU32 seqnum);
290 virtual void setOutstandingReadCount(
unsigned int bucketNum, MAKRti::DtU32 val);
291 virtual void decrementOutstandingReadCount(
unsigned int bucketNum);
292 virtual void setDestinationAddr(
unsigned int bucketNum,
const MAKRti::DtInetAddr& dest);
293 virtual void setPayloadLen(
unsigned int bucketNum, MAKRti::DtU32 len);
294 virtual void setPriority(
unsigned int bucketNum, MAKRti::DtU32 priority);
295 virtual void setReliable(
unsigned int bucketNum, MAKRti::DtU32 val);
296 virtual void setBestEffort(
unsigned int bucketNum, MAKRti::DtU32 val);
297 virtual void setFirst(
unsigned int bucketNum, MAKRti::DtU32 val);
298 virtual void setLast(
unsigned int bucketNum, MAKRti::DtU32 val);
301 virtual pthread_rwlock_t * rdWrLock();
305 virtual pthread_mutex_t * orcMutex();
306 virtual pthread_cond_t * newMessageCond();
307 virtual pthread_mutex_t * newMessageCondMutex();
308 virtual pthread_cond_t * queueFullCond();
309 virtual pthread_mutex_t * queueFullCondMutex();
345 MAKRti::DtU32 sizeInBytes : 24;
346 MAKRti::DtU32 msgPriority : 3;
347 MAKRti::DtU32 reliable : 1;
348 MAKRti::DtU32 besteffort : 1;
349 MAKRti::DtU32 ipv6 : 1;
350 MAKRti::DtU32 last : 1;
351 MAKRti::DtU32 first : 1;
361 MAKRti::DtU32 payloadSize : 24;
362 MAKRti::DtU32 reliable : 1;
363 MAKRti::DtU32 reserved_1 : 7;
371 pthread_rwlock_t rwlock;
388 DtStalledSubscriber(MAKRti::DtU32 bucket, MAKRti::DtU32
id, MAKRti::DtU32 sender, MAKRti::DtU32 seq, MAKRti::DtTime tstamp) :
389 bucketIndex(bucket), stalledId(id), senderId(sender), senderSn(seq), posted(tstamp)
392 MAKRti::DtU32 bucketIndex;
MAKRti::DtU32 senderId
ID of subscriber which hasn't read message.
Definition: smSubscribableMessageQueue.h:394
unsigned int theMaxStallTime
The maximum time to permit a subscriber to be stalled, in seconds.
Definition: smSubscribableMessageQueue.h:332
MAKRti::DtU32 stalledId
Definition: smSubscribableMessageQueue.h:393
#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
#define DtSM_INVALID_SUBSCRIBER_ID
Definition: vlShmSubscribableMessageQueue__rti__.h:37
MAKRti::DtIpcCondVar__rti__ * myNewMessageCond
Definition: smSubscribableMessageQueue.h:316
DtSmqEnqueueMode
Definition: smSubscribableMessageQueue.h:39
This must always be the last item.
Definition: smSubscribableMessageQueue.h:44
DtStalledSubscriber(MAKRti::DtU32 bucket, MAKRti::DtU32 id, MAKRti::DtU32 sender, MAKRti::DtU32 seq, MAKRti::DtTime tstamp)
Definition: smSubscribableMessageQueue.h:388
#define NULL
Definition: baseTypes13.h:8
MAKRti::DtClock myTimestampClock
Definition: smSubscribableMessageQueue.h:315
#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:52
MAKRti::DtU32 senderSn
ID of sender whose message is unread.
Definition: smSubscribableMessageQueue.h:395
~DtStalledSubscriber()
Definition: smSubscribableMessageQueue.h:391
Definition: smSubscribableMessageQueue.h:42
unsigned int theMaxSubscribers
The maximum number of subscribers allowed in this queue.
Definition: smSubscribableMessageQueue.h:328
Definition: smSubscribableMessageQueue.h:41
MAKRti::DtIpcMutex__rti__ * myQueueFullMutex
Definition: smSubscribableMessageQueue.h:319
DtStalledSubscriber keeps track of which queue subscribers have not read a bucket as well as informat...
Definition: smSubscribableMessageQueue.h:385
DtStalledSubscriberList myStalledSubs
!WIN32
Definition: smSubscribableMessageQueue.h:314
MAKRti::DtIpcCondVar__rti__ * myQueueFullCond
Definition: smSubscribableMessageQueue.h:318
Implements a message queue for multiple subscribers.
Definition: smSubscribableMessageQueue.h:78
#define DT_DLL_LRC
Definition: rtiMsConfig.h:155
MAKRti::DtTime posted
Sender-assigned seq. no. of unread message.
Definition: smSubscribableMessageQueue.h:396
MAKRti::DtIpcMutex__rti__ * myNewMessageMutex
Definition: smSubscribableMessageQueue.h:317
std::list< DtStalledSubscriber * > DtStalledSubscriberList
Definition: smSubscribableMessageQueue.h:201