(file) Return to Monitor.cpp CVS log (file) (dir) Up to [Pegasus] / pegasus / src / Pegasus / Common

Diff for /pegasus/src/Pegasus/Common/Monitor.cpp between version 1.103.10.22 and 1.124

version 1.103.10.22, 2006/07/21 18:17:51 version 1.124, 2007/10/29 08:14:34
Line 29 
Line 29 
 // //
 //============================================================================== //==============================================================================
 // //
 // Author: Mike Brasher (mbrasher@bmc.com)  
 //  
 // Modified By: Mike Day (monitor_2) mdday@us.ibm.com  
 //              Amit K Arora (Bug#1153) amita@in.ibm.com  
 //              Alagaraja Ramasubramanian (alags_raj@in.ibm.com) for Bug#1090  
 //              Sushma Fernandes (sushma@hp.com) for Bug#2057  
 //              Josephine Eskaline Joyce (jojustin@in.ibm.com) for PEP#101  
 //              Roger Kumpf, Hewlett-Packard Company (roger_kumpf@hp.com)  
 //  
 //%///////////////////////////////////////////////////////////////////////////// //%/////////////////////////////////////////////////////////////////////////////
  
   #include "Network.h"
 #include <Pegasus/Common/Config.h> #include <Pegasus/Common/Config.h>
   
 #include <cstring> #include <cstring>
 #include "Monitor.h" #include "Monitor.h"
 #include "MessageQueue.h" #include "MessageQueue.h"
Line 51 
Line 42 
 #include <Pegasus/Common/MessageQueueService.h> #include <Pegasus/Common/MessageQueueService.h>
 #include <Pegasus/Common/Exception.h> #include <Pegasus/Common/Exception.h>
 #include "ArrayIterator.h" #include "ArrayIterator.h"
   #include "HostAddress.h"
   #include <errno.h>
   
 //const static DWORD MAX_BUFFER_SIZE = 4096;  // 4 kilobytes  
   
 #ifdef PEGASUS_OS_TYPE_WINDOWS  
 # if defined(FD_SETSIZE) && FD_SETSIZE != 1024  
 #  error "FD_SETSIZE was not set to 1024 prior to the last inclusion \  
 of <winsock.h>. It may have been indirectly included (e.g., by including \  
 <windows.h>). Find inclusion of that header which is visible to this \  
 compilation unit and #define FD_SETZIE to 1024 prior to that inclusion; \  
 otherwise, less than 64 clients (the default) will be able to connect to the \  
 CIMOM. PLEASE DO NOT SUPPRESS THIS WARNING; PLEASE FIX THE PROBLEM."  
   
 # endif  
 # define FD_SETSIZE 1024  
 # include <windows.h>  
 #else  
 # include <sys/types.h>  
 # include <sys/socket.h>  
 # include <sys/time.h>  
 # include <netinet/in.h>  
 # include <netdb.h>  
 # include <arpa/inet.h>  
 #endif  
  
 PEGASUS_USING_STD; PEGASUS_USING_STD;
  
 PEGASUS_NAMESPACE_BEGIN PEGASUS_NAMESPACE_BEGIN
  
 static AtomicInt _connections(0);  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
 Mutex Monitor::_cout_mut;  
 #endif  
   
 #ifdef PEGASUS_OS_TYPE_WINDOWS  
  #define PIPE_INCREMENT 1  
 #endif  
   
 //////////////////////////////////////////////////////////////////////////////// ////////////////////////////////////////////////////////////////////////////////
 // //
 // Monitor  // Tickler
 // //
 //////////////////////////////////////////////////////////////////////////////// ////////////////////////////////////////////////////////////////////////////////
  
 #define MAX_NUMBER_OF_MONITOR_ENTRIES  32  Tickler::Tickler()
 Monitor::Monitor()      : _listenSocket(PEGASUS_INVALID_SOCKET),
    : _stopConnections(0),        _clientSocket(PEGASUS_INVALID_SOCKET),
      _stopConnectionsSem(0),        _serverSocket(PEGASUS_INVALID_SOCKET)
      _solicitSocketCount(0),  
      _tickle_client_socket(-1),  
      _tickle_server_socket(-1),  
      _tickle_peer_socket(-1)  
 { {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG      try
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::Monitor(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
     int numberOfMonitorEntriesToAllocate = MAX_NUMBER_OF_MONITOR_ENTRIES;  
     Socket::initializeInterface();  
     _entries.reserveCapacity(numberOfMonitorEntriesToAllocate);  
   
     // setup the tickler  
     initializeTickler();  
   
     // Start the count at 1 because initilizeTickler()  
     // has added an entry in the first position of the  
     // _entries array  
     for( int i = 1; i < numberOfMonitorEntriesToAllocate; i++ )  
     {     {
        _MonitorEntry entry(0, 0, 0);          _initialize();
        _entries.append(entry);  
     }     }
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG      catch (...)
     {     {
         AutoMutex automut(Monitor::_cout_mut);          _uninitialize();
         PEGASUS_STD(cout) << "Exiting:  Monitor::Monitor(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);          throw;
     }     }
 #endif  
 } }
  
 Monitor::~Monitor()  Tickler::~Tickler()
 { {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG      _uninitialize();
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::~Monitor(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }     }
 #endif  
     Tracer::trace(TRC_HTTP, Tracer::LEVEL4, "uninitializing interface");  
  
     try{  #if defined(PEGASUS_OS_TYPE_UNIX) || defined(PEGASUS_OS_VMS)
         if(_tickle_peer_socket >= 0)  
         {  // Use an anonymous pipe for the tickle connection.
             Socket::close(_tickle_peer_socket);  
         }  void Tickler::_initialize()
         if(_tickle_client_socket >= 0)  
         {         {
             Socket::close(_tickle_client_socket);      int fds[2];
         }  
         if(_tickle_server_socket >= 0)      if (pipe(fds) == -1)
         {         {
             Socket::close(_tickle_server_socket);          MessageLoaderParms parms(
         }              "Common.Monitor.TICKLE_CREATE",
               "Received error number $0 while creating the internal socket.",
               getSocketError());
           throw Exception(parms);
     }     }
     catch(...)  
     {      _serverSocket = fds[0];
         Tracer::trace(TRC_HTTP, Tracer::LEVEL4,      _clientSocket = fds[1];
                   "Failed to close tickle sockets");  
     }     }
  
     Socket::uninitializeInterface();  #else
     Tracer::trace(TRC_HTTP, Tracer::LEVEL4,  
                   "returning from monitor destructor");  // Use an external loopback socket connection to allow the tickle socket to
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  // be included in the select() array on Windows.
   
   void Tickler::_initialize()
     {     {
         AutoMutex automut(Monitor::_cout_mut);      //
         PEGASUS_STD(cout) << "Exiting:  Monitor::~Monitor(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);      // Set up the addresses for the listen, client, and server sockets
       // based on whether IPv6 is enabled.
       //
   
       Socket::initializeInterface();
   
   # ifdef PEGASUS_ENABLE_IPV6
       struct sockaddr_storage listenAddress;
       struct sockaddr_storage clientAddress;
       struct sockaddr_storage serverAddress;
   # else
       struct sockaddr_in listenAddress;
       struct sockaddr_in clientAddress;
       struct sockaddr_in serverAddress;
   # endif
   
       int addressFamily;
       SocketLength addressLength;
   
   # ifdef PEGASUS_ENABLE_IPV6
       if (System::isIPv6StackActive())
       {
           // Use the IPv6 loopback address for the listen sockets
           HostAddress::convertTextToBinary(
               HostAddress::AT_IPV6,
               "::1",
               &reinterpret_cast<struct sockaddr_in6*>(&listenAddress)->sin6_addr);
           listenAddress.ss_family = AF_INET6;
           reinterpret_cast<struct sockaddr_in6*>(&listenAddress)->sin6_port = 0;
   
           addressFamily = AF_INET6;
           addressLength = sizeof(struct sockaddr_in6);
     }     }
       else
 #endif #endif
 }  
   
 void Monitor::initializeTickler(){  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {     {
         AutoMutex automut(Monitor::_cout_mut);          // Use the IPv4 loopback address for the listen sockets
         PEGASUS_STD(cout) << "Entering: Monitor::initializeTickler(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);          HostAddress::convertTextToBinary(
               HostAddress::AT_IPV4,
               "127.0.0.1",
               &reinterpret_cast<struct sockaddr_in*>(
                   &listenAddress)->sin_addr.s_addr);
           reinterpret_cast<struct sockaddr_in*>(&listenAddress)->sin_family =
               AF_INET;
           reinterpret_cast<struct sockaddr_in*>(&listenAddress)->sin_port = 0;
   
           addressFamily = AF_INET;
           addressLength = sizeof(struct sockaddr_in);
     }     }
 #endif  
     /*  
        NOTE: On any errors trying to  
              setup out tickle connection,  
              throw an exception/end the server  
     */  
  
     /* setup the tickle server/listener */      // Use the same address for the client socket as the listen socket
       clientAddress = listenAddress;
  
     // get a socket for the server side      //
     if((_tickle_server_socket = ::socket(PF_INET, SOCK_STREAM, 0)) == PEGASUS_INVALID_SOCKET){      // Set up a listen socket to allow the tickle client and server to connect
         //handle error      //
         MessageLoaderParms parms("Common.Monitor.TICKLE_CREATE",  
       // Create the listen socket
       if ((_listenSocket = Socket::createSocket(addressFamily, SOCK_STREAM, 0)) ==
                PEGASUS_INVALID_SOCKET)
       {
           MessageLoaderParms parms(
               "Common.Monitor.TICKLE_CREATE",
                                  "Received error number $0 while creating the internal socket.",                                  "Received error number $0 while creating the internal socket.",
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)              getSocketError());
                                  errno);  
 #else  
                                  WSAGetLastError());  
 #endif  
         throw Exception(parms);         throw Exception(parms);
     }     }
  
     // initialize the address      // Bind the listen socket to the loopback address
     memset(&_tickle_server_addr, 0, sizeof(_tickle_server_addr));      if (::bind(
 #ifdef PEGASUS_PLATFORM_OS400_ISERIES_IBM              _listenSocket,
 #pragma convert(37)              reinterpret_cast<struct sockaddr*>(&listenAddress),
 #endif              addressLength) < 0)
     _tickle_server_addr.sin_addr.s_addr = inet_addr("127.0.0.1");      {
 #ifdef PEGASUS_PLATFORM_OS400_ISERIES_IBM          MessageLoaderParms parms(
 #pragma convert(0)              "Common.Monitor.TICKLE_BIND",
 #endif  
     _tickle_server_addr.sin_family = PF_INET;  
     _tickle_server_addr.sin_port = 0;  
   
     PEGASUS_SOCKLEN_T _addr_size = sizeof(_tickle_server_addr);  
   
     // bind server side to socket  
     if((::bind(_tickle_server_socket,  
                reinterpret_cast<struct sockaddr*>(&_tickle_server_addr),  
                sizeof(_tickle_server_addr))) < 0){  
         // handle error  
 #ifdef PEGASUS_OS_ZOS  
     MessageLoaderParms parms("Common.Monitor.TICKLE_BIND_LONG",  
                                  "Received error:$0 while binding the internal socket.",strerror(errno));  
 #else  
         MessageLoaderParms parms("Common.Monitor.TICKLE_BIND",  
                                  "Received error number $0 while binding the internal socket.",                                  "Received error number $0 while binding the internal socket.",
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)              getSocketError());
                                  errno);  
 #else  
                                  WSAGetLastError());  
 #endif  
 #endif  
         throw Exception(parms);         throw Exception(parms);
     }     }
  
     // tell the kernel we are a server      // Listen for a connection from the tickle client
     if((::listen(_tickle_server_socket,3)) < 0){      if ((::listen(_listenSocket, 3)) < 0)
         // handle error      {
         MessageLoaderParms parms("Common.Monitor.TICKLE_LISTEN",          MessageLoaderParms parms(
               "Common.Monitor.TICKLE_LISTEN",
                          "Received error number $0 while listening to the internal socket.",                          "Received error number $0 while listening to the internal socket.",
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)              getSocketError());
                                  errno);  
 #else  
                                  WSAGetLastError());  
 #endif  
         throw Exception(parms);         throw Exception(parms);
     }     }
  
     // make sure we have the correct socket for our server      // Verify we have the correct listen socket
     int sock = ::getsockname(_tickle_server_socket,      SocketLength tmpAddressLength = addressLength;
                    reinterpret_cast<struct sockaddr*>(&_tickle_server_addr),      int sock = ::getsockname(
                    &_addr_size);          _listenSocket,
     if(sock < 0){          reinterpret_cast<struct sockaddr*>(&listenAddress),
         // handle error          &tmpAddressLength);
         MessageLoaderParms parms("Common.Monitor.TICKLE_SOCKNAME",      if (sock < 0)
       {
           MessageLoaderParms parms(
               "Common.Monitor.TICKLE_SOCKNAME",
                          "Received error number $0 while getting the internal socket name.",                          "Received error number $0 while getting the internal socket name.",
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)              getSocketError());
                                  errno);  
 #else  
                                  WSAGetLastError());  
 #endif  
         throw Exception(parms);         throw Exception(parms);
     }     }
  
     /* set up the tickle client/connector */      //
       // Set up the client side of the tickle connection.
       //
  
     // get a socket for our tickle client      // Create the client socket
     if((_tickle_client_socket = ::socket(PF_INET, SOCK_STREAM, 0)) == PEGASUS_INVALID_SOCKET){      if ((_clientSocket = Socket::createSocket(addressFamily, SOCK_STREAM, 0)) ==
         // handle error               PEGASUS_INVALID_SOCKET)
         MessageLoaderParms parms("Common.Monitor.TICKLE_CLIENT_CREATE",      {
                          "Received error number $0 while creating the internal client socket.",          MessageLoaderParms parms(
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)              "Common.Monitor.TICKLE_CLIENT_CREATE",
                                  errno);              "Received error number $0 while creating the internal client "
 #else                  "socket.",
                                  WSAGetLastError());              getSocketError());
 #endif  
         throw Exception(parms);         throw Exception(parms);
     }     }
  
     // setup the address of the client      // Bind the client socket to the loopback address
     memset(&_tickle_client_addr, 0, sizeof(_tickle_client_addr));      if (::bind(
 #ifdef PEGASUS_PLATFORM_OS400_ISERIES_IBM              _clientSocket,
 #pragma convert(37)              reinterpret_cast<struct sockaddr*>(&clientAddress),
 #endif              addressLength) < 0)
     _tickle_client_addr.sin_addr.s_addr = inet_addr("127.0.0.1");      {
 #ifdef PEGASUS_PLATFORM_OS400_ISERIES_IBM          MessageLoaderParms parms(
 #pragma convert(0)              "Common.Monitor.TICKLE_CLIENT_BIND",
 #endif              "Received error number $0 while binding the internal client "
     _tickle_client_addr.sin_family = PF_INET;                  "socket.",
     _tickle_client_addr.sin_port = 0;              getSocketError());
           throw Exception(parms);
       }
  
     // bind socket to client side      // Connect the client socket to the listen socket address
     if((::bind(_tickle_client_socket,      if (::connect(
                reinterpret_cast<struct sockaddr*>(&_tickle_client_addr),              _clientSocket,
                sizeof(_tickle_client_addr))) < 0){              reinterpret_cast<struct sockaddr*>(&listenAddress),
         // handle error              addressLength) < 0)
         MessageLoaderParms parms("Common.Monitor.TICKLE_CLIENT_BIND",      {
                          "Received error number $0 while binding the internal client socket.",          MessageLoaderParms parms(
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)              "Common.Monitor.TICKLE_CLIENT_CONNECT",
                                  errno);              "Received error number $0 while connecting the internal client "
 #else                  "socket.",
                                  WSAGetLastError());              getSocketError());
 #endif  
         throw Exception(parms);         throw Exception(parms);
     }     }
  
     // connect to server side      //
     if((::connect(_tickle_client_socket,      // Set up the server side of the tickle connection.
                   reinterpret_cast<struct sockaddr*>(&_tickle_server_addr),      //
                   sizeof(_tickle_server_addr))) < 0){  
         // handle error      tmpAddressLength = addressLength;
         MessageLoaderParms parms("Common.Monitor.TICKLE_CLIENT_CONNECT",  
                          "Received error number $0 while connecting the internal client socket.",      // Accept the client socket connection.
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)      _serverSocket = ::accept(
                                  errno);          _listenSocket,
 #else          reinterpret_cast<struct sockaddr*>(&serverAddress),
                                  WSAGetLastError());          &tmpAddressLength);
 #endif  
       if (_serverSocket == PEGASUS_SOCKET_ERROR)
       {
           MessageLoaderParms parms(
               "Common.Monitor.TICKLE_ACCEPT",
               "Received error number $0 while accepting the internal socket "
                   "connection.",
               getSocketError());
         throw Exception(parms);         throw Exception(parms);
     }     }
  
     /* set up the slave connection */      //
     memset(&_tickle_peer_addr, 0, sizeof(_tickle_peer_addr));      // Close the listen socket and make the other sockets non-blocking
     PEGASUS_SOCKLEN_T peer_size = sizeof(_tickle_peer_addr);      //
     pegasus_sleep(1);  
       Socket::close(_listenSocket);
     // this call may fail, we will try a max of 20 times to establish this peer connection      _listenSocket = PEGASUS_INVALID_SOCKET;
     if((_tickle_peer_socket = ::accept(_tickle_server_socket,  
             reinterpret_cast<struct sockaddr*>(&_tickle_peer_addr),      Socket::disableBlocking(_serverSocket);
             &peer_size)) < 0){      Socket::disableBlocking(_clientSocket);
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)  
         // Only retry on non-windows platforms.  
         if(_tickle_peer_socket == -1 && errno == EAGAIN)  
         {  
           int retries = 0;  
           do  
           {  
             pegasus_sleep(1);  
             _tickle_peer_socket = ::accept(_tickle_server_socket,  
                 reinterpret_cast<struct sockaddr*>(&_tickle_peer_addr),  
                 &peer_size);  
             retries++;  
           } while(_tickle_peer_socket == -1 && errno == EAGAIN && retries < 20);  
         }         }
   
 #endif #endif
   
   void Tickler::_uninitialize()
   {
       PEG_TRACE_CSTRING(TRC_HTTP, Tracer::LEVEL4, "uninitializing interface");
   
       try
       {
           if (_serverSocket != PEGASUS_INVALID_SOCKET)
           {
               Socket::close(_serverSocket);
               _serverSocket = PEGASUS_INVALID_SOCKET;
     }     }
     if(_tickle_peer_socket == -1){          if (_clientSocket != PEGASUS_INVALID_SOCKET)
         // handle error          {
         MessageLoaderParms parms("Common.Monitor.TICKLE_ACCEPT",              Socket::close(_clientSocket);
                          "Received error number $0 while accepting the internal socket connection.",              _clientSocket = PEGASUS_INVALID_SOCKET;
 #if !defined(PEGASUS_OS_TYPE_WINDOWS)  
                                  errno);  
 #else  
                                  WSAGetLastError());  
 #endif  
         throw Exception(parms);  
     }     }
     // add the tickler to the list of entries to be monitored and set to IDLE because Monitor only          if (_listenSocket != PEGASUS_INVALID_SOCKET)
     // checks entries with IDLE state for events  
     _MonitorEntry entry(_tickle_peer_socket, 1, INTERNAL);  
     entry._status = _MonitorEntry::IDLE;  
     _entries.append(entry);  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {     {
         AutoMutex automut(Monitor::_cout_mut);              Socket::close(_listenSocket);
         PEGASUS_STD(cout) << "Exiting:  Monitor::initializeTickler(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);              _listenSocket = PEGASUS_INVALID_SOCKET;
     }     }
 #endif  
 } }
       catch (...)
 void Monitor::tickle(void)  
 {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {     {
         AutoMutex automut(Monitor::_cout_mut);          PEG_TRACE_CSTRING(TRC_HTTP, Tracer::LEVEL4,
         PEGASUS_STD(cout) << "Entering: Monitor::tickle(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);              "Failed to close tickle sockets");
     }     }
 #endif      Socket::uninitializeInterface();
     static char _buffer[] =  }
   
   
   ////////////////////////////////////////////////////////////////////////////////
   //
   // Monitor
   //
   ////////////////////////////////////////////////////////////////////////////////
   
   #define MAX_NUMBER_OF_MONITOR_ENTRIES  32
   Monitor::Monitor()
      : _stopConnections(0),
        _stopConnectionsSem(0),
        _solicitSocketCount(0)
     {     {
       '0','0'      int numberOfMonitorEntriesToAllocate = MAX_NUMBER_OF_MONITOR_ENTRIES;
     };      _entries.reserveCapacity(numberOfMonitorEntriesToAllocate);
  
     AutoMutex autoMutex(_tickle_mutex);      // Create a MonitorEntry for the Tickler and set its state to IDLE so the
     Socket::disableBlocking(_tickle_client_socket);      // Monitor will watch for its events.
     Socket::write(_tickle_client_socket,&_buffer, 2);      _MonitorEntry entry(_tickler.getServerSocket(), 1, INTERNAL);
     Socket::enableBlocking(_tickle_client_socket);      entry._status = _MonitorEntry::IDLE;
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG      _entries.append(entry);
   
       // Start the count at 1 because _entries[0] is the Tickler
       for (int i = 1; i < numberOfMonitorEntriesToAllocate; i++)
     {     {
         AutoMutex automut(Monitor::_cout_mut);         _MonitorEntry entry(0, 0, 0);
         PEGASUS_STD(cout) << "Exiting:  Monitor::tickle(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);         _entries.append(entry);
     }     }
 #endif  
 } }
  
 void Monitor::setState( Uint32 index, _MonitorEntry::entry_status status )  Monitor::~Monitor()
 { {
     // Set the state to requested state      PEG_TRACE_CSTRING(TRC_HTTP, Tracer::LEVEL4,
     _entries[index]._status = status;                    "returning from monitor destructor");
 } }
  
 Boolean Monitor::run(Uint32 milliseconds)  void Monitor::tickle()
 { {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG      AutoMutex autoMutex(_tickleMutex);
     {      Socket::write(_tickler.getClientSocket(), "\0\0", 2);
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::run(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }     }
 #endif  
  
     Boolean handled_events = false;  void Monitor::setState(
     int i = 0;      Uint32 index,
       _MonitorEntry::entry_status status)
   {
       AutoMutex autoEntryMutex(_entry_mut);
       // Set the state to requested state
       _entries[index]._status = status;
   }
  
   void Monitor::run(Uint32 milliseconds)
   {
     struct timeval tv = {milliseconds/1000, milliseconds%1000*1000};     struct timeval tv = {milliseconds/1000, milliseconds%1000*1000};
  
     fd_set fdread;     fd_set fdread;
Line 424 
Line 382 
  
     ArrayIterator<_MonitorEntry> entries(_entries);     ArrayIterator<_MonitorEntry> entries(_entries);
  
     // Check the stopConnections flag.  If set, clear the Acceptor monitor entries      // Check the stopConnections flag.  If set, clear the Acceptor monitor
       // entries
     if (_stopConnections.get() == 1)     if (_stopConnections.get() == 1)
     {     {
         for ( int indx = 0; indx < (int)entries.size(); indx++)         for ( int indx = 0; indx < (int)entries.size(); indx++)
Line 472 
Line 431 
  
                                         if (h._responsePending == true)                                         if (h._responsePending == true)
                                         {                                         {
                         if (!entry.namedPipeConnection)                  PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
                         {                      "Monitor::run - Ignoring connection delete request "
                             Tracer::trace(TRC_HTTP, Tracer::LEVEL4, "Monitor::run - "                          "because responses are still pending. "
                                                                                                         "Ignoring connection delete request because "  
                                                                                                         "responses are still pending. "  
                                                                                                         "connection=0x%p, socket=%d\n",                                                                                                         "connection=0x%p, socket=%d\n",
                                                                                                         (void *)&h, h.getSocket());                      (void *)&h, h.getSocket()));
                         }  
                         else  
                         {  
                             Tracer::trace(TRC_HTTP, Tracer::LEVEL4, "Monitor::run - "  
                                                                                                         "Ignoring connection delete request because "  
                                                                                                         "responses are still pending. "  
                                                                                                         "connection=0x%p, NamedPipe=%d\n",  
                                                                                                         (void *)&h, h.getNamedPipe().getPipe());  
                         }  
                                                 continue;                                                 continue;
                                         }                                         }
                                         h._connectionClosePending = false;                                         h._connectionClosePending = false;
           MessageQueue &o = h.get_owner();           MessageQueue &o = h.get_owner();
           Message* message;              Message* message= new CloseConnectionMessage(entry.socket);
           if (!entry.namedPipeConnection)  
           {  
               message= new CloseConnectionMessage(entry.socket);  
           }  
           else  
           {  
               message= new CloseConnectionMessage(entry.namedPipe);  
   
           }  
           message->dest = o.getQueueId();           message->dest = o.getQueueId();
  
           // HTTPAcceptor is responsible for closing the connection.           // HTTPAcceptor is responsible for closing the connection.
Line 516 
Line 455 
           // unlocked will not result in an ArrayIndexOutOfBounds           // unlocked will not result in an ArrayIndexOutOfBounds
           // exception.           // exception.
  
           autoEntryMutex.unlock();              _entry_mut.unlock();
           o.enqueue(message);           o.enqueue(message);
           autoEntryMutex.lock();              _entry_mut.lock();
           // After enqueue a message and the autoEntryMutex has been released and locked again,  
           // the array of _entries can be changed. The ArrayIterator has be reset with the original _entries.              // After enqueue a message and the autoEntryMutex has been
               // released and locked again, the array of _entries can be
               // changed. The ArrayIterator has be reset with the original
               // _entries.
           entries.reset(_entries);           entries.reset(_entries);
        }        }
     }     }
Line 533 
Line 475 
         place to calculate the max file descriptor (maximum socket number)         place to calculate the max file descriptor (maximum socket number)
         because we have to traverse the entire array.         because we have to traverse the entire array.
     */     */
     //Array<HANDLE> pipeEventArray;      SocketHandle maxSocketCurrentPass = 0;
         PEGASUS_SOCKET maxSocketCurrentPass = 0;      for (int indx = 0; indx < (int)entries.size(); indx++)
     int indx;  
   
   
 #ifdef PEGASUS_OS_TYPE_WINDOWS  
   
     //This array associates named pipe connections to their place in [indx]  
     //in the entries array. The value in poition zero of the array is the  
     //index of the fist named pipe connection in the entries array  
     Array <Uint32> indexPipeCountAssociator;  
     int pipeEntryCount=0;  
     int MaxPipes = PIPE_INCREMENT;  
     HANDLE* hEvents = new HANDLE[PIPE_INCREMENT];  
   
 #endif  
   
     for( indx = 0; indx < (int)entries.size(); indx++)  
     {  
   
   
 #ifdef PEGASUS_OS_TYPE_WINDOWS  
        if(entries[indx].isNamedPipeConnection())  
        {  
   
            //entering this clause mean that a Named Pipe connection is at entries[indx]  
            //cout << "In Monitor::run in clause to to create array of for WaitformultipuleObjects" << endl;  
   
            //cout << "In Monitor::run - pipe being added to array is " << entries[indx].namedPipe.getName() << endl;  
   
             entries[indx].pipeSet = false;  
   
            // We can Keep a counter in the Monitor class for the number of named pipes ...  
            //  Which can be used here to create the array size for hEvents..( obviously before this for loop.:-) )  
             if (pipeEntryCount >= MaxPipes)  
             {  
                // cout << "Monitor::run 'if (pipeEntryCount >= MaxPipes)' begining - pipeEntryCount=" <<  
                    // pipeEntryCount << " MaxPipes=" << MaxPipes << endl;  
                  MaxPipes += PIPE_INCREMENT;  
                  HANDLE* temp_hEvents = new HANDLE[MaxPipes];  
   
                  for (Uint32 i =0;i<pipeEntryCount;i++)  
                  {  
                      temp_hEvents[i] = hEvents[i];  
                  }  
   
                  delete [] hEvents;  
   
                  hEvents = temp_hEvents;  
                 // cout << "Monitor::run 'if (pipeEntryCount >= MaxPipes)' ending"<< endl;  
   
             }  
   
            //pipeEventArray.append((entries[indx].namedPipe.getOverlap()).hEvent);  
            hEvents[pipeEntryCount] = entries[indx].namedPipe.getOverlap()->hEvent;  
   
            indexPipeCountAssociator.append(indx);  
   
        pipeEntryCount++;  
   
   
   
        }  
        else  
   
 #endif  
        {        {
   
            if(maxSocketCurrentPass < entries[indx].socket)            if(maxSocketCurrentPass < entries[indx].socket)
             maxSocketCurrentPass = entries[indx].socket;             maxSocketCurrentPass = entries[indx].socket;
  
Line 609 
Line 486 
                _idleEntries++;                _idleEntries++;
                FD_SET(entries[indx].socket, &fdread);                FD_SET(entries[indx].socket, &fdread);
            }            }
   
        }  
   }   }
  
     /*     /*
Line 619 
Line 494 
     */     */
     maxSocketCurrentPass++;     maxSocketCurrentPass++;
  
     autoEntryMutex.unlock();      _entry_mut.unlock();
  
     //     //
     // The first argument to select() is ignored on Windows and it is not     // The first argument to select() is ignored on Windows and it is not
     // a socket value.  The original code assumed that the number of sockets     // a socket value.  The original code assumed that the number of sockets
     // and a socket value have the same type.  On Windows they do not.     // and a socket value have the same type.  On Windows they do not.
     //     //
   
     int events;  
     int pEvents;  
   
 #ifdef PEGASUS_OS_TYPE_WINDOWS #ifdef PEGASUS_OS_TYPE_WINDOWS
       int events = select(0, &fdread, NULL, NULL, &tv);
    // events = select(0, &fdread, NULL, NULL, &tv);  
   
     //if (events == NULL)  
     //{  // This connection uses namedPipes  
   
         events = 0;  
         DWORD dwWait=NULL;  
         pEvents = 0;  
   
   
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
        {  
         AutoMutex automut(Monitor::_cout_mut);  
         cout << "Monitor::run - Calling WaitForMultipleObjects\n";  
         }  
 #endif  
    // }  
         //this should be in a try block  
   
     dwWait = WaitForMultipleObjects(  
                  MaxPipes,  
                  hEvents,               //ABB:- array of event objects  
                  FALSE,                 // ABB:-does not wait for all  
                  milliseconds);        //ABB:- timeout value   //WW this may need be shorter  
   
     if(dwWait == WAIT_TIMEOUT)  
         {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
         {  
             AutoMutex automut(Monitor::_cout_mut);  
         cout << "Wait WAIT_TIMEOUT\n";  
         cout << "Monitor::run before the select in TIMEOUT clause events = " << events << endl;  
         }  
 #endif  
                 events = select(0, &fdread, NULL, NULL, &tv);  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
             AutoMutex automut(Monitor::_cout_mut);  
            cout << "Monitor::run after the select in TIMEOUT clause events = " << events << endl;  
 #endif  
   
   
                    // Sleep(2000);  
             //continue;  
   
              //return false;  // I think we do nothing.... Mybe there is a socket connection... so  
              // cant return.  
         }  
         else if (dwWait == WAIT_FAILED)  
         {  
             if (GetLastError() == 6) //WW this may be too specific  
             {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                 AutoMutex automut(Monitor::_cout_mut);  
                 cout << "Monitor::run about to call 'select since waitForMultipleObjects failed\n";  
 #endif  
                 /********* NOTE  
                 this time (tv) combined with the waitForMulitpleObjects timeout is  
                 too long it will cause the client side to time out  
                 ******************/  
                 events = select(0, &fdread, NULL, NULL, &tv);  
   
             }  
             else  
             {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                 AutoMutex automut(Monitor::_cout_mut);  
                 cout << "Wait Failed returned\n";  
                 cout << "failed with " << GetLastError() << "." << endl;  
 #endif  
                 pEvents = -1;  
                 return false;  
             }  
         }  
         else  
         {  
             int pCount = dwWait - WAIT_OBJECT_0;  // determines which pipe  
             {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                  {  
                      AutoMutex automut(Monitor::_cout_mut);  
                      // cout << endl << "****************************" <<  
                      //  "Monitor::run WaitForMultiPleObject returned activity on server pipe: "<<  
                      //  pCount<< endl <<  endl;  
                      cout << "Monitor::run WaitForMultiPleObject returned activity pipeEntrycount is " <<  
                      pipeEntryCount <<  
                      " this is the type " << entries[indexPipeCountAssociator[pCount]]._type << " this is index " << indexPipeCountAssociator[pCount] << endl;  
                  }  
 #endif  
   
                /* There is a timeing problem here sometimes the wite in HTTPConnection i s  
              not all the way done (has not _monitor->setState (_entry_index, _MonitorEntry::IDLE) )  
              there for that should be done here if it is not done alread*/  
   
                if (entries[indexPipeCountAssociator[pCount]]._status.get() != _MonitorEntry::IDLE)  
                {  
                    this->setState(indexPipeCountAssociator[pCount], _MonitorEntry::IDLE);  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
             AutoMutex automut(Monitor::_cout_mut);  
   
                    cout << "setting state of index " << indexPipeCountAssociator[pCount]  << " to IDLE" << endl;  
 #endif  
                }  
   
   
             }  
   
             pEvents = 1;  
   
             //this statment gets the pipe entry that was trigered  
             entries[indexPipeCountAssociator[pCount]].pipeSet = true;  
   
         }  
 #else #else
     events = select(maxSocketCurrentPass, &fdread, NULL, NULL, &tv);      int events = select(maxSocketCurrentPass, &fdread, NULL, NULL, &tv);
 #endif #endif
     autoEntryMutex.lock();      _entry_mut.lock();
     // After enqueue a message and the autoEntryMutex has been released and locked again,  
     // the array of _entries can be changed. The ArrayIterator has be reset with the original _entries  
     entries.reset(_entries);  
  
 #ifdef PEGASUS_OS_TYPE_WINDOWS      // After enqueue a message and the autoEntryMutex has been released and
     if(pEvents == -1)      // locked again, the array of _entries can be changed. The ArrayIterator
     {      // has be reset with the original _entries
         Tracer::trace(TRC_HTTP, Tracer::LEVEL4,      entries.reset(_entries);
           "Monitor::run - errorno = %d has occurred on select.",GetLastError() );  
        // The EBADF error indicates that one or more or the file  
        // descriptions was not valid. This could indicate that  
        // the entries structure has been corrupted or that  
        // we have a synchronization error.  
   
         // We need to generate an assert  here...  
        PEGASUS_ASSERT(GetLastError()!= EBADF);  
   
   
     }  
  
     if(events == SOCKET_ERROR)      if (events == PEGASUS_SOCKET_ERROR)
 #else  
     if(events == -1)  
 #endif  
     {     {
           PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
         Tracer::trace(TRC_HTTP, Tracer::LEVEL4,              "Monitor::run - errorno = %d has occurred on select.", errno));
           "Monitor::run - errorno = %d has occurred on select.", errno);  
        // The EBADF error indicates that one or more or the file        // The EBADF error indicates that one or more or the file
        // descriptions was not valid. This could indicate that        // descriptions was not valid. This could indicate that
        // the entries structure has been corrupted or that        // the entries structure has been corrupted or that
Line 783 
Line 524 
  
        PEGASUS_ASSERT(errno != EBADF);        PEGASUS_ASSERT(errno != EBADF);
     }     }
     else if ((events)||(pEvents))      else if (events)
     {  
   
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
         {         {
                  AutoMutex automut(Monitor::_cout_mut);          PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
      cout << "IN Monior::run events= " << events << " pEvents= " << pEvents<< endl;              "Monitor::run select event received events = %d, monitoring %d "
         }                  "idle entries",
 #endif              events, _idleEntries));
   
      Tracer::trace(TRC_HTTP, Tracer::LEVEL4,  
           "Monitor::run select event received events = %d, monitoring %d idle entries",  
            events, _idleEntries);  
        for( int indx = 0; indx < (int)entries.size(); indx++)        for( int indx = 0; indx < (int)entries.size(); indx++)
        {        {
            //cout << "Monitor::run at start of 'for( int indx = 0; indx ' - index = " << indx << endl;              // The Monitor should only look at entries in the table that are
           // The Monitor should only look at entries in the table that are IDLE (i.e.,              // IDLE (i.e., owned by the Monitor).
           // owned by the Monitor).  
         // cout << endl << " status of entry " << indx << " is " << entries[indx]._status.get() << endl;  
           if((entries[indx]._status.get() == _MonitorEntry::IDLE) &&           if((entries[indx]._status.get() == _MonitorEntry::IDLE) &&
              ((FD_ISSET(entries[indx].socket, &fdread)&& (events)) ||                  (FD_ISSET(entries[indx].socket, &fdread)))
              (entries[indx].isNamedPipeConnection() && entries[indx].pipeSet && (pEvents))))  
           {           {
                   MessageQueue *q = MessageQueue::lookup(entries[indx].queueId);
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG                  PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
               {  
                  AutoMutex automut(Monitor::_cout_mut);  
                  cout <<"Monitor::run - index  " << indx << " just got into 'if' statement" << endl;  
               }  
 #endif  
               MessageQueue *q;  
            try{  
   
                  q = MessageQueue::lookup(entries[indx].queueId);  
               }  
              catch (Exception e)  
              {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                  AutoMutex automut(Monitor::_cout_mut);  
                  cout << " this is what lookup gives - " << e.getMessage() << endl;  
 #endif  
                  exit(1);  
              }  
              catch(...)  
              {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                  AutoMutex automut(Monitor::_cout_mut);  
                  cout << "MessageQueue::lookup gives strange exception " << endl;  
 #endif  
                  exit(1);  
              }  
   
   
   
   
               Tracer::trace(TRC_HTTP, Tracer::LEVEL4,  
                   "Monitor::run indx = %d, queueId =  %d, q = %p",                   "Monitor::run indx = %d, queueId =  %d, q = %p",
                   indx, entries[indx].queueId, q);                      indx, entries[indx].queueId, q));
            //  printf("Monitor::run indx = %d, queueId =  %d, q = %p",  
              //     indx, entries[indx].queueId, q);  
              //cout << "Monitor::run before PEGASUS_ASSerT(q !=0) " << endl;  
              PEGASUS_ASSERT(q !=0);              PEGASUS_ASSERT(q !=0);
  
   
              try              try
              {              {
                 /* {  
                  AutoMutex automut(Monitor::_cout_mut);  
                   cout <<" this is the type " << entries[indx]._type <<  
                       " for index " << indx << endl;  
                cout << "IN Monior::run right before entries[indx]._type == Monitor::CONNECTION" << endl;  
                  }*/  
                if(entries[indx]._type == Monitor::CONNECTION)                if(entries[indx]._type == Monitor::CONNECTION)
                 {                 {
                           PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG                              "entries[indx].type for indx = %d is "
                     {                                  "Monitor::CONNECTION",
                     cout << "In Monitor::run Monitor::CONNECTION clause" << endl;                              indx));
                     AutoMutex automut(Monitor::_cout_mut);  
                     }  
 #endif  
   
                                       Tracer::trace(TRC_HTTP, Tracer::LEVEL4,  
                      "entries[indx].type for indx = %d is Monitor::CONNECTION", indx);  
                    static_cast<HTTPConnection *>(q)->_entry_index = indx;                    static_cast<HTTPConnection *>(q)->_entry_index = indx;
  
                    // Do not update the entry just yet. The entry gets updated once                          // Do not update the entry just yet. The entry gets
                    // the request has been read.                          // updated once the request has been read.
                    //entries[indx]._status = _MonitorEntry::BUSY;                    //entries[indx]._status = _MonitorEntry::BUSY;
  
                    // If allocate_and_awaken failure, retry on next iteration                          // If allocate_and_awaken failure, retry on next
                           // iteration
 /* Removed for PEP 183. /* Removed for PEP 183.
                    if (!MessageQueueService::get_thread_pool()->allocate_and_awaken(                          if (!MessageQueueService::get_thread_pool()->
                            (void *)q, _dispatch))                                  allocate_and_awaken((void *)q, _dispatch))
                    {                    {
                       Tracer::trace(TRC_DISCARDED_DATA, Tracer::LEVEL2,                              PEG_TRACE_CSTRING(TRC_DISCARDED_DATA,
                           "Monitor::run: Insufficient resources to process request.");                                  Tracer::LEVEL2,
                                   "Monitor::run: Insufficient resources to "
                                       "process request.");
                       entries[indx]._status = _MonitorEntry::IDLE;                       entries[indx]._status = _MonitorEntry::IDLE;
                       return true;                       return true;
                    }                    }
 */ */
 // Added for PEP 183 // Added for PEP 183
                    HTTPConnection *dst = reinterpret_cast<HTTPConnection *>(q);                          HTTPConnection *dst =
                    Tracer::trace(TRC_HTTP, Tracer::LEVEL4,                              reinterpret_cast<HTTPConnection *>(q);
                          "Monitor::_dispatch: entering run() for indx  = %d, queueId = %d, q = %p",                          PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
                    dst->_entry_index, dst->_monitor->_entries[dst->_entry_index].queueId, dst);                              "Monitor::_dispatch: entering run() for "
                                   "indx = %d, queueId = %d, q = %p",
                    /*In the case of named Pipes, the request has already been read from the pipe                              dst->_entry_index,
                    therefor this section passed the request data to the HTTPConnection                              dst->_monitor->_entries[dst->_entry_index].queueId,
                    NOTE: not sure if this would be better suited in a sparate private method                              dst));
                    */  
  
                    dst->setNamedPipe(entries[indx].namedPipe); //this step shouldn't be needd  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                    {  
                        AutoMutex automut(Monitor::_cout_mut);  
                    cout << "In Monitor::run after dst->setNamedPipe string read is " <<  entries[indx].namedPipe.raw << endl;  
                    }  
 #endif  
                    try                    try
                    {                    {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                        {  
                        AutoMutex automut(Monitor::_cout_mut);  
                        cout << "In Monitor::run about to call 'dst->run(1)' "  << endl;  
                        }  
 #endif  
                        dst->run(1);                        dst->run(1);
                    }                    }
                    catch (...)                    catch (...)
                    {                    {
                               PEG_TRACE_CSTRING(TRC_HTTP, Tracer::LEVEL4,
                        Tracer::trace(TRC_HTTP, Tracer::LEVEL4,  
                        "Monitor::_dispatch: exception received");                        "Monitor::_dispatch: exception received");
                    }                    }
                    Tracer::trace(TRC_HTTP, Tracer::LEVEL4,                          PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
                    "Monitor::_dispatch: exited run() for index %d", dst->_entry_index);                              "Monitor::_dispatch: exited run() for index %d",
                               dst->_entry_index));
                    if (entries[indx].isNamedPipeConnection())  
                    {                          // It is possible the entry status may not be set to
                        entries[indx]._type = Monitor::ACCEPTOR;                          // busy.  The following will fail in that case.
                    }                          // PEGASUS_ASSERT(dst->_monitor->_entries[
                           //     dst->_entry_index]._status.get() ==
                    // It is possible the entry status may not be set to busy.                          //    _MonitorEntry::BUSY);
                    // The following will fail in that case.                          // Once the HTTPConnection thread has set the status
                    // PEGASUS_ASSERT(dst->_monitor->_entries[dst->_entry_index]._status.get() == _MonitorEntry::BUSY);                          // value to either Monitor::DYING or Monitor::IDLE,
                    // Once the HTTPConnection thread has set the status value to either                          // it has returned control of the connection to the
                    // Monitor::DYING or Monitor::IDLE, it has returned control of the connection                          // Monitor.  It is no longer permissible to access
                    // to the Monitor.  It is no longer permissible to access the connection                          // the connection or the entry in the _entries table.
                    // or the entry in the _entries table.  
                           // The following is not relevant as the worker thread
                    // The following is not relevant as the worker thread or the                          // or the reader thread will update the status of the
                    // reader thread will update the status of the entry.                          // entry.
                    //if (dst->_connectionClosePending)                    //if (dst->_connectionClosePending)
                    //{                    //{
                    //  dst->_monitor->_entries[dst->_entry_index]._status = _MonitorEntry::DYING;                          //  dst->_monitor->_entries[dst->_entry_index]._status =
                           //    _MonitorEntry::DYING;
                    //}                    //}
                    //else                    //else
                    //{                    //{
                    //  dst->_monitor->_entries[dst->_entry_index]._status = _MonitorEntry::IDLE;                          //  dst->_monitor->_entries[dst->_entry_index]._status =
                           //    _MonitorEntry::IDLE;
                    //}                    //}
 // end Added for PEP 183 // end Added for PEP 183
                 }                 }
                 else if( entries[indx]._type == Monitor::INTERNAL){                      else if (entries[indx]._type == Monitor::INTERNAL)
                       {
                         // set ourself to BUSY,                         // set ourself to BUSY,
                         // read the data                         // read the data
                         // and set ourself back to IDLE                         // and set ourself back to IDLE
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
             AutoMutex automut(Monitor::_cout_mut);  
  
             cout << endl << " in - entries[indx]._type == Monitor::INTERNAL- " << endl << endl;  
 #endif  
             if (!entries[indx].isNamedPipeConnection())  
             {  
                             entries[indx]._status = _MonitorEntry::BUSY;                             entries[indx]._status = _MonitorEntry::BUSY;
                             static char buffer[2];                             static char buffer[2];
                         Socket::disableBlocking(entries[indx].socket);                          Sint32 amt =
                         Sint32 amt = Socket::read(entries[indx].socket,&buffer, 2);                              Socket::read(entries[indx].socket,&buffer, 2);
                         Socket::enableBlocking(entries[indx].socket);  
                             entries[indx]._status = _MonitorEntry::IDLE;                             entries[indx]._status = _MonitorEntry::IDLE;
             }             }
                 }  
                 else                 else
                 {                 {
                           PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG                              "Non-connection entry, indx = %d, has been "
       {                                  "received.",
             AutoMutex automut(Monitor::_cout_mut);                              indx));
             cout << "In Monitor::run else clause of CONNECTION if statments" << endl;  
       }  
 #endif  
                                Tracer::trace(TRC_HTTP, Tracer::LEVEL4,  
                      "Non-connection entry, indx = %d, has been received.", indx);  
                    int events = 0;                    int events = 0;
            Message *msg;  
   
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
           {  
            AutoMutex automut(Monitor::_cout_mut);  
            cout << " In Monitor::run Just before checking if NamedPipeConnection" << "for Index "<<indx<< endl;  
            }  
 #endif  
            if (entries[indx].isNamedPipeConnection())  
            {  
                if(!entries[indx].namedPipe.isConnectionPipe)  
                { /*if we enter this clasue it means that the named pipe that we are  
                    looking at has recived a connection but is not the pipe we get connection requests over.  
                    therefore we need to change the _type to CONNECTION and wait for a CIM Operations request*/  
                    entries[indx]._type = Monitor::CONNECTION;  
   
   
      /* This is a test  - this shows that the read file needs to be done  
      before we call wiatForMultipleObjects*/  
     /******************************************************  
     ********************************************************/  
   
   
   
         memset(entries[indx].namedPipe.raw,'\0',NAMEDPIPE_MAX_BUFFER_SIZE);  
         BOOL rc = ::ReadFile(  
                 entries[indx].namedPipe.getPipe(),  
                 &entries[indx].namedPipe.raw,  
                 NAMEDPIPE_MAX_BUFFER_SIZE,  
                 &entries[indx].namedPipe.bytesRead,  
                 entries[indx].namedPipe.getOverlap());  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
         {  
          AutoMutex automut(Monitor::_cout_mut);  
          cout << "Monitor::run just called read on index " << indx << endl;  
         }  
 #endif  
   
          //&entries[indx].namedPipe.bytesRead = &size;  
         if(!rc)  
         {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
            AutoMutex automut(Monitor::_cout_mut);  
            cout << "ReadFile failed for : "  << GetLastError() << "."<< endl;  
 #endif  
         }  
   
   
   
     /******************************************************  
     ********************************************************/  
   
   
   
   
                  continue;  
   
   
                }  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                {  
                    AutoMutex automut(Monitor::_cout_mut);  
                     cout << " In Monitor::run about to create a Pipe message" << endl;  
   
                }  
 #endif  
                events |= NamedPipeMessage::READ;  
                msg = new NamedPipeMessage(entries[indx].namedPipe, events);  
            }  
            else  
            {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                {  
                AutoMutex automut(Monitor::_cout_mut);  
                cout << " In Monitor::run ..its a socket message" << endl;  
                }  
 #endif  
                events |= SocketMessage::READ;                events |= SocketMessage::READ;
                        msg = new SocketMessage(entries[indx].socket, events);                          Message* msg = new SocketMessage(
            }                              entries[indx].socket, events);
   
                    entries[indx]._status = _MonitorEntry::BUSY;                    entries[indx]._status = _MonitorEntry::BUSY;
                    autoEntryMutex.unlock();                          _entry_mut.unlock();
                    q->enqueue(msg);                    q->enqueue(msg);
                    autoEntryMutex.lock();                          _entry_mut.lock();
            // After enqueue a message and the autoEntryMutex has been released and locked again,  
            // the array of entries can be changed. The ArrayIterator has be reset with the original _entries                          // After enqueue a message and the autoEntryMutex has
                           // been released and locked again, the array of
                           // entries can be changed. The ArrayIterator has be
                           // reset with the original _entries
            entries.reset(_entries);            entries.reset(_entries);
                    entries[indx]._status = _MonitorEntry::IDLE;                    entries[indx]._status = _MonitorEntry::IDLE;
   
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
                    {  
                        AutoMutex automut(Monitor::_cout_mut);  
                        PEGASUS_STD(cout) << "Exiting:  Monitor::run(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
                    }  
 #endif  
                    return true;  
                 }                 }
              }              }
              catch(...)              catch(...)
              {              {
              }              }
              handled_events = true;  
           }           }
        }        }
     }     }
   
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Exiting:  Monitor::run(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
     return(handled_events);  
 } }
  
 void Monitor::stopListeningForConnections(Boolean wait) void Monitor::stopListeningForConnections(Boolean wait)
 { {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::stopListeningForConnections(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
     PEG_METHOD_ENTER(TRC_HTTP, "Monitor::stopListeningForConnections()");     PEG_METHOD_ENTER(TRC_HTTP, "Monitor::stopListeningForConnections()");
     // set boolean then tickle the server to recognize _stopConnections     // set boolean then tickle the server to recognize _stopConnections
     _stopConnections = 1;     _stopConnections = 1;
Line 1115 
Line 680 
     }     }
  
     PEG_METHOD_EXIT();     PEG_METHOD_EXIT();
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Exiting:  Monitor::stopListeningForConnections(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
 } }
  
  
 int  Monitor::solicitSocketMessages( int  Monitor::solicitSocketMessages(
     PEGASUS_SOCKET socket,      SocketHandle socket,
     Uint32 events,     Uint32 events,
     Uint32 queueId,     Uint32 queueId,
     int type)     int type)
 { {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::solicitSocketMessages(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
    PEG_METHOD_ENTER(TRC_HTTP, "Monitor::solicitSocketMessages");    PEG_METHOD_ENTER(TRC_HTTP, "Monitor::solicitSocketMessages");
    AutoMutex autoMut(_entry_mut);    AutoMutex autoMut(_entry_mut);
    // Check to see if we need to dynamically grow the _entries array    // Check to see if we need to dynamically grow the _entries array
Line 1143 
Line 696 
    // current connections requested    // current connections requested
    _solicitSocketCount++;  // bump the count    _solicitSocketCount++;  // bump the count
    int size = (int)_entries.size();    int size = (int)_entries.size();
    if((int)_solicitSocketCount >= (size-1)){      if ((int)_solicitSocketCount >= (size-1))
         for(int i = 0; i < ((int)_solicitSocketCount - (size-1)); i++){      {
           for (int i = 0; i < ((int)_solicitSocketCount - (size-1)); i++)
           {
                 _MonitorEntry entry(0, 0, 0);                 _MonitorEntry entry(0, 0, 0);
                 _entries.append(entry);                 _entries.append(entry);
         }         }
Line 1162 
Line 717 
             _entries[index]._type = type;             _entries[index]._type = type;
             _entries[index]._status = _MonitorEntry::IDLE;             _entries[index]._status = _MonitorEntry::IDLE;
  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
             {  
                 AutoMutex automut(Monitor::_cout_mut);  
                 PEGASUS_STD(cout) << "Exiting:  Monitor::solicitSocketMessages(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
             }  
 #endif  
             return index;             return index;
          }          }
       }       }
Line 1175 
Line 724 
       {       {
       }       }
    }    }
    _solicitSocketCount--;  // decrease the count, if we are here we didnt do anything meaningful      // decrease the count, if we are here we didn't do anything meaningful
       _solicitSocketCount--;
    PEG_METHOD_EXIT();    PEG_METHOD_EXIT();
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
    {  
        AutoMutex automut(Monitor::_cout_mut);  
        PEGASUS_STD(cout) << "Exiting:  Monitor::solicitSocketMessages(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
    }  
 #endif  
    return -1;    return -1;
   
 } }
  
 void Monitor::unsolicitSocketMessages(PEGASUS_SOCKET socket)  void Monitor::unsolicitSocketMessages(SocketHandle socket)
 { {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::unsolicitSocketMessages(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
   
     PEG_METHOD_ENTER(TRC_HTTP, "Monitor::unsolicitSocketMessages");     PEG_METHOD_ENTER(TRC_HTTP, "Monitor::unsolicitSocketMessages");
     AutoMutex autoMut(_entry_mut);     AutoMutex autoMut(_entry_mut);
  
     /*     /*
         Start at index = 1 because _entries[0] is the tickle entry which never needs          Start at index = 1 because _entries[0] is the tickle entry which
         to be EMPTY;          never needs to be EMPTY;
     */     */
     unsigned int index;     unsigned int index;
     for(index = 1; index < _entries.size(); index++)     for(index = 1; index < _entries.size(); index++)
Line 1217 
Line 753 
  
     /*     /*
         Dynamic Contraction:         Dynamic Contraction:
         To remove excess entries we will start from the end of the _entries array          To remove excess entries we will start from the end of the _entries
         and remove all entries with EMPTY status until we find the first NON EMPTY.          array and remove all entries with EMPTY status until we find the
         This prevents the positions, of the NON EMPTY entries, from being changed.          first NON EMPTY.  This prevents the positions, of the NON EMPTY
           entries, from being changed.
     */     */
     index = _entries.size() - 1;     index = _entries.size() - 1;
     while(_entries[index]._status.get() == _MonitorEntry::EMPTY){      while (_entries[index]._status.get() == _MonitorEntry::EMPTY)
       {
         if(_entries.size() > MAX_NUMBER_OF_MONITOR_ENTRIES)         if(_entries.size() > MAX_NUMBER_OF_MONITOR_ENTRIES)
                 _entries.remove(index);                 _entries.remove(index);
         index--;         index--;
     }     }
     PEG_METHOD_EXIT();     PEG_METHOD_EXIT();
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Exiting:  Monitor::unsolicitSocketMessages(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
 } }
  
 // Note: this is no longer called with PEP 183. // Note: this is no longer called with PEP 183.
 PEGASUS_THREAD_RETURN PEGASUS_THREAD_CDECL Monitor::_dispatch(void *parm)  ThreadReturnType PEGASUS_THREAD_CDECL Monitor::_dispatch(void* parm)
 { {
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::_dispatch(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
    HTTPConnection *dst = reinterpret_cast<HTTPConnection *>(parm);    HTTPConnection *dst = reinterpret_cast<HTTPConnection *>(parm);
    Tracer::trace(TRC_HTTP, Tracer::LEVEL4,      PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
         "Monitor::_dispatch: entering run() for indx  = %d, queueId = %d, q = %p",          "Monitor::_dispatch: entering run() for indx  = %d, queueId = %d, "
         dst->_entry_index, dst->_monitor->_entries[dst->_entry_index].queueId, dst);              "q = %p",
           dst->_entry_index,
           dst->_monitor->_entries[dst->_entry_index].queueId,
           dst));
   
    try    try
    {    {
       dst->run(1);       dst->run(1);
    }    }
    catch (...)    catch (...)
    {    {
       Tracer::trace(TRC_HTTP, Tracer::LEVEL4,          PEG_TRACE_CSTRING(TRC_HTTP, Tracer::LEVEL4,
           "Monitor::_dispatch: exception received");           "Monitor::_dispatch: exception received");
    }    }
    Tracer::trace(TRC_HTTP, Tracer::LEVEL4,      PEG_TRACE((TRC_HTTP, Tracer::LEVEL4,
           "Monitor::_dispatch: exited run() for index %d", dst->_entry_index);          "Monitor::_dispatch: exited run() for index %d", dst->_entry_index));
  
    PEGASUS_ASSERT(dst->_monitor->_entries[dst->_entry_index]._status.get() == _MonitorEntry::BUSY);      PEGASUS_ASSERT(dst->_monitor->_entries[dst->_entry_index]._status.get() ==
           _MonitorEntry::BUSY);
  
    // Once the HTTPConnection thread has set the status value to either    // Once the HTTPConnection thread has set the status value to either
    // Monitor::DYING or Monitor::IDLE, it has returned control of the connection      // Monitor::DYING or Monitor::IDLE, it has returned control of the
    // to the Monitor.  It is no longer permissible to access the connection      // connection to the Monitor.  It is no longer permissible to access the
    // or the entry in the _entries table.      // connection or the entry in the _entries table.
    if (dst->_connectionClosePending)    if (dst->_connectionClosePending)
    {    {
       dst->_monitor->_entries[dst->_entry_index]._status = _MonitorEntry::DYING;          dst->_monitor->_entries[dst->_entry_index]._status =
               _MonitorEntry::DYING;
    }    }
    else    else
    {    {
       dst->_monitor->_entries[dst->_entry_index]._status = _MonitorEntry::IDLE;          dst->_monitor->_entries[dst->_entry_index]._status =
               _MonitorEntry::IDLE;
    }    }
    return 0;    return 0;
 } }
  
   
 //This method is anlogsu to solicitSocketMessages. It does the same thing for named Pipes  
 int  Monitor::solicitPipeMessages(  
     NamedPipe namedPipe,  
     Uint32 events,  //not sure what has to change for this enum  
     Uint32 queueId,  
     int type)  
 {  
    PEG_METHOD_ENTER(TRC_HTTP, "Monitor::solicitPipeMessages");  
   
    AutoMutex autoMut(_entry_mut);  
    // Check to see if we need to dynamically grow the _entries array  
    // We always want the _entries array to 2 bigger than the  
    // current connections requested  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
    AutoMutex automut(Monitor::_cout_mut);  
    PEGASUS_STD(cout) << "In Monitor::solicitPipeMessages at the begining" << PEGASUS_STD(endl);  
 #endif  
   
   
    _solicitSocketCount++;  // bump the count  
    int size = (int)_entries.size();  
    if((int)_solicitSocketCount >= (size-1)){  
         for(int i = 0; i < ((int)_solicitSocketCount - (size-1)); i++){  
                 _MonitorEntry entry(0, 0, 0);  
                 _entries.append(entry);  
         }  
    }  
   
    int index;  
    for(index = 1; index < (int)_entries.size(); index++)  
    {  
       try  
       {  
          if(_entries[index]._status.get() == _MonitorEntry::EMPTY)  
          {  
             _entries[index].socket = NULL;  
             _entries[index].namedPipe = namedPipe;  
             _entries[index].namedPipeConnection = true;  
             _entries[index].queueId  = queueId;  
             _entries[index]._type = type;  
             _entries[index]._status = _MonitorEntry::IDLE;  
   #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
             AutoMutex automut(Monitor::_cout_mut);  
             PEGASUS_STD(cout) << "In Monitor::solicitPipeMessages after seting up  _entries[index] index = " << index << PEGASUS_STD(endl);  
   #endif  
             return index;  
          }  
       }  
       catch(...)  
       {  
       }  
   
    }  
    _solicitSocketCount--;  // decrease the count, if we are here we didnt do anything meaningful  
    PEGASUS_STD(cout) << "In Monitor::solicitPipeMessages nothing happed - it didn't work" << PEGASUS_STD(endl);  
   
    PEG_METHOD_EXIT();  
    return -1;  
   
 }  
   
 void Monitor::unsolicitPipeMessages(NamedPipe namedPipe)  
 {  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Entering: Monitor::unsolicitPipeMessages(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
   
     PEG_METHOD_ENTER(TRC_HTTP, "Monitor::unsolicitPipeMessages");  
     AutoMutex autoMut(_entry_mut);  
   
     /*  
         Start at index = 1 because _entries[0] is the tickle entry which never needs  
         to be EMPTY;  
     */  
     unsigned int index;  
     for(index = 1; index < _entries.size(); index++)  
     {  
        if(_entries[index].namedPipe.getPipe() == namedPipe.getPipe())  
        {  
           _entries[index]._status = _MonitorEntry::EMPTY;  
           //_entries[index].namedPipe = PEGASUS_INVALID_SOCKET;  
           _solicitSocketCount--;  
           break;  
        }  
     }  
   
     /*  
         Dynamic Contraction:  
         To remove excess entries we will start from the end of the _entries array  
         and remove all entries with EMPTY status until we find the first NON EMPTY.  
         This prevents the positions, of the NON EMPTY entries, from being changed.  
     */  
     index = _entries.size() - 1;  
     while(_entries[index]._status.get() == _MonitorEntry::EMPTY){  
         if((_entries[index].namedPipe.getPipe() == namedPipe.getPipe()) ||  
             (_entries.size() > MAX_NUMBER_OF_MONITOR_ENTRIES))  
         {  
             _entries.remove(index);  
         }  
         index--;  
     }  
     PEG_METHOD_EXIT();  
 #ifdef PEGASUS_LOCALDOMAINSOCKET_DEBUG  
     {  
         AutoMutex automut(Monitor::_cout_mut);  
         PEGASUS_STD(cout) << "Exiting:  Monitor::unsolicitPipeMessages(): (tid:" << Uint32(pegasus_thread_self()) << ")" << PEGASUS_STD(endl);  
     }  
 #endif  
 }  
   
   
   
 PEGASUS_NAMESPACE_END PEGASUS_NAMESPACE_END


Legend:
Removed from v.1.103.10.22  
changed lines
  Added in v.1.124

No CVS admin address has been configured
Powered by
ViewCVS 0.9.2