wangzhengquan
2021-02-23 8df1ff06b931b0e414ed435a033f508867b345b7
src/socket/shm_socket.cpp
@@ -6,6 +6,7 @@
#include "bus_error.h"
#include "sole.h"
#include "shm_mm.h"
#include "key_def.h"
static Logger *logger = LoggerFactory::getLogger();
@@ -24,7 +25,7 @@
static int shm_recvpakfrom(shm_socket_t *sockt, shm_packet_t *recvpak ,  const struct timespec *timeout,  int flag);
   
static int shm_sendpakto(shm_socket_t *sockt, const shm_packet_t *sendpak,
static int shm_sendpakto(shm_socket_t *sockt,  shm_packet_t *sendpak,
               const int key, const struct timespec *timeout, const int flag);
@@ -109,26 +110,33 @@
}
int shm_socket_close(shm_socket_t *sockt) {
static int _shm_socket_close_(shm_socket_t *sockt) {
  
  int rv;
  logger->debug("shm_socket_close\n");
  // hashtable_remove(hashtable, mkey);
  // if(sockt->queue != NULL) {
  //   delete sockt->queue;
  //   sockt->queue = NULL;
  // }
  pthread_mutex_destroy(&(sockt->mutex) );
  free(sockt);
  auto it =  shmQueueStMap.find(key);
  if(it != shmQueueStMap.end()) {
    it->second.status = SHM_QUEUE_ST_CLOSED
    it->second.closeTime = time(NULL);
  if(sockt->key != 0) {
    auto it =  shmQueueStMap->find(sockt->key);
    if(it != shmQueueStMap->end()) {
      it->second.status = SHM_QUEUE_ST_CLOSED;
      it->second.closeTime = time(NULL);
    }
  }
  pthread_mutex_destroy(&(sockt->mutex) );
  free(sockt);
  return 0;
}
int shm_socket_close(shm_socket_t *sockt) {
  return _shm_socket_close_(sockt);
}
@@ -246,12 +254,8 @@
  if (rv != 0) {
    if(rv == ETIMEDOUT)
      return EBUS_TIMEOUT;
    else {
      logger->debug("%d shm_recvfrom failed %s", shm_socket_get_key(sockt), bus_strerror(rv));
      return rv;
    }
    logger->debug("%d shm_recvfrom failed %s", shm_socket_get_key(sockt), bus_strerror(rv));
    return rv;
   
  } 
@@ -268,6 +272,9 @@
  if(key != NULL)
    *key = recvpak.key;
  if(recvpak.key == 0) {
    err_exit(0, "key = %d, pid= %d, recvpak.key == 0",  shm_socket_get_key(sockt), getpid());
  }
  mm_free(recvpak.buf);
  return 0;
}
@@ -283,7 +290,7 @@
    return;
  logger->debug("%d destroy tmp socket\n", pthread_self()); 
  shm_socket_close((shm_socket_t *)tmp_socket);
  _shm_socket_close_((shm_socket_t *)tmp_socket);
  rv =  pthread_setspecific(_perthread_socket_key_, NULL);
  if ( rv != 0) {
    logger->error(rv, "shm_sendandrecv : pthread_setspecific");
@@ -311,7 +318,7 @@
                    const int send_size, const int key, void **recv_buf,
                    int *recv_size,  const struct timespec *timeout,  int flags) {
  int rv, tryn = 3;
  int rv, tryn = 6;
  shm_packet_t sendpak;
  shm_packet_t recvpak;
  std::map<std::string, shm_packet_t>::iterator recvbufIter;
@@ -391,15 +398,13 @@
                    int *recv_size,  const struct timespec *timeout,  int flags) {
  
 
  int rv, tryn = 6;
  int rv = 0, tryn = 16;
  shm_packet_t sendpak;
  shm_packet_t recvpak;
  std::map<int, shm_packet_t>::iterator recvbufIter;
  // 用thread local 保证每个线程用一个独占的socket接受对方返回的信息
  shm_socket_t *tmp_socket;
 
  /* If first call from this thread, allocate buffer for thread, and save its location */
  // logger->debug("%d create tmp socket\n", pthread_self() );
  rv = pthread_once(&_once_, _create_socket_key_perthread);
  if (rv != 0) {
    logger->error(rv, "shm_sendandrecv pthread_once");
@@ -420,7 +425,6 @@
    }
  }
 
  sendpak.key = tmp_socket->key;
  sendpak.size = send_size;
  if(send_buf != NULL) {
@@ -447,11 +451,6 @@
    rv = shm_recvpakfrom(tmp_socket, &recvpak, timeout, flags);
    if (rv != 0) {
      if(rv == ETIMEDOUT) {
        return EBUS_TIMEOUT;
      }
      logger->debug("%d shm_recvfrom failed %s", shm_socket_get_key(tmp_socket), bus_strerror(rv));
      return rv;
    } 
@@ -520,26 +519,27 @@
    return EBUS_RECVFROM_WRONG_END;
  } 
   
  shm_socket_close(tmp_socket);
  _shm_socket_close_(tmp_socket);
  return rv;
 
}
 
   
static int shm_sendpakto(shm_socket_t *sockt, const shm_packet_t *sendpak,
static int shm_sendpakto(shm_socket_t *sockt,  shm_packet_t *sendpak,
               const int key, const struct timespec *timeout, const int flag) {
  int rv;
  shm_queue_status_t stRecord;
  LockFreeQueue<shm_packet_t> *remoteQueue;
  hashtable_t *hashtable = mm_get_hashtable();
  if( sockt->queue != NULL) 
    goto LABEL_PUSH;
  if(hashtable_get_queue_count(hashtable) > QUEUE_COUNT_LIMIT) {
    return EBUS_EXCEED_LIMIT;
  }
  // if(hashtable_get_queue_count(hashtable) > QUEUE_COUNT_LIMIT) {
  //   return EBUS_EXCEED_LIMIT;
  // }
 
  {
    if ((rv = pthread_mutex_lock(&(sockt->mutex))) != 0)
@@ -552,13 +552,15 @@
      sockt->queue = shm_socket_bind_queue( sockt->key, sockt->force_bind);
      if(sockt->queue  == NULL ) {
        logger->error("%s. key = %d", bus_strerror(EBUS_KEY_INUSED), sockt->key);
        if ((rv = pthread_mutex_unlock(&(sockt->mutex))) != 0)
          err_exit(rv, "shm_sendto : pthread_mutex_unlock");
        return EBUS_KEY_INUSED;
      }
      // 标记key对应的状态 ,为opened
      stRecord.status = SHM_QUEUE_ST_OPENED;
      stRecord.createTime = time(NULL);
      shmQueueStMap.insert({sockt->key, stRecord});
      shmQueueStMap->insert({sockt->key, stRecord});
      
    }
@@ -575,21 +577,25 @@
  }
  // 检查key标记的状态
  auto it =  shmQueueStMap.find(key);
  if(it != shmQueueStMap.end()) {
  auto it =  shmQueueStMap->find(key);
  if(it != shmQueueStMap->end()) {
    if(it->second.status == SHM_QUEUE_ST_CLOSED) {
      // key对应的状态是关闭的
      goto ERR_CLOSED;
    }
  }
  LockFreeQueue<shm_packet_t> *remoteQueue = shm_socket_attach_queue(key);
  remoteQueue = shm_socket_attach_queue(key);
  if (remoteQueue == NULL ) {
    goto ERR_CLOSED;
  }
  sendpak->key = sockt->key;
  rv = remoteQueue->push(*sendpak, timeout, flag);
  if(rv == ETIMEDOUT) {
    return EBUS_TIMEOUT;
  }
  return rv;
ERR_CLOSED:
@@ -608,9 +614,9 @@
  if( sockt->queue != NULL) 
    goto LABEL_POP;
  if(hashtable_get_queue_count(hashtable) > QUEUE_COUNT_LIMIT) {
    return EBUS_EXCEED_LIMIT;
  }
  // if(hashtable_get_queue_count(hashtable) > QUEUE_COUNT_LIMIT) {
  //   return EBUS_EXCEED_LIMIT;
  // }
  {
    if ((rv = pthread_mutex_lock(&(sockt->mutex))) != 0)
@@ -623,13 +629,15 @@
    sockt->queue = shm_socket_bind_queue( sockt->key, sockt->force_bind);
    if(sockt->queue  == NULL ) {
      logger->error("%s. key = %d", bus_strerror(EBUS_KEY_INUSED), sockt->key);
      if ((rv = pthread_mutex_unlock(&(sockt->mutex))) != 0)
        err_exit(rv, "shm_recvfrom : pthread_mutex_unlock");
      return EBUS_KEY_INUSED;
    }
    
    // 标记key对应的状态 ,为opened
    stRecord.status = SHM_QUEUE_ST_OPENED;
    stRecord.createTime = time(NULL);
    shmQueueStMap.insert({sockt->key, stRecord});
    shmQueueStMap->insert({sockt->key, stRecord});
    
    if ((rv = pthread_mutex_unlock(&(sockt->mutex))) != 0)
      err_exit(rv, "shm_recvfrom : pthread_mutex_unlock");
@@ -638,25 +646,20 @@
  
LABEL_POP:
  // 检查key标记的状态
  // auto shmQueueMapIter =  shmQueueStMap.find(sockt->key);
  // if(shmQueueMapIter != shmQueueStMap.end()) {
  //   stRecord = shmQueueMapIter->second;
  //   if(stRecord.status = SHM_QUEUE_ST_CLOSED) {
  //     // key对应的状态是关闭的
  //     goto ERR_CLOSED;
  //   }
  // }
 
  rv = sockt->queue->pop(recvpak, timeout, flag);
  if(rv != 0)
  if(rv != 0) {
    if(rv == ETIMEDOUT) {
      return EBUS_TIMEOUT;
    }
    return rv;
  }
  
  if(recvpak.action == BUS_ACTION_STOP) {
    return EBUS_STOPED;
  }
  *_recvpak = recvpak;
  return rv;
  return 0;
}