/****************************************************************************/
/* Summary  :                                                               */
/* Filename : MessageQueue.cc                                               */
/* Author   :                                                               */
/* Project  :                                                               */
/* Version  : $Revisions:$                                            */
/* Created  : 04/05/2000                                                    */
/* Modified :                                                               */
/* Archived :                                                               */
/****************************************************************************/
/* Modification History:                                                    */
/****************************************************************************/

#include <errno.h>
#include <malloc.h>
#include <string.h>

#include "Syslog.h"
#include "MessageQueue.h"

/*
CLASS 
MessageQueue

DESCRIPTION
MessageQueue is a base class which implements an OO abstraction of
a message queue (analagous to SharedData for shared memory)

AUTHOR
Marsh
*/

//===============================================================
MessageQueue::MessageQueue( const char *name,
			    int msgSize,
			    int maxQueueSize,
			    Access access,
			    Boolean blocking )
{

     debug = False;

     _queue = -1;

     _name = NULL;
     _name = strdup( name );

     openQueue( name, msgSize, maxQueueSize, access, blocking );
}

//===============================================================
MessageQueue::~MessageQueue()
{
//     unlink();
     closeQueue();

     if( _name ) free( (void *)_name );
}

//===============================================================
int MessageQueue::write( void *data )
{

     if( _queue < 0 ) {
	  dprintf("MessageQueue::write -- Write called when queue not open!" );
	  return -1;
     }

     if( mq_send( _queue,
		  data,
		  _msgSize,
		  0 ) < 0 ) {
	  dprintf("MessageQueue::write -- Write failed (%d): %s",
		  errno, strerror( errno ) );
	  return -1;
     }

     return 0;
}

//===============================================================
int MessageQueue::read( void *data )
{
     unsigned int priority;
     int charsRead;

     if( _queue < 0 ) {
	  dprintf("MessageQueue::read -- Read called when queue not open!");
	  return -1;
     }

     if( (charsRead = mq_receive( _queue,
		      data,
		      _msgSize,
		      &priority )) < 0 ) {
	  dprintf("MessageQueue::read -- Read error (%d):%s",
		  errno, strerror( errno ) );
	  return -1;
     }

     return charsRead;
}

//===============================================================
int MessageQueue::msgsPending( void )
{
     if( _queue < 0 ) return -1;

     mq_attr attr;

     if( mq_getattr( _queue, &attr ) ) {
	  // Something bad happened
	  dprintf("MessageQueue::msgsPending -- Error %s", errno );
	  return -1;
     }

     return attr.mq_curmsgs;
}

//===============================================================
int MessageQueue::setAccess( Access access )
{
     return openQueue( _name,
		_msgSize,
		_maxQueueSize,
		access );
}

//===============================================================
int MessageQueue::flush( void )
{
     int msgsFlushed = 0;
     unsigned int priority;

     if( _queue < 0 ) return -1;

     char *blankData = new char[ _msgSize + 1 ];

     while( msgsPending() )
     {
	  mq_receive( _queue,
		      blankData,
		      _msgSize,
		      &priority );
	  msgsFlushed++;

     }

     delete blankData;

     return msgsFlushed;
}

//===============================================================
int MessageQueue::openQueue( const char *name,
			     int msgSize,
			     int maxQueueSize,
			     Access access,
			     Boolean blocking )
{
     int oflag;
     mode_t mode;
     mq_attr attr;

     if( _queue > 0 ) {
	  dprintf("MessageQueue::openQueue -- openQueue called with queue open.  Closing queue" );
	  if( closeQueue() ) 
	  {
	       // Something bad happened
	       return -1;
	  }
     }

     // Set up the attributes

     attr.mq_maxmsg = maxQueueSize;
     attr.mq_msgsize = msgSize;
     attr.mq_flags = MQ_MULT_NOTIFY | MQ_NOTIFY_ALWAYS;
     if( !blocking )
	  attr.mq_flags |= MQ_NONBLOCK;

//     attr.mq_flags |= MQ_READBUF_DYNAMIC;

     // Set up the output flags
     switch( access ) {
     case ReadWrite:
	  oflag = O_RDWR;
	  break;
     case Read:
	  oflag = O_RDONLY;
	  break;
     case Write:
	  oflag= O_WRONLY;
	  break;
     };

     oflag |= O_CREAT;
     if( !blocking )
	  oflag |= O_NONBLOCK;

     // Set output mode
     mode = 0664;

     if( (_queue = mq_open( name, oflag, mode, &attr )) < 0  ) 
     {
	  // Something bad happend
	  switch( errno ) {
	  case EACCES:
	       Syslog::write("MessageQueue::openQueue -- "
			     "mq_open access error.  Is Mqueue running?");
	       break;
	  }
	  dprintf("MessageQueue::openQueue -- mq_open failed on (%d): %s",
		  errno, strerror( errno ) );
	  _queue = -1;
     }

     _msgSize = msgSize;
     _maxQueueSize = maxQueueSize;
     _access = access;
     _blocking = blocking;

     return 0;
}

//===============================================================
int MessageQueue::closeQueue()
{
     if( _queue < 0 )
	  return 0;

     if( mq_close( _queue ) ) 
     {
	  // Something bad happened
	  dprintf("MessageQueue::closeQueue -- Error on mq_close");
	  _queue = -1;
	  return 0;
     }

     _queue = -1;
     return 0;
}


//===============================================================
int MessageQueue::unlink()
{
     if( _queue < 0 )
	  return 0;

     if( mq_unlink( _name ) ) {
	  dprintf("MessageQueue::unlink -- Error on mq_unlink");
	  return -1;
     }

     return 0;
}

//===============================================================


