MAK RTIspy API Documentation for HLA 1.3
 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 
23 using namespace MAKRti;
24 
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
40 
42 {
45 
47 };
48 
54 #define DtRtiSubMsgQ_theMaxSubscribers 16384
55 
80 class DT_DLL_LRC DtRtiSubscribableMessageQueue: public MAKRti::DtSmSubscribableMessageQueue__rti__
81 {
82 
83 public:
84 
89  DtRtiSubscribableMessageQueue(const MAKRti::DtString &name);
91 
92 private:
93 
96 
99 
100 public:
101 
107  virtual int subscribe(MAKRti::DtU32 subscriberId = DtSM_INVALID_SUBSCRIBER_ID);
108 
117  virtual void unsubscribeSelf();
118 
127  virtual void unsubscribe(DtU32 id);
128 
133  virtual bool isEmpty(int *nextSize, DtU32 flagsMask = DtSMQ_FETCH_DEFAULT);
134 
136  virtual bool queueIsReliable() const;
137 
139  virtual bool isSubscribed(DtU32 subscriberId) const;
140 
142  virtual int subscriberCount() const;
143 
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());
152 
155  virtual void getRcvMsgTransportInfo(MAKRti::DtInetAddr *destAddr, DtU32 *flags);
156 
159  virtual void getRcvMsgSenderInfo(DtU32* senderId, DtU32* msgSn = NULL);
160 
163  virtual int waitForMessages(MAKRti::DtTime waitTime = 0);
164 
167  virtual int waitUntilFull(MAKRti::DtTime waitTime = 0);
168 
173  virtual void resetQueue(DtU32 numberOfBuckets, DtU32 payloadSize);
174 
178  virtual void setEnqueueMode(DtSmqEnqueueMode mode);
179 
183  virtual void initSyncVars();
184 
187  virtual void clearSyncVars();
188 
190  virtual std::ostream &printDataToStream(std::ostream &str) const;
191 
192 public:
193 
198  static unsigned int poolSize(DtU32 numberOfBuckets, DtU32 payloadSize);
199 
200 protected:
201  struct DtRtiQueueHeader;
202  struct DtRtiBucketHeader;
204  typedef std::list< DtStalledSubscriber * > DtStalledSubscriberList;
205 
207  virtual unsigned int maxSubscribers() const;
208 
210  volatile DtRtiQueueHeader *queueHeader() const;
211 
213  DtU32 queueHeaderSize() const;
214 
216  virtual DtU32 numberOfBuckets() const;
217 
219  virtual DtU32 payloadSize() const;
220 
222  virtual volatile DtU32 queueState() const;
223 
226  virtual volatile DtU32 queueFront(unsigned int sub) const;
227 
229  virtual volatile DtU32 queueBack() const;
230 
233  virtual volatile DtU32 queueBackstop() const;
234 
237  virtual DtU32 updateFront(unsigned int sub, unsigned int val);
238 
241  virtual DtU32 updateBack(unsigned int val);
242 
245  virtual DtU32 updateBackstop(unsigned int val);
246 
248  virtual void setSubscribed(unsigned int sub, unsigned int val);
249 
251  virtual int setSubscriberCount(DtU32 val);
252 
254  virtual int incrementSubscriberCount();
255 
257  virtual int decrementSubscriberCount();
258 
260  virtual void setState(DtU32 val);
261 
263  virtual void lockForRead();
265  virtual void unlockForRead();
266 
268  virtual bool stalledSubscribers(unsigned int bucketNum, unsigned int numStalled);
269 
272  virtual bool handleStall(unsigned int bucketNum, unsigned int sub);
273 
275  virtual DtU32 bucketHeaderSize() const;
276 
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;
288 
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);
301 
302 #ifndef WIN32
303 
304 
305 
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();
312 #endif
313 
314 protected:
315 
317  MAKRti::DtClock myTimestampClock;
322 
330  unsigned int theMaxSubscribers;
331 
334  unsigned int theMaxStallTime;
335 
336 protected:
337 
343  {
344  DtU32 senderId;
345  DtU32 msgSn;
347  DtU32 sizeInBytes : 24;
348  DtU32 msgPriority : 3;
349  DtU32 reliable : 1;
350  DtU32 besteffort : 1;
351  DtU32 ipv6 : 1;
352  DtU32 last : 1;
353  DtU32 first : 1;
354  DtU128 destAddr;
355  };
356 
361  {
362  DtU32 state;
363  DtU32 payloadSize : 24;
364  DtU32 reliable : 1;
365  DtU32 reserved_1 : 7;
370  DtU32 back;
371  DtU32 backstop;
372 #ifndef WIN32
373 
374 
375  pthread_rwlock_t rwlock;
376  pthread_cond_t newMessage;
377  pthread_mutex_t newMessageMutex;
378  pthread_cond_t queueFull;
379  pthread_mutex_t queueFullMutex;
380  pthread_mutex_t ptOrcMutex;
381 #endif
382  };
383 
388  {
389  public:
390  DtStalledSubscriber(DtU32 bucket, DtU32 id, DtU32 sender, DtU32 seq, MAKRti::DtTime tstamp) :
391  bucketIndex(bucket), stalledId(id), senderId(sender), senderSn(seq), posted(tstamp)
392  {};
394  DtU32 bucketIndex;
395  DtU32 stalledId;
396  DtU32 senderId;
397  DtU32 senderSn;
398  MAKRti::DtTime posted;
399  };
400 
401 };
402 
403 #endif
404 
405 

Document ID: Generated on Fri Sep 27 00:57:56 EDT 2019 from SVN revision 201708
Copyright © 2005-2018 VT MÄK Inc. All Rights Reserved (www.mak.com)