MAK RTIspy API Documentation for HLA Evolved
 All Classes Namespaces Files Functions Variables Typedefs Enumerations Enumerator Properties Friends Macros Groups Pages
smSubscribableMessageQueue.h
Go to the documentation of this file.
1 /*********************************************************************
2  ** Copyright (c) 2005 MaK Technologies, Inc.
3  ** All rights reserved.
4  *********************************************************************/
5 /*********************************************************************
6  ** $RCSfile: smSubscribableMessageQueue.h,v $ $Revision: 1.14 $ $State: Exp $
7  *********************************************************************/
8 #ifndef smSubscribableMessageQueue_H_
9 #define smSubscribableMessageQueue_H_
10 
13 
14 #include "rtiMsConfig.h"
16 #include "vlIpcCondVar__rti__.h"
17 #include "vlIpcMutex__rti__.h"
18 #include <vlutil/vlTime.h>
19 #include <vlutil/vlUtil.h>
20 #include <vlutil/vlInetAddr.h>
21 #include <list>
22 
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
38 
40 {
43 
45 };
46 
52 #define DtRtiSubMsgQ_theMaxSubscribers 16384
53 
78 class DT_DLL_LRC DtRtiSubscribableMessageQueue: public MAKRti::DtSmSubscribableMessageQueue__rti__
79 {
80 
81 public:
82 
87  DtRtiSubscribableMessageQueue(const MAKRti::DtString &name);
89 
90 private:
91 
94 
97 
98 public:
99 
105  virtual int subscribe(MAKRti::DtU32 subscriberId = DtSM_INVALID_SUBSCRIBER_ID);
106 
115  virtual void unsubscribeSelf();
116 
125  virtual void unsubscribe(MAKRti::DtU32 id);
126 
131  virtual bool isEmpty(int *nextSize, MAKRti::DtU32 flagsMask = DtSMQ_FETCH_DEFAULT);
132 
134  virtual bool queueIsReliable() const;
135 
137  virtual bool isSubscribed(MAKRti::DtU32 subscriberId) const;
138 
140  virtual int subscriberCount() const;
141 
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());
150 
153  virtual void getRcvMsgTransportInfo(MAKRti::DtInetAddr *destAddr, MAKRti::DtU32 *flags);
154 
157  virtual void getRcvMsgSenderInfo(MAKRti::DtU32* senderId, MAKRti::DtU32* msgSn = NULL);
158 
161  virtual int waitForMessages(MAKRti::DtTime waitTime = 0);
162 
165  virtual int waitUntilFull(MAKRti::DtTime waitTime = 0);
166 
171  virtual void resetQueue(MAKRti::DtU32 numberOfBuckets, MAKRti::DtU32 payloadSize);
172 
176  virtual void setEnqueueMode(DtSmqEnqueueMode mode);
177 
181  virtual void initSyncVars();
182 
185  virtual void clearSyncVars();
186 
188  virtual std::ostream &printDataToStream(std::ostream &str) const;
189 
190 public:
191 
196  static unsigned int poolSize(MAKRti::DtU32 numberOfBuckets, MAKRti::DtU32 payloadSize);
197 
198 protected:
199  struct DtRtiQueueHeader;
200  struct DtRtiBucketHeader;
202  typedef std::list< DtStalledSubscriber * > DtStalledSubscriberList;
203 
205  virtual unsigned int maxSubscribers() const;
206 
208  volatile DtRtiQueueHeader *queueHeader() const;
209 
211  MAKRti::DtU32 queueHeaderSize() const;
212 
214  virtual MAKRti::DtU32 numberOfBuckets() const;
215 
217  virtual MAKRti::DtU32 payloadSize() const;
218 
220  virtual volatile MAKRti::DtU32 queueState() const;
221 
224  virtual volatile MAKRti::DtU32 queueFront(unsigned int sub) const;
225 
227  virtual volatile MAKRti::DtU32 queueBack() const;
228 
231  virtual volatile MAKRti::DtU32 queueBackstop() const;
232 
235  virtual MAKRti::DtU32 updateFront(unsigned int sub, unsigned int val);
236 
239  virtual MAKRti::DtU32 updateBack(unsigned int val);
240 
243  virtual MAKRti::DtU32 updateBackstop(unsigned int val);
244 
246  virtual void setSubscribed(unsigned int sub, unsigned int val);
247 
249  virtual int setSubscriberCount(MAKRti::DtU32 val);
250 
252  virtual int incrementSubscriberCount();
253 
255  virtual int decrementSubscriberCount();
256 
258  virtual void setState(MAKRti::DtU32 val);
259 
261  virtual void lockForRead();
263  virtual void unlockForRead();
264 
266  virtual bool stalledSubscribers(unsigned int bucketNum, unsigned int numStalled);
267 
270  virtual bool handleStall(unsigned int bucketNum, unsigned int sub);
271 
273  virtual MAKRti::DtU32 bucketHeaderSize() const;
274 
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;
286 
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);
299 
300 #ifndef WIN32
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();
310 #endif
311 
312 protected:
313 
315  MAKRti::DtClock myTimestampClock;
316  MAKRti::DtIpcCondVar__rti__ * myNewMessageCond;
317  MAKRti::DtIpcMutex__rti__ * myNewMessageMutex;
318  MAKRti::DtIpcCondVar__rti__ * myQueueFullCond;
319  MAKRti::DtIpcMutex__rti__ * myQueueFullMutex;
320 
328  unsigned int theMaxSubscribers;
329 
332  unsigned int theMaxStallTime;
333 
334 protected:
335 
341  {
342  MAKRti::DtU32 senderId;
343  MAKRti::DtU32 msgSn;
344  MAKRti::DtU32 outstandingReadCount;
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;
352  MAKRti::DtU128 destAddr;
353  };
354 
359  {
360  MAKRti::DtU32 state;
361  MAKRti::DtU32 payloadSize : 24;
362  MAKRti::DtU32 reliable : 1;
363  MAKRti::DtU32 reserved_1 : 7;
364  MAKRti::DtU32 numberOfBuckets;
365  MAKRti::DtU32 numberOfSubscribers;
366  MAKRti::DtU32 front[DtRtiSubMsgQ_theMaxSubscribers];
367  MAKRti::DtU32 subscribed[DtRtiSubMsgQ_theMaxSubscribers];
368  MAKRti::DtU32 back;
369  MAKRti::DtU32 backstop;
370 #ifndef WIN32
371  pthread_rwlock_t rwlock;
374  pthread_cond_t newMessage;
375  pthread_mutex_t newMessageMutex;
376  pthread_cond_t queueFull;
377  pthread_mutex_t queueFullMutex;
378  pthread_mutex_t ptOrcMutex;
379 #endif
380  };
381 
386  {
387  public:
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)
390  {};
392  MAKRti::DtU32 bucketIndex;
393  MAKRti::DtU32 stalledId;
394  MAKRti::DtU32 senderId;
395  MAKRti::DtU32 senderSn;
396  MAKRti::DtTime posted;
397  };
398 
399 };
400 
401 #endif
402 
403 
MAKRti::DtU32 senderId
ID of subscriber which hasn&#39;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
pthread_mutex_t newMessageMutex
signal on arrival of new messages
Definition: smSubscribableMessageQueue.h:375
MAKRti::DtU128 destAddr
if set start of message in bucket
Definition: smSubscribableMessageQueue.h:352
#define DtSM_INVALID_SUBSCRIBER_ID
Definition: vlShmSubscribableMessageQueue__rti__.h:37
DtQueueHeader exists at the front of the message queue and contains useful information about the stat...
Definition: smSubscribableMessageQueue.h:358
MAKRti::DtIpcCondVar__rti__ * myNewMessageCond
Definition: smSubscribableMessageQueue.h:316
MAKRti::DtU32 numberOfSubscribers
Definition: smSubscribableMessageQueue.h:365
DtSmqEnqueueMode
Definition: smSubscribableMessageQueue.h:39
MAKRti::DtU32 backstop
Definition: smSubscribableMessageQueue.h:369
MAKRti::DtU32 senderId
Definition: smSubscribableMessageQueue.h:342
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
MAKRti::DtU32 msgSn
the ID of the subscriber
Definition: smSubscribableMessageQueue.h:343
MAKRti::DtClock myTimestampClock
Definition: smSubscribableMessageQueue.h:315
MAKRti::DtU32 state
Definition: smSubscribableMessageQueue.h:360
#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
pthread_cond_t newMessage
Definition: smSubscribableMessageQueue.h:374
MAKRti::DtU32 back
Definition: smSubscribableMessageQueue.h:368
~DtStalledSubscriber()
Definition: smSubscribableMessageQueue.h:391
MAKRti::DtU32 numberOfBuckets
Definition: smSubscribableMessageQueue.h:364
Definition: smSubscribableMessageQueue.h:42
unsigned int theMaxSubscribers
The maximum number of subscribers allowed in this queue.
Definition: smSubscribableMessageQueue.h:328
pthread_cond_t queueFull
Definition: smSubscribableMessageQueue.h:376
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
pthread_mutex_t ptOrcMutex
Definition: smSubscribableMessageQueue.h:378
Implements a message queue for multiple subscribers.
Definition: smSubscribableMessageQueue.h:78
MAKRti::DtU32 outstandingReadCount
sequence number (assigned by sender)
Definition: smSubscribableMessageQueue.h:344
#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
pthread_mutex_t queueFullMutex
Definition: smSubscribableMessageQueue.h:377
DtBucketHeader is a structure found at the front of each message to keep track of the number of bytes...
Definition: smSubscribableMessageQueue.h:340

Document ID: Generated on Sun Jul 20 16:15:30 EDT 2025 from SVN revision 277985
Copyright © 2005-2025 MAK Technologies Inc. All Rights Reserved (www.mak.com)