Procházet zdrojové kódy

添加ntp超时机制

guohui před 1 rokem
rodič
revize
6cf8f1eed5
6 změnil soubory, kde provedl 169 přidání a 25 odebrání
  1. 2 11
      pro/src/app.c
  2. 1 1
      pro/src/common.c
  3. 99 4
      pro/src/ntpclient.c
  4. 9 3
      pro/src/ntpclient.h
  5. 57 6
      pro/src/sqlite_handle.c
  6. 1 0
      pro/src/sqlite_handle.h

+ 2 - 11
pro/src/app.c

@@ -29,15 +29,6 @@
 #include <unistd.h>
 #include "pduMIB_trap.h"
 
-static void get_time(char *ts)
-{
-    struct timeval tv;
-    struct tm      *tm;
-
-    gettimeofday(&tv, NULL);
-    tm = localtime(&tv.tv_sec);
-    sprintf(ts, "%04d-%02d-%02d %02d:%02d:%02d.%03ld", tm->tm_year+1900, tm->tm_mon+1, tm->tm_mday, tm->tm_hour, tm->tm_min, tm->tm_sec, tv.tv_usec/1000);
-}
 
 #include "file.h"
 #define REG_UPG_READ        1501
@@ -1811,7 +1802,7 @@ static void* ntp_thread(void *arg)
     }
     dev_Alarm_Run_message(dm, language_alarm_Init_Success[0], "NTP");
 
-
+    #define NTP_TIMEOUT   5
     #define NTP_INTERVAL_TIME  300
     time1 = time(NULL);
     while (h->quit == 0)
@@ -1820,7 +1811,7 @@ static void* ntp_thread(void *arg)
         {
             if (ntp_sync_flag || ((time1 - time2) >= NTP_INTERVAL_TIME))
             {
-                pc = createNTPClient(pntp->ntpaddress, pntp->port);
+                pc = createNTPClient(pntp->ntpaddress, pntp->port, NTP_TIMEOUT);
                 if (pc)
                 {
                     if (getCurrentTimeFromNTP(pc, &utc_time) == 0)

+ 1 - 1
pro/src/common.c

@@ -504,7 +504,7 @@ int set_utctime_zone(time_t t, int zone)
 int get_time_str(char *s, int len, time_t t)
 {
     struct tm *tm=localtime(&t);
-    return snprintf(s, len, "%04d-%02d-%02d %02d:%02d:%02d",tm->tm_year,tm->tm_mon,tm->tm_mday,tm->tm_hour,tm->tm_min,tm->tm_sec);
+    return snprintf(s, len, "%04d-%02d-%02d %02d:%02d:%02d.000",tm->tm_year,tm->tm_mon,tm->tm_mday,tm->tm_hour,tm->tm_min,tm->tm_sec);
 }
 
 

+ 99 - 4
pro/src/ntpclient.c

@@ -1,4 +1,9 @@
+#define _GNU_SOURCE
 #include <netdb.h>
+#include <signal.h>
+#include <unistd.h>
+#include <unistd.h>
+#include <fcntl.h>
 #include "ntpclient.h"
 
   
@@ -59,7 +64,49 @@ static int get_addr(ipaddr_t *addr, struct addrinfo *info)
     }
     return 0;
 }
-NTPClient* createNTPClient(const char* ntp_server, int port) {  
+
+
+//#define GETADDR_ASYNC
+#define GETADDR_SYNC_TIMEOUT
+
+#ifdef GETADDR_ASYNC
+static timer_t tmrId=NULL;
+static struct gaicb hostReq;
+static void timer_fn(union sigval value)
+{
+    printf("___ XXXXXXXXXXXXXXXXXXX\n");
+    gai_cancel(&hostReq);
+    timer_delete(tmrId);
+}
+static int start_timer(int flag)
+{
+    struct sigevent sev;
+    struct itimerspec its={0};
+
+    memset(&sev, 0, sizeof(sev));
+    sev.sigev_value.sival_ptr = &tmrId;
+    sev.sigev_notify = SIGEV_THREAD;
+    sev.sigev_notify_function = timer_fn;
+    sev.sigev_value.sival_int = 11;
+    sev.sigev_notify_attributes = NULL;
+    if (timer_create(CLOCK_REALTIME, &sev, &tmrId) == -1) {
+        perror("timer_create");
+        return -1;
+    }
+
+    its.it_value.tv_sec = 2;
+    its.it_value.tv_nsec = 0;
+    its.it_interval.tv_sec = 0;
+    its.it_interval.tv_nsec = 0;
+    if (timer_settime(tmrId, 0, &its, NULL) == -1) {
+        perror("timer_settime");
+        return -1;
+    }
+    return 0;
+}
+#endif
+
+NTPClient* createNTPClient(const char* ntp_server, int port, int timeout_sec) {  
     NTPClient* client = (NTPClient*)malloc(sizeof(NTPClient));  
     if (!client) {  
         perror("createNTPClient, malloc failed\n");  
@@ -78,15 +125,63 @@ NTPClient* createNTPClient(const char* ntp_server, int port) {
     hints.ai_protocol = 0;//IPPROTO_UDP;
 
     sprintf(tmp, "%d", port);
-    err = getaddrinfo(ntp_server, tmp, &hints, &res);
+#ifdef GETADDR_ASYNC
+    struct sigevent sig;
+    struct gaicb *pcb=&hostReq;
+
+    hostReq.ar_name = strdup(ntp_server);
+    hostReq.ar_service = strdup(tmp); //the port we will listen on
+    hostReq.ar_request = &hints;
+
+    sig.sigev_notify = SIGEV_SIGNAL;
+    sig.sigev_value.sival_ptr = &hostReq;
+    sig.sigev_signo = SIGRTMIN;
+
+    err = getaddrinfo_a(GAI_NOWAIT, &pcb, 1, &sig);
     if(err){
         printf(" getaddrinfo from %sfailed, %s\n", ntp_server, gai_strerror(err));  
         free(client); return NULL;
     }
+#else
+    
+#ifdef GETADDR_SYNC_TIMEOUT
+    fd_set fds;
+    struct timeval tv;
+    fd = socket(AF_INET, SOCK_DGRAM, 0);
+    if (fd < 0) {
+        perror("socket error");
+        free(client);
+        return NULL;
+    }
+    // 设置套接字为非阻塞模式
+    int flags = fcntl(fd, F_GETFL, 0);
+    fcntl(fd, F_SETFL, flags | O_NONBLOCK);
+#endif
+
+    err = getaddrinfo(ntp_server, tmp, &hints, &res);
+    if(err){
+        printf(" getaddrinfo from %s failed, %s\n", ntp_server, gai_strerror(err));
+        close(fd);
+        free(client); return NULL;
+    }
+
+#ifdef GETADDR_SYNC_TIMEOUT
+    FD_ZERO(&fds); FD_SET(fd, &fds);
+    tv.tv_sec = timeout_sec; tv.tv_usec = 0;
+    int n = select(fd, NULL, &fds, NULL, &tv);
+    if (n <= 0) { // 0:超时  <0:出错
+        close(fd);
+        freeaddrinfo(res);
+        free(client);
+        printf("___ %s !!\n", (n==0)?"timeout":"error");
+        return NULL;
+    }
+    close(fd);
+#endif
 
     for(rp=res; rp; rp=rp->ai_next){
         get_addr(&ipa, rp);
-        printf("22 socket to %s:%d ", ipa.ip, ipa.port);
+        printf("socket to %s:%d ", ipa.ip, ipa.port);
 
         fd = socket(rp->ai_family, rp->ai_socktype, 0);
         if (fd>=0) {
@@ -98,10 +193,10 @@ NTPClient* createNTPClient(const char* ntp_server, int port) {
         printf("failed!\n");
     }
     freeaddrinfo(res);
-
     if(fd<0) {
         free(client); return NULL;
     }
+#endif
     printf("___createNTPClient ok!\n");
   
     return client;  

+ 9 - 3
pro/src/ntpclient.h

@@ -13,14 +13,20 @@
     
 #define NTP_PORT 123  
 #define NTP_PACKET_SIZE 48  
-// 定义NTP客户端类  
+
+
+typedef struct {
+    timer_t    tmr;
+}timer_handle_t;
+
 typedef struct NTPClient_s {  
     int sockfd;  
     struct sockaddr server_addr;  
-    char ntp_packet[NTP_PACKET_SIZE];  
+    char ntp_packet[NTP_PACKET_SIZE]; 
+    timer_t tmr; 
 } NTPClient;  
     
-NTPClient* createNTPClient(const char* ntp_server, int port);
+NTPClient* createNTPClient(const char* ntp_server, int port, int timeout_sec);
 
 // 销毁NTP客户端实例  
 void destroyNTPClient(NTPClient* client) ;

+ 57 - 6
pro/src/sqlite_handle.c

@@ -1262,8 +1262,7 @@ int dev_search_power_info(sqlite3 *db,
     char** pResult = NULL;  //用来指向sql执行结果的指针
     char* err_msg=NULL;
  
-    sprintf(select_sql,"SELECT* FROM Table_PowerInfo WHERE (product_id=%d AND product_ch_id=%d) AND (product_samp_time BETWEEN '%s' AND '%s') LIMIT %d OFFSET %d;",
-                        product_id,
+    sprintf(select_sql,"SELECT* FROM Table_PowerInfo WHERE (product_ch_id=%d) AND (product_samp_time BETWEEN '%s' AND '%s') LIMIT %d OFFSET %d;",
                         product_ch_id,
                         start_time,
                         stop_time,
@@ -4670,15 +4669,20 @@ const char *tabNamePool[TAB_ID_MAX]={
 
 };
 
-static int sql_requery(sqlite3 *db, int tab_id, char *condition, void *data, int *cnt)
+static int sql_requery(sqlite3 *db, int tab_id, time_t from, time_t to, void *data, int *cnt)
 {
     int r=-1;
     char *p,temp[1024];
+    char time_s[64],time_e[64];
     sqlite3_stmt *stmt=NULL;
+    GlobalPowerInfo *pwrInfo;
 
-    if(tab_id==TAB_ID_POWER_INFO) {
-        //snprintf(temp, sizeof(temp), "SELECT * FROM %s %s;", tabNamePool[tab_id], product_ch_id, start_time, end_time);
+    get_time_str(time_s, sizeof(time_s), from);
+    get_time_str(time_e, sizeof(time_e), to);
 
+    if(tab_id==TAB_ID_POWER_INFO) {
+        snprintf(temp, sizeof(temp), "SELECT * FROM %s WHERE (product_samp_time BETWEEN '%s' AND '%s';", tabNamePool[tab_id], time_s, time_e);
+        
         r = sqlite3_prepare_v2(db, temp, -1, &stmt, 0);
         if (r != SQLITE_OK) {
             log_e("___mqtt_load, sqlite3_prepare_v2 failed\n");
@@ -4716,7 +4720,54 @@ static int sql_requery(sqlite3 *db, int tab_id, char *condition, void *data, int
 
     }
     else if(tab_id==TAB_ID_POWER3_INFO) {
-        //
+        snprintf(temp, sizeof(temp), "SELECT * FROM %s WHERE (product_ph_samp_time BETWEEN '%s' AND '%s') ORDER BY product_index DESC;", tabNamePool[tab_id], time_s, time_e);
+
+        r = sqlite3_prepare_v2(db, temp, -1, &stmt, 0);
+        if (r != SQLITE_OK) {
+            log_e("___mqtt_load, sqlite3_prepare_v2 failed\n");
+            return -1;
+        }
+    
+#if 0
+        strcpy(_globalPowerInfoTemp->samp_time,pResult[i*ncolumn+0]);//i=0 10 
+        _globalPowerInfoTemp->product_id = atoi(pResult[i*ncolumn+1]);
+        _globalPowerInfoTemp->product_ch_id =  atoi(pResult[i*ncolumn+2]);
+        _globalPowerInfoTemp->product_status =  atoi(pResult[i*ncolumn+3]);
+        _globalPowerInfoTemp->_power_info.voltage = atof(pResult[i*ncolumn+4]);
+        _globalPowerInfoTemp->_power_info.current = atof(pResult[i*ncolumn+5]);
+        _globalPowerInfoTemp->_power_info.power = atof(pResult[i*ncolumn+6]);
+        _globalPowerInfoTemp->_power_info.freq = atof(pResult[i*ncolumn+7]);
+        _globalPowerInfoTemp->_power_info.consumption = atof(pResult[i*ncolumn+8]);
+        _globalPowerInfoTemp->_power_info.factor = atof(pResult[i*ncolumn+9]);
+#endif
+        while (sqlite3_step(stmt) == SQLITE_ROW) {
+            //ser[idx].id = sqlite3_column_int(stmt, 0);
+            //ser[idx].mode = sqlite3_column_int(stmt, 1);
+
+            p = (char*)sqlite3_column_text(stmt, 2);
+            //if(p) strcpy(ser[idx].name, p);
+
+            p = (char*)sqlite3_column_text(stmt, 3);
+            //if(p) strcpy(ser[idx].server, p);
+
+            p = (char*)sqlite3_column_text(stmt, 4);
+           // if(p) strcpy(ser[idx].port, p);
+
+            p = (char*)sqlite3_column_text(stmt, 5);
+            //if(p) strcpy(ser[idx].cid, p);
+
+            p = (char*)sqlite3_column_text(stmt, 6);
+            //if(p) strcpy(ser[idx].user, p);
+
+            p = (char*)sqlite3_column_text(stmt, 7);
+            //if(p) strcpy(ser[idx].password, p);
+
+            p = (char*)sqlite3_column_text(stmt, 8);
+            //if(p) strcpy(ser[idx].cert, p);
+
+            
+        }
+        sqlite3_finalize(stmt);
     }
     
     return r;

+ 1 - 0
pro/src/sqlite_handle.h

@@ -196,6 +196,7 @@ enum {
     
     TAB_ID_MAX
 };
+
 int dev_data_xport(sqlite3 *db, int tab_id, int io_type, uint32_t time_s, uint32_t time_e, void *data, int *cnt);