MAK RTIspy API Documentation for HLA 1.3
smSubscribableMessageQueue.h
Go to the documentation of this file.
00001 /*********************************************************************
00002  ** Copyright (c) 2005 MaK Technologies, Inc.
00003  ** All rights reserved.
00004  *********************************************************************/
00005 /*********************************************************************
00006  ** $RCSfile: smSubscribableMessageQueue.h,v $ $Revision: 1.14 $ $State: Exp $
00007  *********************************************************************/
00008 #ifndef smSubscribableMessageQueue_H_
00009 #define smSubscribableMessageQueue_H_
00010 
00013 
00014 #include "rtiMsConfig.h"
00015 #include "vlShmSubscribableMessageQueue__rti__.h"
00016 #include "vlIpcCondVar__rti__.h"
00017 #include "vlIpcMutex__rti__.h"
00018 #include <vlutil/vlTime.h>
00019 #include <vlutil/vlUtil.h>
00020 #include <vlutil/vlInetAddr.h>
00021 #include <list>
00022 
00023 using namespace MAKRti;
00024 
00026 #define DtSMQ_PRIO_0          0x00  //! lowest priority messages
00027 #define DtSMQ_PRIO_1          0x01  //! 
00028 #define DtSMQ_PRIO_2          0x02  //! 
00029 #define DtSMQ_PRIO_3          0x03  //! 
00030 #define DtSMQ_PRIO_4          0x04  //! 
00031 #define DtSMQ_PRIO_5          0x05  //! 
00032 #define DtSMQ_PRIO_6          0x06  //! 
00033 #define DtSMQ_PRIO_7          0x07  //! highest priority messages
00034 #define DtSMQ_PRIO_MASK       0x07  //! mask for message priority levels
00035 #define DtSMQ_MSG_RELIABLE    0x10  //! mark message for reliable transport
00036 #define DtSMQ_MSG_BESTEFFORT  0x20  //! mark message for best-effort transport
00037 //#define DtSMQ_MSG_TBD         0x40   //! for future use - leave undefined for now
00038 #define DtSMQ_SKIP_RELIABLE   0x10  //! fetch all messages except those marked for reliable transport
00039 #define DtSMQ_SKIP_BESTEFFORT 0x20  //! fetch all messages except those marked for best-effort transport
00040 
00041 enum DtSmqEnqueueMode
00042 {
00043    DtSmqEnqueueBestEffortMode = 0,
00044    DtSmqEnqueueReliableMode   = 1,
00045 
00046    DtSmqEnqueueInvalidMode        
00047 };
00048 
00054 #define DtRtiSubMsgQ_theMaxSubscribers 16384
00055 
00080 class DT_DLL_LRC DtRtiSubscribableMessageQueue: public MAKRti::DtSmSubscribableMessageQueue__rti__
00081 {
00082 
00083 public:
00084 
00089    DtRtiSubscribableMessageQueue(const MAKRti::DtString &name);
00090    virtual ~DtRtiSubscribableMessageQueue();
00091 
00092 private:
00093 
00095    DtRtiSubscribableMessageQueue(const DtRtiSubscribableMessageQueue &other);
00096 
00098    DtRtiSubscribableMessageQueue& operator=(const DtRtiSubscribableMessageQueue &other);
00099 
00100 public:
00101 
00107    virtual int subscribe(MAKRti::DtU32 subscriberId = DtSM_INVALID_SUBSCRIBER_ID);
00108 
00117    virtual void unsubscribeSelf();
00118 
00127    virtual void unsubscribe(DtU32 id);
00128 
00133    virtual bool isEmpty(int *nextSize, DtU32 flagsMask = DtSMQ_FETCH_DEFAULT);
00134 
00136    virtual bool queueIsReliable() const;
00137 
00139    virtual bool isSubscribed(DtU32 subscriberId) const;
00140 
00142    virtual int subscriberCount() const;
00143 
00149    virtual bool sendMessage(const void *message, unsigned int size);
00150    virtual bool sendMessage(const void *message, unsigned int size, unsigned int flags,
00151       const MAKRti::DtInetAddr& destAddr = MAKRti::DtInetAddr::inaddrAny());
00152 
00155    virtual void getRcvMsgTransportInfo(MAKRti::DtInetAddr *destAddr, DtU32 *flags);
00156 
00159    virtual void getRcvMsgSenderInfo(DtU32* senderId, DtU32* msgSn = NULL);
00160 
00163    virtual int waitForMessages(MAKRti::DtTime waitTime = 0);
00164 
00167    virtual int waitUntilFull(MAKRti::DtTime waitTime = 0);
00168 
00173    virtual void resetQueue(DtU32 numberOfBuckets, DtU32 payloadSize);
00174 
00178    virtual void setEnqueueMode(DtSmqEnqueueMode mode);
00179 
00183    virtual void initSyncVars();
00184 
00187    virtual void clearSyncVars();
00188 
00190    virtual std::ostream &printDataToStream(std::ostream &str) const;
00191 
00192 public:
00193 
00198    static unsigned int poolSize(DtU32 numberOfBuckets, DtU32 payloadSize);
00199 
00200 protected:
00201    struct DtRtiQueueHeader;
00202    struct DtRtiBucketHeader;
00203    class DtStalledSubscriber;
00204    typedef std::list< DtStalledSubscriber * > DtStalledSubscriberList;
00205 
00207    virtual unsigned int maxSubscribers() const;
00208 
00210    volatile DtRtiQueueHeader *queueHeader() const;
00211 
00213    DtU32 queueHeaderSize() const;
00214 
00216    virtual DtU32 numberOfBuckets() const;
00217 
00219    virtual DtU32 payloadSize() const;
00220 
00222    virtual volatile DtU32 queueState() const;
00223 
00226    virtual volatile DtU32 queueFront(unsigned int sub) const;
00227 
00229    virtual volatile DtU32 queueBack() const;
00230 
00233    virtual volatile DtU32 queueBackstop() const;
00234 
00237    virtual DtU32 updateFront(unsigned int sub, unsigned int val);
00238 
00241    virtual DtU32 updateBack(unsigned int val);
00242 
00245    virtual DtU32 updateBackstop(unsigned int val);
00246 
00248    virtual void setSubscribed(unsigned int sub, unsigned int val);
00249 
00251    virtual int setSubscriberCount(DtU32 val);
00252 
00254    virtual int incrementSubscriberCount();
00255 
00257    virtual int decrementSubscriberCount();
00258 
00260    virtual void setState(DtU32 val);
00261 
00263    virtual void lockForRead();
00265    virtual void unlockForRead();
00266 
00268    virtual bool stalledSubscribers(unsigned int bucketNum, unsigned int numStalled);
00269 
00272    virtual bool handleStall(unsigned int bucketNum, unsigned int sub);
00273 
00275    virtual DtU32 bucketHeaderSize() const;
00276 
00278    virtual DtU32 bucketSenderId(unsigned int bucketNum) const;
00279    virtual DtU32 bucketMsgSn(unsigned int bucketNum) const;
00280    virtual DtU32 bucketOutstandingReadCount(unsigned int bucketNum) const;
00281    virtual DtInetAddr bucketDestinationAddr(unsigned int bucketNum) const;
00282    virtual DtU32 bucketPayloadLen(unsigned int bucketNum) const;
00283    virtual DtU32 bucketPriority(unsigned int bucketNum) const;
00284    virtual bool  bucketIsReliable(unsigned int bucketNum) const;
00285    virtual bool  bucketIsBestEffort(unsigned int bucketNum) const;
00286    virtual bool  bucketIsFirst(unsigned int bucketNum) const;
00287    virtual bool  bucketIsLast(unsigned int bucketNum) const;
00288 
00290    virtual void  setSenderId(unsigned int bucketNum, DtU32 sid);
00291    virtual void  setMsgSn(unsigned int bucketNum, DtU32 seqnum);
00292    virtual void  setOutstandingReadCount(unsigned int bucketNum, DtU32 val);
00293    virtual void  decrementOutstandingReadCount(unsigned int bucketNum);
00294    virtual void  setDestinationAddr(unsigned int bucketNum, const DtInetAddr& dest);
00295    virtual void  setPayloadLen(unsigned int bucketNum, DtU32 len);
00296    virtual void  setPriority(unsigned int bucketNum, DtU32 priority);
00297    virtual void  setReliable(unsigned int bucketNum, DtU32 val);
00298    virtual void  setBestEffort(unsigned int bucketNum, DtU32 val);
00299    virtual void  setFirst(unsigned int bucketNum, DtU32 val);
00300    virtual void  setLast(unsigned int bucketNum, DtU32 val);
00301 
00302 #ifndef WIN32
00303 
00304 
00305 
00306    virtual pthread_rwlock_t * rdWrLock();
00307    virtual pthread_mutex_t  * orcMutex();
00308    virtual pthread_cond_t   * newMessageCond();
00309    virtual pthread_mutex_t  * newMessageCondMutex();
00310    virtual pthread_cond_t   * queueFullCond();
00311    virtual pthread_mutex_t  * queueFullCondMutex();
00312 #endif //! !WIN32
00313 
00314 protected:
00315 
00316    DtStalledSubscriberList   myStalledSubs;
00317    MAKRti::DtClock                   myTimestampClock;
00318    DtIpcCondVar__rti__           *  myNewMessageCond;
00319    DtIpcMutex__rti__             *  myNewMessageMutex;
00320    DtIpcCondVar__rti__           *  myQueueFullCond;
00321    DtIpcMutex__rti__             *  myQueueFullMutex;
00322 
00330    unsigned int theMaxSubscribers;
00331 
00334    unsigned int theMaxStallTime;
00335 
00336 protected:
00337 
00342    struct DtRtiBucketHeader
00343    {
00344       DtU32 senderId;         
00345       DtU32 msgSn;            
00346       DtU32 outstandingReadCount;  
00347       DtU32 sizeInBytes : 24; 
00348       DtU32 msgPriority : 3;  
00349       DtU32 reliable    : 1;  
00350       DtU32 besteffort  : 1;  
00351       DtU32 ipv6        : 1;  
00352       DtU32 last        : 1;  
00353       DtU32 first       : 1;  
00354       DtU128 destAddr;        
00355    };
00356 
00360    struct DtRtiQueueHeader
00361    {
00362       DtU32 state;             
00363       DtU32 payloadSize : 24;  
00364       DtU32 reliable    :  1;
00365       DtU32 reserved_1  :  7;
00366       DtU32 numberOfBuckets;
00367       DtU32 numberOfSubscribers;
00368       DtU32 front[DtRtiSubMsgQ_theMaxSubscribers];
00369       DtU32 subscribed[DtRtiSubMsgQ_theMaxSubscribers];
00370       DtU32 back;
00371       DtU32 backstop;
00372 #ifndef WIN32
00373 
00374 
00375       pthread_rwlock_t  rwlock;
00376       pthread_cond_t    newMessage;   
00377       pthread_mutex_t   newMessageMutex;
00378       pthread_cond_t    queueFull;
00379       pthread_mutex_t   queueFullMutex;
00380       pthread_mutex_t   ptOrcMutex;   
00381 #endif //! !WIN32
00382    };
00383 
00387    class DtStalledSubscriber
00388    {
00389    public:
00390       DtStalledSubscriber(DtU32 bucket, DtU32 id, DtU32 sender, DtU32 seq, MAKRti::DtTime tstamp) :
00391          bucketIndex(bucket), stalledId(id), senderId(sender), senderSn(seq), posted(tstamp)
00392          {};
00393       ~DtStalledSubscriber(){};
00394       DtU32    bucketIndex; 
00395       DtU32    stalledId;  
00396       DtU32    senderId;   
00397       DtU32    senderSn;   
00398       MAKRti::DtTime   posted;     
00399    };
00400 
00401 };
00402 
00403 #endif
00404 
00405 

Document ID: Generated on Thu Jun 14 14:15:04 EDT 2012 from SVN revision 116116
Copyright © 2005-2012 VT MÄK Inc. All Rights Reserved (www.mak.com)