#include <sys/select.h>
#include <sys/time.h>
#include <sys/types.h>
#include <unistd.h>
#include <sys/syscall.h>
#include <stdlib.h>
#include <stdio.h>
#include <string.h>
#include <sys/stat.h>
#include <sys/types.h>
#include <errno.h>

#include "s7plc_common.h"
#if 1
/*=============================================================================|
|  PROJECT SNAP7                                                         1.4.0 |
|==============================================================================|
|  Copyright (C) 2013, 2014 Davide Nardella                                    |
|  All rights reserved.                                                        |
|==============================================================================|
|  SNAP7 is free software: you can redistribute it and/or modify               |
|  it under the terms of the Lesser GNU General Public License as published by |
|  the Free Software Foundation, either version 3 of the License, or           |
|  (at your option) any later version.                                         |
|                                                                              |
|  It means that you can distribute your commercial software linked with       |
|  SNAP7 without the requirement to distribute the source code of your         |
|  application and without the requirement that your application be itself     |
|  distributed under LGPL.                                                     |
|                                                                              |
|  SNAP7 is distributed in the hope that it will be useful,                    |
|  but WITHOUT ANY WARRANTY; without even the implied warranty of              |
|  MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the               |
|  Lesser GNU General Public License for more details.                         |
|                                                                              |
|  You should have received a copy of the GNU General Public License and a     |
|  copy of Lesser GNU General Public License along with Snap7.                 |
|  If not, see  http://www.gnu.org/licenses/                                   |
|==============================================================================|
|                                                                              |
|  Client Example                                                              |
|                                                                              |
|=============================================================================*/
#include <stdio.h>
#include <stdlib.h>
#include <ctype.h>
#include "snap7.h"


#ifdef OS_WINDOWS
# define WIN32_LEAN_AND_MEAN
# include <windows.h>
#endif
    S7Object Client;
    unsigned char Buffer[65536]; // 64 K buffer
    int SampleDBNum = 1000;

    char *Address;     // PLC IP Address
    int Rack=0,Slot=2; // Default Rack and Slot

    int ok = 0; // Number of test pass
    int ko = 0; // Number of test failure

    int JobDone=false;
    int JobResult=0;

//------------------------------------------------------------------------------
//  Async completion callback 
//------------------------------------------------------------------------------
// This is a simply text demo, we use callback only to set an internal flag...
void S7API CliCompletion(void *usrPtr, int opCode, int opResult)
{
    JobResult=opResult;
    JobDone = true;
}
//------------------------------------------------------------------------------
// SysSleep (copied from snap_sysutils.cpp) multiplatform millisec sleep
//------------------------------------------------------------------------------
void SysSleep(longword Delay_ms)
{
#ifdef OS_WINDOWS
	Sleep(Delay_ms);
#else
    struct timespec ts;
    ts.tv_sec = (time_t)(Delay_ms / 1000);
    ts.tv_nsec =(long)((Delay_ms - ts.tv_sec) * 1000000);
    nanosleep(&ts, (struct timespec *)0);
#endif
}
//------------------------------------------------------------------------------
//  Usage Syntax
//------------------------------------------------------------------------------
void Usage()
{
    printf("Usage\n");
    printf("  client <IP> [Rack=0 Slot=2]\n");
    printf("Example\n");
    printf("  client 192.168.1.101 0 2\n");
    printf("or\n");
    printf("  client 192.168.1.101\n");
    getchar();
}
//------------------------------------------------------------------------------
// hexdump, a very nice function, it's not mine.
// I found it on the net somewhere some time ago... thanks to the author ;-)
//------------------------------------------------------------------------------
#ifndef HEXDUMP_COLS
#define HEXDUMP_COLS 16
#endif
void hexdump(void *mem, unsigned int len)
{
        unsigned int i, j;

        for(i = 0; i < len + ((len % HEXDUMP_COLS) ? (HEXDUMP_COLS - len % HEXDUMP_COLS) : 0); i++)
        {
                /* print offset */
                if(i % HEXDUMP_COLS == 0)
                {
                        printf("0x%04x: ", i);
                }

                /* print hex data */
                if(i < len)
                {
                        printf("%02x ", 0xFF & ((char*)mem)[i]);
                }
                else /* end of block, just aligning for ASCII dump */
                {
                        printf("   ");
                }

                /* print ASCII dump */
                if(i % HEXDUMP_COLS == (HEXDUMP_COLS - 1))
                {
                        for(j = i - (HEXDUMP_COLS - 1); j <= i; j++)
                        {
                                if(j >= len) /* end of block, not really printing */
                                {
                                        putchar(' ');
                                }
                                else if(isprint((((char*)mem)[j] & 0x7F))) /* printable char */
                                {
                                        putchar(0xFF & ((char*)mem)[j]);
                                }
                                else /* other char */
                                {
                                        putchar('.');
                                }
                        }
                        putchar('\n');
                }
        }
}
//------------------------------------------------------------------------------
// Check error
//------------------------------------------------------------------------------
int Check(int Result, char * function)
{
    int ExecTime;
	char text[1024];
    printf("\n");
    printf("+-----------------------------------------------------\n");
    printf("| %s\n",function);
    printf("+-----------------------------------------------------\n");
    if (Result==0) {
        Cli_GetExecTime(Client, &ExecTime);
        printf("| Result         : OK\n");
        printf("| Execution time : %d ms\n",ExecTime);
        printf("+-----------------------------------------------------\n");
        ok++;
    }
    else {
        printf("| ERROR !!! \n");
        if (Result<0)
            printf("| Library Error (-1)\n");
        else
		{
			Cli_ErrorText(Result, text, 1024);
			printf("| %s\n",text);
		}
        printf("+-----------------------------------------------------\n");
        ko++;
    }
    return !Result;
}
//------------------------------------------------------------------------------
// Multi Read
//------------------------------------------------------------------------------
void MultiRead()
{
     int res;

     // Multiread buffers
     byte MB[16]; // 16 Merker bytes
     byte EB[16]; // 16 Digital Input bytes
     byte AB[16]; // 16 Digital Output bytes
     word TM[8];  // 8 timers
     word CT[8];  // 8 counters

     // Prepare struct
     TS7DataItem Items[5];

     // NOTE : *AMOUNT IS NOT SIZE* , it's the number of items

     // Merkers
     Items[0].Area     =S7AreaMK;
     Items[0].WordLen  =S7WLByte;
     Items[0].DBNumber =0;        // Don't need DB
     Items[0].Start    =0;        // Starting from 0
     Items[0].Amount   =16;       // 16 Items (bytes)
     Items[0].pdata    =&MB;
     // Digital Input bytes
     Items[1].Area     =S7AreaPE;
     Items[1].WordLen  =S7WLByte;
     Items[1].DBNumber =0;        // Don't need DB
     Items[1].Start    =0;        // Starting from 0
     Items[1].Amount   =16;       // 16 Items (bytes)
     Items[1].pdata    =&EB;
     // Digital Output bytes
     Items[2].Area     =S7AreaPA;
     Items[2].WordLen  =S7WLByte;
     Items[2].DBNumber =0;        // Don't need DB
     Items[2].Start    =0;        // Starting from 0
     Items[2].Amount   =16;       // 16 Items (bytes)
     Items[2].pdata    =&AB;
     // Timers
     Items[3].Area     =S7AreaTM;
     Items[3].WordLen  =S7WLTimer;
     Items[3].DBNumber =0;        // Don't need DB
     Items[3].Start    =0;        // Starting from 0
     Items[3].Amount   =8;        // 8 Timers
     Items[3].pdata    =&TM;
     // Counters
     Items[4].Area     =S7AreaCT;
     Items[4].WordLen  =S7WLCounter;
     Items[4].DBNumber =0;        // Don't need DB
     Items[4].Start    =0;        // Starting from 0
     Items[4].Amount   =8;        // 8 Counters
     Items[4].pdata    =&CT;

     res=Cli_ReadMultiVars(Client, &Items[0], 5);
     if (Check(res,"Multiread Vars"))
     {
        // Result of Client->ReadMultivars is the "global result" of
        // the function, it's OK if something was exchanged.

        // But we need to check single Var results.
        // Let shall suppose that we ask for 5 vars, 4 of them are ok but
        // the 5th is inexistent, we will have 4 results ok and 1 not ok.

        printf("Dump MB0..MB15 - Var Result : %d\n",Items[0].Result);
        if (Items[0].Result==0)
            hexdump(&MB,16);
        printf("Dump EB0..EB15 - Var Result : %d\n",Items[1].Result);
        if (Items[1].Result==0)
            hexdump(&EB,16);
        printf("Dump AB0..AB15 - Var Result : %d\n",Items[2].Result);
        if (Items[2].Result==0)
            hexdump(&AB,16);
        printf("Dump T0..T7 - Var Result : %d\n",Items[3].Result);
        if (Items[3].Result==0)
            hexdump(&TM,16);         // 8 Timers -> 16 bytes
        printf("Dump Z0..Z7 - Var Result : %d\n",Items[4].Result);
        if (Items[4].Result==0)
            hexdump(&CT,16);         // 8 Counters -> 16 bytes
     };
}
//------------------------------------------------------------------------------
// List blocks in AG
//------------------------------------------------------------------------------
void ListBlocks()
{
    TS7BlocksList List;
    int res=Cli_ListBlocks(Client, &List);
    if (Check(res,"List Blocks in AG"))
    {
        printf("  OBCount  : %d\n",List.OBCount);
	    printf("  FBCount  : %d\n",List.FBCount);
   		printf("  FCCount  : %d\n",List.FCCount);
   		printf("  SFBCount : %d\n",List.SFBCount);
   		printf("  SFCCount : %d\n",List.SFCCount);
   		printf("  DBCount  : %d\n",List.DBCount);
   		printf("  SDBCount : %d\n",List.SDBCount);
    };
}
//------------------------------------------------------------------------------
// CPU Info : catalog
//------------------------------------------------------------------------------
void OrderCode()
{
     TS7OrderCode Info;
     int res=Cli_GetOrderCode(Client, &Info);
     if (Check(res,"Catalog"))
     {
          printf("  Order Code : %s\n",Info.Code);
          printf("  Version    : %d.%d.%d\n",Info.V1,Info.V2,Info.V3);
     };
}
//------------------------------------------------------------------------------
// CPU Info : unit info
//------------------------------------------------------------------------------
void CpuInfo()
{
     TS7CpuInfo Info;
     int res=Cli_GetCpuInfo(Client, &Info);
     if (Check(res,"Unit Info"))
     {
          printf("  Module Type Name : %s\n",Info.ModuleTypeName);
          printf("  Seriel Number    : %s\n",Info.SerialNumber);
          printf("  AS Name          : %s\n",Info.ASName);
          printf("  Module Name      : %s\n",Info.ModuleName);
     };
}
//------------------------------------------------------------------------------
// CP Info
//------------------------------------------------------------------------------
void CpInfo()
{
     TS7CpInfo Info;
     int res=Cli_GetCpInfo(Client, &Info);
     if (Check(res,"Communication processor Info"))
     {
          printf("  Max PDU Length   : %d bytes\n",Info.MaxPduLengt);
          printf("  Max Connections  : %d \n",Info.MaxConnections);
          printf("  Max MPI Rate     : %d bps\n",Info.MaxMpiRate);
          printf("  Max Bus Rate     : %d bps\n",Info.MaxBusRate);
     };
}
//------------------------------------------------------------------------------
// PLC Status
//------------------------------------------------------------------------------
void UnitStatus()
{
     int res=0;

     int Status;
     Cli_GetPlcStatus(Client, &Status);
     if (Check(res,"CPU Status"))
     {
          switch (Status)
          {
              case S7CpuStatusRun : printf("  RUN\n"); break;
              case S7CpuStatusStop: printf("  STOP\n"); break;
              default             : printf("  UNKNOWN\n"); break;
          }
     };

}
//------------------------------------------------------------------------------
// Upload DB0 (surely exists in AG)
//------------------------------------------------------------------------------
void UploadDB0()
{
     int Size = sizeof(Buffer); // Size is IN/OUT par
                                // In input it tells the client the size available
                                // In output it tells us how many bytes were uploaded.
     int res=Cli_Upload(Client, Block_SDB, 0, &Buffer, &Size);
     if (Check(res,"Block Upload (SDB 0)"))
     {
          printf("Dump (%d bytes) :\n",Size);
          hexdump(&Buffer,Size);
     }
}
//------------------------------------------------------------------------------
// Async Upload DB0 (using callback as completion trigger)
//------------------------------------------------------------------------------
void AsCBUploadDB0()
{
     int Size = sizeof(Buffer); // Size is IN/OUT par
                                // In input it tells the client the size available
                                // In output it tells us how many bytes were uploaded.
     int res;
	 JobDone=false;
     
	 res=Cli_AsUpload(Client, Block_SDB, 0, &Buffer, &Size);
     
     if (res==0)
     {
         while (!JobDone)
         {
             SysSleep(100);
         }
         res=JobResult;
     }    

     if (Check(res,"Async (callback) Block Upload (SDB 0)"))
     {
          printf("Dump (%d bytes) :\n",Size);
          hexdump(&Buffer,Size);
     }
}
//------------------------------------------------------------------------------
// Async Upload DB0 (using event wait as completion trigger)
//------------------------------------------------------------------------------
void AsEWUploadDB0()
{
     int Size = sizeof(Buffer); // Size is IN/OUT par
                                // In input it tells the client the size available
                                // In output it tells us how many bytes were uploaded.
     int res;
	 JobDone=false;
     
	 res=Cli_AsUpload(Client, Block_SDB, 0, &Buffer, &Size);
     
     if (res==0)
     {
         res=Cli_WaitAsCompletion(Client,3000);
     }    

     if (Check(res,"Async (Wait event) Block Upload (SDB 0)"))
     {
          printf("Dump (%d bytes) :\n",Size);
          hexdump(&Buffer,Size);
     }
}
//------------------------------------------------------------------------------
// Async Upload DB0 (using polling as completion trigger)
//------------------------------------------------------------------------------
void AsPOUploadDB0()
{
     int Size = sizeof(Buffer); // Size is IN/OUT par
                                // In input it tells the client the size available
                                // In output it tells us how many bytes were uploaded.
     int res;
	 JobDone=false;
     
	 res=Cli_AsUpload(Client, Block_SDB, 0, &Buffer, &Size);
     
     if (res==0)
     {
		 while (Cli_CheckAsCompletion(Client,&res)!=JobComplete)
         {
             SysSleep(100);
         };         
     }    

     if (Check(res,"Async (polling) Block Upload (SDB 0)"))
     {
          printf("Dump (%d bytes) :\n",Size);
          hexdump(&Buffer,Size);
     }
}
//------------------------------------------------------------------------------
// Read a sample SZL Block
//------------------------------------------------------------------------------
void ReadSzl_0011_0000()
{
     PS7SZL SZL = (PS7SZL)(&Buffer);  // use our buffer casted as TS7SZL
     int Size = sizeof(Buffer);
     // Block ID 0x0011 IDX 0x0000 normally exists in every CPU
     int res=Cli_ReadSZL(Client, 0x0011, 0x0000, SZL, &Size);
     if (Check(res,"Read SZL - ID : 0x0011, IDX 0x0000"))
     {
        printf("  LENTHDR : %d\n",SZL->Header.LENTHDR);
        printf("  N_DR    : %d\n",SZL->Header.N_DR);
        printf("Dump (%d bytes) :\n",Size);
        hexdump(&Buffer,Size);
     }
}
//------------------------------------------------------------------------------
// Unit Connection
//------------------------------------------------------------------------------
int CliConnect()
{
    int Requested, Negotiated, res;

    res = Cli_ConnectTo(Client, Address,Rack,Slot);
    if (Check(res,"UNIT Connection")) {
          Cli_GetPduLength(Client, &Requested, &Negotiated);
          printf("  Connected to   : %s (Rack=%d, Slot=%d)\n",Address,Rack,Slot);
          printf("  PDU Requested  : %d bytes\n",Requested);
          printf("  PDU Negotiated : %d bytes\n",Negotiated);
    };
    return !res;
}
//------------------------------------------------------------------------------
// Unit Disconnection
//------------------------------------------------------------------------------
void CliDisconnect()
{
     Cli_Disconnect(Client);
}
//------------------------------------------------------------------------------
// Perform readonly tests, no cpu status modification
//------------------------------------------------------------------------------
void PerformTests()
{
     OrderCode();
     CpuInfo();
     CpInfo();
     UnitStatus();
     ReadSzl_0011_0000();
     UploadDB0();
     AsCBUploadDB0();
     AsEWUploadDB0();
     AsPOUploadDB0();
     MultiRead();
}

void Summary()
{
    printf("\n");
    printf("+-----------------------------------------------------\n");
    printf("| Test Summary \n");
    printf("+-----------------------------------------------------\n");
    printf("| Performed : %d\n",(ok+ko));
    printf("| Passed    : %d\n",ok);
    printf("| Failed    : %d\n",ko);
    printf("+----------------------------------------[press a key]\n");
    getchar();
}

int test_t(void)
{
#if 0
// Get Progran args
    if (argc!=2 && argc!=4)
    {
        Usage();
        return 1;
    }
    Address=argv[1];
    if (argc==4)
    {
        Rack=atoi(argv[2]);
        Slot=atoi(argv[3]);
    }
#endif
	Address = "192.168.0.1";
	Rack = 0;
	Slot = 1;
// Client Creation
	Client=Cli_Create();
    Cli_SetAsCallback(Client, CliCompletion,NULL);

// Connection
    if (CliConnect())
    {
        PerformTests();
        CliDisconnect();
    };

// Deletion
    Cli_Destroy(&Client);
    Summary();

    return 0;
}

#else
int s7_plc_period_rs_triger(s7_plc_var_t *var);

gw_port_e ports[] = {S7_PLC};

int s7_plc_client_init(s7_plc_var_t *var, tcp_node_list_t *tcp_node, char *ip_addr, int port)
{
    uint32_t old_response_to_sec;
    uint32_t old_response_to_usec;
    uint32_t new_response_to_sec;
    uint32_t new_response_to_usec;

    tcp_node->ctx = modbus_new_tcp(ip_addr, port);
    if (tcp_node->ctx == NULL)
    {
        dy_syslog(LOG_ERR, "Unable to allocate libmodbus context\n");
        return -1;
    }
    modbus_set_debug(tcp_node->ctx, TRUE);
    modbus_set_error_recovery(tcp_node->ctx,
                              MODBUS_ERROR_RECOVERY_LINK |
                              MODBUS_ERROR_RECOVERY_PROTOCOL);

    modbus_get_response_timeout(tcp_node->ctx, &old_response_to_sec, &old_response_to_usec);
    if (modbus_connect(tcp_node->ctx) == -1)
    {
        dy_syslog(LOG_ERR, "Connection failed: %s\n", modbus_strerror(errno));
        modbus_free(tcp_node->ctx);
        tcp_node->ctx = NULL;
        return -1;
    }
    else
        dy_syslog(LOG_DEBUG, "ip_add %s port %d Connection succeed ctx %p...", ip_addr, port, tcp_node->ctx);
    modbus_get_response_timeout(tcp_node->ctx, &new_response_to_sec, &new_response_to_usec);

    return 0;
}

static int s7_plc_msg_tag_data_ctrl(s7_plc_var_t *var, mqtt_message_t *mqtt_msg)
{
    char find = 0;
    pp_regulate_signal_t regulate_data;
    s7_plc_data_list_t *s7_plc_data;
    int ret;
    int port;

    //先检查参数
    ret = get_rglt_data_from_json(&regulate_data, mqtt_msg->payload);
    if (ret < 0 || regulate_data.len <= 0)
    {
        dy_syslog(LOG_ERR, "parse real data structure failed");
        return -1;
    }

    s7_plc_data = calloc(sizeof(s7_plc_data_list_t), 1);
    if (s7_plc_data)
    {
        memcpy(&s7_plc_data->rs, &regulate_data, sizeof(regulate_data));

        {
            tcp_node_list_t *tcp_node = NULL;
            list_for_each_entry(tcp_node, &var->tcp_node_list, list)
            {
                var->period_triggered = 1;
                s7_plc_data_list_t *cmd = NULL;
                dy_syslog(LOG_DEBUG, "tcp_port(%d %d) tcp_ip_addr(%s %s) \n", tcp_node->tcp_port, s7_plc_data->rs.tcp_port, tcp_node->tcp_ip_addr, s7_plc_data->rs.tcp_ip_addr);
                if (tcp_node->tcp_port == s7_plc_data->rs.tcp_port && strcmp(tcp_node->tcp_ip_addr, s7_plc_data->rs.tcp_ip_addr) == 0)
                {
                    find = 0;
                    list_for_each_entry(cmd, &tcp_node->data_list, list)
                    {
                        if (strcmp(cmd->rs.sn, regulate_data.sn) == 0 &&
                                cmd->rs.period == s7_plc_data->rs.period &&
                                cmd->rs.len == s7_plc_data->rs.len &&
                                (memcmp(cmd->rs.data, s7_plc_data->rs.data, cmd->rs.len) == 0))
                        {
                            find = 1;
                            break;
                        }
                    }
                    if (!find)
                    {
                        list_add_tail(&s7_plc_data->list, &tcp_node->data_list);
                        dy_syslog(LOG_DEBUG, "list_add_tail tcp_port(%d) tcp_ip_addr(%s) sn:%s\n", s7_plc_data->rs.tcp_port, s7_plc_data->rs.tcp_ip_addr, s7_plc_data->rs.sn);
                    }
                }
            }
        }

        dy_syslog(LOG_INFO, "sn %s data[0] %d identifier %s port %d period %d", s7_plc_data->rs.sn, s7_plc_data->rs.data[0], regulate_data.src_identifier, S7_PLC, regulate_data.period);
        if (find)
        {
            dy_syslog(LOG_WARNING, "Repeated commands !!!");
            free(s7_plc_data);
        }
    }

    return 0;
}

int s7_plc_tag_data_report(s7_plc_var_t *var, char *data, int len, unsigned char port)
{
    pp_real_time_data_t real_data = {0};
    regulate_cmd_backup_t *fields = NULL;

    //restore the identifier
    {
        char key_str[64] = {0};
        data_info_t *data_info = NULL;

        sprintf(key_str, "S7_PLC");
        data_info = var->identifier_backup.get(&var->identifier_backup, key_str);
        if (data_info)
        {
            fields = (regulate_cmd_backup_t *)data_info->data;
        }
    }
    real_data.port = port;
    real_data.len = len;
    real_data.instruction_code = data[1];
    real_data.data = malloc(real_data.len);
    if (fields != NULL)
    {
        real_data.mi = fields->mi;
        strcpy(real_data.sn, fields->sn);
        strncpy(real_data.src_identifier, fields->src_identifier, sizeof(real_data.src_identifier));
    }

    memcpy(real_data.data, data, real_data.len);

    send_to_proto_parser(var->session, &real_data);
    flush_device_status_by_sn(var->ds, ACTION_RCV, real_data.sn);

    free(real_data.data);

    return 0;
}

void s7_plc_bus_send(s7_plc_var_t *var, pp_regulate_signal_t *rs, unsigned char port, unsigned int communication_timeout, tcp_node_list_t *tcp_node)
{
    int ret;
    char *data = rs->data;
    int len = rs->len;
    long long cur_ms = clock_get_ms();
    int try_cnt = 0;

    //store the identifier
    {
        char key_str[64] = {0};
        regulate_cmd_backup_t fields;

        strcpy(fields.src_identifier, rs->src_identifier);
        strcpy(fields.sn, rs->sn);
        fields.mi = rs->mi;

        sprintf(key_str, "S7_PLC");
        var->identifier_backup.insert(&var->identifier_backup, key_str, &fields, sizeof(fields));
    }
    flush_device_status_by_sn(var->ds, ACTION_SND, rs->sn);
s7_plc_flag:
    if (port == S7_PLC)
    {
        //tcp发送数据,阻塞等待执行数据返回
        if (!tcp_node->ctx)
            s7_plc_client_init(var, tcp_node, tcp_node->tcp_ip_addr, tcp_node->tcp_port);

        if (tcp_node->ctx)
        {
            char read_nBytes[512] = {0};

            len = rs->len - 2;
            dy_syslog_hex(LOG_DEBUG, data, len, "sn:%s  len %d  identifier:%s    网关->PORT[%d]", rs->sn, len, rs->src_identifier, port);
            int rc = modbus_communication_transparent(tcp_node->ctx, data, len, read_nBytes);
            if (rc > 0)
            {
                dy_syslog_hex(LOG_DEBUG, read_nBytes, rc, "	len %d			网关<-PORT[%d]		", rc, port);
                s7_plc_tag_data_report(var, read_nBytes, rc, port);
            }
            else
            {
                dy_syslog(LOG_DEBUG, "ctx %p try_cnt %d rc %d ", tcp_node->ctx, try_cnt, rc);
                if (try_cnt < 2)
                {
                    try_cnt += 1;
                    modbus_close(tcp_node->ctx);
                    usleep(200000);
                    modbus_connect(tcp_node->ctx);
                    goto s7_plc_flag;
                }
                else
                {
                    modbus_close(tcp_node->ctx);
                    free(tcp_node->ctx);
                    tcp_node->ctx = NULL;
                }
            }
        }
        else
            dy_syslog(LOG_WARNING, "tcp_ip_addr:%s port:%d ctx is NULL!!!", tcp_node->tcp_ip_addr, tcp_node->tcp_port);
    }
}

static int ds_flush_timer(s7_plc_var_t *var)
{
    time_t now = time(NULL);

    TIMER_CONFIRM(var->node_status_timer);
    if (var->period_triggered == 0)
    {
        s7_plc_period_rs_triger(var);
    }

    if (now >= var->nodes_status_sync_date + 30)
    {
        static time_t lasttime = 0;
        var->nodes_status_sync_date = now;
        char *node_status = ds_print(var->ds);

        dy_syslog(LOG_INFO, "online_status_changed:%d strlen(node_status):%d, (now-lasttime):%d",
                  var->ds->online_status_changed, strlen(node_status), now - lasttime);
        if (node_status && strlen(node_status) > 20 &&
                (var->ds->online_status_changed || now > lasttime + 3600)) //在线状态变化上报，1小时也会上报
        {
            char topic[TOPIC_MAX_LEN];
            int ret = 0;

            snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/%s", var->sn_str, var->proc_name, TOPIC_EVT_NODESSTATUS);

            ret = dy_mqtt_session_publish(var->session, topic, node_status, strlen(node_status));
            if (ret == 0)
            {
                lasttime = now;
                var->ds->online_status_changed = false;
            }
            node_sta_backup(var->ds, S7_PLC_NODE_STATUS_BAK_FILE);
        }
        if (node_status)
        {
            write_file_data(NODES_CACHE "/s7_plc_status.json", node_status, strlen(node_status));
            free(node_status);
        }
    }
}

static int send_poll(s7_plc_var_t *var, struct list_head *data_list, int port, tcp_node_list_t *tcp_node)
{
    int ret, find = 0;
    unsigned int communication_timeout = 200;
    s7_plc_data_list_t *send_data = NULL;
    time_t now = time(NULL);

    list_for_each_entry(send_data, data_list, list)
    {
        if (!send_data->rs.period)
        {
            find = 1;
            if (communication_timeout < send_data->rs.communication_timeout)
            {
                communication_timeout = send_data->rs.communication_timeout;
            }
            s7_plc_bus_send(var, &send_data->rs, port, communication_timeout, tcp_node);

            list_del(&send_data->list);
            free(send_data->rs.data);
            free(send_data);
            break;
        }
        else
        {
            //第一次发送延时
            if (send_data->last_send < 5)
            {
                send_data->last_send++;
                find = 1;
                break;
            }
            if (var->period_flag == 0) continue;
            ret = check_timeout_second(now, &send_data->last_send, send_data->rs.period);
            if (ret)
            {
                find = 1;
                if (communication_timeout < send_data->rs.communication_timeout)
                {
                    communication_timeout = send_data->rs.communication_timeout;
                }
                s7_plc_bus_send(var, &send_data->rs, port, communication_timeout, tcp_node);
                break;
            }
        }
    }

    return find;
}

static int s7_plc_send_poll(s7_plc_var_t *var)
{
    TIMER_CONFIRM(var->send_timer);

    if (!var->start_flag)
    {
        return 0;
    }

    tcp_node_list_t *tcp_node = NULL;
    list_for_each_entry(tcp_node, &var->tcp_node_list, list)
    {
        if (send_poll(var, &tcp_node->data_list, S7_PLC, tcp_node) == 1)
            break;
    }

    return 0;
}

static void s7_plc_loop(s7_plc_var_t *var)
{
    int ret = -1, maxfd, i;
    fd_set rset;
    struct timeval timeout;

    while (1)
    {
        SELECT_INIT();
        SELECT_ADD_FD(var->send_timer);
        SELECT_ADD_FD(var->node_status_timer);

        timeout.tv_usec = 0;
        timeout.tv_sec = 5;

        ret = select(maxfd + 1, &rset, 0, 0, &timeout);
        if (ret < 0)
        {
            dy_syslog(LOG_INFO, "errno %d\n", errno);

            if (errno == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret > 0)
        {
            if (var->send_timer > 0 && FD_ISSET(var->send_timer, &rset))
            {
                FD_CLR(var->send_timer, &rset);
                s7_plc_send_poll(var);
            }
            if (var->node_status_timer > 0 && FD_ISSET(var->node_status_timer, &rset))
            {
                FD_CLR(var->node_status_timer, &rset);
                ds_flush_timer(var);
            }
        }
    }
}

static void s7_plc_subscribe_all(s7_plc_var_t *var)
{
    mqtt_session_t *mqtt_session = var->session;
    char topic[TOPIC_MAX_LEN] = {0};
    int i = 0;

    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/%s/device/+/data/%s", port_enum2char(S7_PLC), TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);
    dy_mqtt_session_subscribe(mqtt_session, topic);
    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/%s/%s", port_enum2char(S7_PLC), TOPIC_NOTIFY_UPGRADE);
    dy_mqtt_session_subscribe(mqtt_session, topic);
}

static int s7_plc_mqtt_handle_recv_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    s7_plc_var_t *var = (s7_plc_var_t *)obj;
    dy_syslog(LOG_DEBUG, "received MQTT topic:%s payload length:%d", mqtt_msg->topic, mqtt_msg->payloadLen);
    if (strstr(mqtt_msg->topic, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA))
    {
        //发送下行数据
        s7_plc_msg_tag_data_ctrl(var, mqtt_msg);
    }
}

// 建立与内部broker之间的MQTT连接
int s7_plc_mqtt_client_init(s7_plc_var_t *var)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_s7_plc_%s", var->sn_str);
    var->session = dy_mqtt_session_new(clientId);
    if (var->session == NULL) return -1;

    dy_mqtt_session_set_address(var->session, INTERNAL_BROKER_ADDR, INTERNAL_BROKER_PORT, INTERNAL_BROKER_USER, INTERNAL_BROKER_PASS);
    dy_mqtt_session_set_opts(var->session, DEFAULT_IPC_QOS, KEEP_ALIVE_MAX);
    dy_mqtt_session_set_callbacks(var->session, s7_plc_mqtt_handle_recv_msg, NULL);

    dy_mqtt_session_init(var->session, (void*)var);
    s7_plc_subscribe_all(var);
}

int s7_plc_period_rs_triger(s7_plc_var_t *var)
{
    int i, j, k;

    //建立TCP连接
    for (k = 0; k < var->nodes_cfg_table->node_cnt; k++)
    {
        if (var->nodes_cfg_table->node[k].port == S7_PLC && strlen(var->nodes_cfg_table->node[k].tcp_ip_addr) > 5)
        {
            int find = 0;
            tcp_node_list_t *tcp_node = NULL;

            list_for_each_entry(tcp_node, &var->tcp_node_list, list)
            {
                if (strcmp(tcp_node->tcp_ip_addr, var->nodes_cfg_table->node[k].tcp_ip_addr) == 0 &&
                        tcp_node->tcp_port == var->nodes_cfg_table->node[k].tcp_port)
                {
                    find = 1;
                    break;
                }
            }

            if (find == 0)
            {
                tcp_node = calloc(sizeof(tcp_node_list_t), 1);

                s7_plc_client_init(var, tcp_node, var->nodes_cfg_table->node[k].tcp_ip_addr, var->nodes_cfg_table->node[k].tcp_port);
                INIT_LIST_HEAD(&tcp_node->data_list);
                strncpy(tcp_node->tcp_ip_addr, var->nodes_cfg_table->node[k].tcp_ip_addr, SN_MAX_LEN);
                tcp_node->tcp_port = var->nodes_cfg_table->node[k].tcp_port;
                list_add_tail(&tcp_node->list, &var->tcp_node_list);
            }
        }
    }

    for (i = 0; i < var->template_table->template_cnt; i++)
    {
        //解析周期命令
        for (j = 0; j < var->template_table->template[i].service_tab->serviceCnt; j++)
        {
            if (var->template_table->template[i].service_tab->service[j].server_period && !var->template_table->template[i].service_tab->service[j].direction)
            {
                //节点信息
                for (k = 0; k < var->nodes_cfg_table->node_cnt; k++)
                {
                    if (var->nodes_cfg_table->node[k].port == S7_PLC && strcmp(var->nodes_cfg_table->node[k].template_id, var->template_table->template[i].template_id) == 0)
                    {
                        char *data = service_rs_data(&var->template_table->template[i].service_tab->service[j], var->nodes_cfg_table->node[k].sn);
                        char topic[TOPIC_MAX_LEN] = {0};
                        snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/%s", var->sn_str, port_enum2char(var->nodes_cfg_table->node[k].port), var->nodes_cfg_table->node[k].sn, TOPIC_EVT_SET_RGLT);
                        dy_mqtt_session_publish(var->session, topic, data, strlen(data));
                        free(data);
                    }
                }
            }
        }
    }
}

int s7_plc_init(s7_plc_var_t *var)
{
    int i, j;
    const char *board_name = NULL;

    check_make_dir(NODES_CACHE);
    check_make_dir(NODES_CFG);
    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    board_name = get_board_name();

    dy_syslog(LOG_INFO, "board name:%s", board_name);
    dy_syslog(LOG_INFO, "board SN:%s", var->sn_str);

    if (load_nodes_cfg(&var->nodes_cfg_table, NODES_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load nodes cfg fail");
    }
    if (load_templates_cfg(&var->template_table, TEMPLATES_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load template cfg fail");
    }
    dev_status_init(&var->ds, var->nodes_cfg_table, var->template_table, ports, ARRAY_SIZE(ports));
    node_sta_recovery(var->ds, S7_PLC_NODE_STATUS_BAK_FILE);

    kv_array_init(&var->identifier_backup, 32);

    INIT_LIST_HEAD(&var->tcp_node_list);

    var->send_timer = my_timer_create();
    if (var->send_timer > 0)
    {
        my_timer_set(var->send_timer, 1, 200);
    }
    var->node_status_timer = my_timer_create();
    if (var->node_status_timer > 0)
    {
        my_timer_set(var->node_status_timer, 1, 5000);
    }

    s7_plc_mqtt_client_init(var);
    s7_plc_period_rs_triger(var);

    var->start_flag = 1;
    var->period_flag = 1;

    return 0;
}
#endif

int main(int argc, char *argv[])
{
    s7_plc_var_t var = {0};
    
    openlog("s7_plc", LOG_PID, LOG_DAEMON);

    memset(&var, 0, sizeof(var));
	test_t();
    //s7_plc_init(&var);
    //s7_plc_loop(&var);

    return 0;
}
