name = "TCP"; $workerTcp::$logFile = __DIR__.'/workerman.log'; $num = 1; $workerTcp->onWorkerStart = function ($workerTcp) { global $db; //数据库 global $con; //访问本地创建的webserver服务器 global $connection_to_tcp; //访问接口服务器(远程的) global $connection_count; global $send_data; //存储insert数据 global $index; //insert管理数据序号 global $send_order; //insert全局的订单号 global $seqnum; global $rollback_data; //rollback数据 global $rollback_index; //rollback 管理数据序号 global $retry_data; //retry数据 global $retry_index; //retry 管理数据的序号 global $retry_rollback_data; //retry表中rollback数据 global $retry_rollback_index; //retry表rollback管理数据的序号 $seqnum = 1; $db = new Connection(DB_HOST, DB_PORT, DB_USERNAME, DB_PASSWORD, DB_NAME); //本地数据库的配置 //访问本地的websocket服务器============================================= $con = new AsyncTcpConnection("ws://127.0.0.1:12380"); $con->onConnect = function ($on) { global $index; global $rollback_index; global $retry_index; global $retry_rollback_index; $index = 0; $rollback_index = 0; $retry_index = 0; $retry_rollback_index = 0; sendToWebServer([ 'type' => 'login', 'uid' => 'tcp' ]); }; // websocke服务器发送信息过来触发的函数 $con->onMessage = function ($con, $data) { global $send_data; global $connection_count; $msg = json_decode($data, true); if ($msg['type'] == "heartbeat") { $connection_count = $msg["connection_count"]; } if ($msg['type'] == "get_order_record") { //查询订单状态 sendTo_tcp_Server($msg); } if ($msg['type'] == "set_current_order") { //设置当前跟单订单 sendTo_tcp_Server($msg); } if ($msg['type'] == "insert") { //Insert请求 global $index; //传递的index global $send_data; //全局的数据 global $send_order; //全局的订单号码 $insert_data = $msg['data']; $insert_data = array_chunk($insert_data, 12); //每十个是一个数组 $send_data = $insert_data; //传递的数据 $send_order = (int)$msg['orderid']; send_insert($send_data, $index, $send_order); } if ($msg['type'] == "rollback") { //Rollback请求 global $rollback_index; //传递index global $rollback_data; //全局的数据 $rollback_data = array_chunk($msg['data'], 5); //每几个是一个数组 send_rollback($rollback_data, $rollback_index); } if ($msg['type'] == "except_rollback") { //retry表中的 Rollback请求 global $retry_rollback_data; //retry表中rollback数据 global $retry_rollback_index; //retry表rollback管理数据的序号 $retry_rollback_data = array_chunk($msg['data'], 5); send_retry_rollback($retry_rollback_data, $retry_rollback_index); } if ($msg['type'] == "retry") { //Retry请求 global $retry_data; //retry数据 global $retry_index; //retry 管理数据的序号 $retry_data = array_chunk($msg['data'], 5); send_retry($retry_data, $retry_index); } }; $con->onClose = function ($con) { echo "connection closed\n"; $con->reConnect(1); }; $con->onError = function ($con, $code, $msg) { echo "Error code:$code msg:$msg\n"; }; $con->connect(); // 访问本地的websocket服务器======================================================= // 向远处的服务器发送信息=========================================================== $connection_to_tcp = new AsyncTcpConnection('tcp://119.23.51.113:10008'); //$connection_to_tcp = new AsyncTcpConnection('tcp://127.0.0.1:1235'); //$connection_to_tcp = new AsyncTcpConnection('tcp://192.168.5.111:10009'); $connection_to_tcp->onConnect = function ($connection_to_tcp) { // 发送心跳信息 $connection_to_tcp->send('{"type":"ping"}' . "\n"); }; // 接口服务器向我发送的信息 $connection_to_tcp->onMessage = function ($connection_to_tcp, $data) { tcpMessHandle($data); }; $connection_to_tcp->onClose = function($connection_to_tcp) { echo "connection closed\n"; log_file("to_tcp connection closed\n"); // $connection_to_tcp->reConnect(); //执行连接 }; $connection_to_tcp->onError = function($connection_to_tcp, $code, $msg) { echo "Error code:$code msg:$msg\n"; }; // 向远处的服务器发送信息============================================================ $connection_to_tcp->connect(); //执行连接 //发送给 接口tcp服务器 心跳======================================================== Timer::add(30, function () use ($connection_to_tcp) { //定时发送心跳信息 $connection_to_tcp->send('{"type":"ping"}' . "\n"); }); // 发送给 webserver 服务器 心跳===================================================== Timer::add(1, function () use ($con,$connection_to_tcp) { // 定时发送心跳信息 global $connection_count; if($connection_count){ if($connection_count>1){ $connection_to_tcp->connect(); //客户端上线重新连接接口 $connection_to_tcp->onConnect = function ($connection_to_tcp) { // 发送心跳信息 $connection_to_tcp->send('{"type":"ping"}' . "\n"); }; // 接口服务器向我发送的信息 $connection_to_tcp->onMessage = function ($connection_to_tcp, $data) { tcpMessHandle($data); }; $connection_to_tcp->onClose = function($connection_to_tcp) { echo "connection closed\n"; log_file("to_tcp connection closed\n"); // $connection_to_tcp->reConnect(); //执行连接 }; $connection_to_tcp->onError = function($connection_to_tcp, $code, $msg) { echo "Error code:$code msg:$msg\n"; }; }else{ //客户端不在线 $connection_to_tcp->close(); } } $con->send(json_encode(['type' => 'heartbeat']). "\n"); }); // 心跳检测(TCP 所有客户端) 心跳=================================================== Timer::add(1, function () use ($workerTcp) { $time_now = time(); //当前时间 foreach ($workerTcp->connections as $connection) { // 有可能该connection还没收到过消息,则lastMessageTime设置为当前时间 if (empty($connection->lastMessageTime)) { $connection->lastMessageTime = $time_now; continue; } // 上次通讯时间间隔大于心跳间隔,则认为客户端已经下线,关闭连接 if ($time_now - $connection->lastMessageTime > HEARTBEAT_TIME) { $connection->close(); } } }); }; $workerTcp->onConnect = function ($connection) { echo "connection onConnect\n"; }; $workerTcp->onError = function ($connection, $code, $msg) { echo "Error code:$code msg:$msg\n"; }; $workerTcp->onClose = function ($connection) { echo "connection closed\n"; }; $workerTcp->onMessage = function ($connection, $data) { echo "connection onMessage\n"; }; //远程服务器返回消息处理 function tcpMessHandle($data) { global $num; //这里显示的是从接口服务器返回的数据(需要根据数据的类似判断) global $db; //操作数据的全局变量 global $recv_buffer; //存储不完整的数据 global $index; //全局的变量 global $send_data; //全局的数据 global $send_order; //全局的订单号码 global $rollback_data; //rollback数据 global $rollback_index; //rollback 管理数据序号 global $retry_data; //retry数据 global $retry_index; //retry 管理数据的序号 global $retry_rollback_data; //retry表中rollback数据 global $retry_rollback_index; //retry表rollback管理数据的序号 global $seqnum; $bufExplode = explode("\n", $data); //通过指定的分隔符,把字符串打散为数组 foreach ($bufExplode as $key => $value) { if ($value) { $msg = json_decode($value, true); //对信息进行处理 if (!$msg) { $recv_buffer .= $value; //把数据组装起来 $msg1 = json_decode($recv_buffer, true); if ($msg1) { $msg = $msg1; } else { continue; //解析的不对就跳出循环 } } } else { continue; //信息不存在就跳出循环 } // 接口服务器返回的信息 echo "接口返回数据"; var_dump($msg); if ($msg['type'] == "get_order_record") { //查询订单状态 $recv_buffer = ""; /* if (!isset($msg['seqnum'])) { return false; } ack($msg['seqnum']); $seqnum++;*/ } if ($msg['type'] == "set_current_order") { //设置当前跟单订单 $recv_buffer = ""; /*if (!isset($msg['seqnum'])) { return false; }*/ sendToWebServer($msg); /* ack($msg['seqnum']); $seqnum++;*/ } if ($msg['type'] == "insert") { //Insert请求 (有数据库操作,插入) $recv_buffer = ""; /* if (!isset($msg['seqnum'])) { return false; } if ($seqnum != $msg["seqnum"]) { return false; }*/ $index++; //将数据加一 $total = count($send_data); //总数据的长度 $insert_data = $msg['dest']; /* ack($msg['seqnum']); $seqnum++;*/ if ($index < $total) { send_insert($send_data, $index, $send_order); } else { $index = 0; } sendToWebServer($msg); //发送给前端 $arr = ""; foreach ($insert_data as $key => $value) { if ($value['error_code'] == -2) { return false; } else { /* $msg_child['type'] = $msg['type']; $msg_child['orig_order'] = $msg['orig_order']; $msg_child['orig_login'] = $msg['orig_login']; $msg_child['dest_login'] = $value['dest_login']; $msg_child['dest_order'] = $value['dest_order']; $msg_child['percentage'] = $value['percentage']; $msg_child['profit'] = $value['profit']; $msg_child['error_code'] = $value['error_code']; $msg_child['addtime']=time(); //添加时间 $result = $db->insert('order_progress')->cols($msg_child)->query(); //向order_progress插入数据 $result1 = $db->insert('order_save')->cols($msg_child)->query(); //向order_save插入数据*/ $time = time(); $arr .= "(" . "'{$msg['type']}'" . "," . $msg['orig_order'] . "," . $msg['orig_login'] . "," . $value['dest_login'] . "," . $value['dest_order'] . "," . $value['percentage'] . "," . $value['profit'] . "," . $value['error_code'] . "," . $time . ")" . ","; } } $filed = "type,orig_order,orig_login,dest_login,dest_order,percentage,profit,error_code,addtime"; $arr = substr($arr, 0, -1); $sql = sprintf("INSERT INTO %s(%s) VALUES %s", "order_progress", $filed, $arr); $sql2 = sprintf("INSERT INTO %s(%s) VALUES %s", "order_save", $filed, $arr); $db->query($sql); $db->query($sql2); } if ($msg['type'] == "rollback") { //Rollback请求 (有数据库操作,插入并看看是否需要直接返回给前端) $recv_buffer = ""; /* if (!isset($msg['seqnum'])) { return false; } if ($seqnum != $msg['seqnum']) { return false; }*/ $rollback_index++; //将数据加一 $total = count($rollback_data); //总数据的长度 $rollbackdata = $msg['desc']; /* ack($msg["seqnum"]); $seqnum++;*/ if ($rollback_index < $total) { send_rollback($rollback_data, $rollback_index); }else{ $rollback_index = 0; } sendToWebServer($msg); //发送给前端 $delete_sql = 'delete from order_progress where '; $insert_sql = "insert into order_save(type,orig_order,orig_login,dest_login,dest_order,percentage,profit,error_code,addtime) values"; $update = []; foreach ($rollbackdata as $key => $value) { if ($value['error_code'] == -2) { return false; } else { $time = time(); if ($value['error_code'] == 0) { $delete_sql .= "(dest_login = $value[dest_login] and orig_order = $msg[orig_order]) or "; $insert_sql .= "('$msg[type]',$msg[orig_order],$msg[orig_login],$value[dest_login],$value[dest_order],$value[percentage],$value[profit],$value[error_code],$time),"; } else { //不成功把参数组装数组以便批量修改 $rollbackdata[$key]['type'] = $msg['type']; $rollbackdata[$key]['orig_order'] = $msg['orig_order']; $rollbackdata[$key]['orig_login'] = $msg['orig_login']; $rollbackdata[$key]['addtime'] = time(); array_push($update, $value); } } } //执行批量修改数据 if (!empty($update)) { $sql = batchUpdate($rollbackdata, 'dest_login', ['orig_order' => $msg['orig_order']],"order_progress"); $db->query($sql); } //成功执行删除后增加数据 if ($delete_sql != 'delete from order_progress where ') { $delete_sql = trim($delete_sql, 'or '); $insert_sql = substr($insert_sql, 0, -1); $db->query($delete_sql); $db->query($insert_sql); } } if ($msg['type'] == "except_rollback") { //retry表中的Rollback请求 (有数据库操作,插入并看看是否需要直接返回给前端) $recv_buffer = ""; if (!isset($msg['seqnum'])) { return false; } if ($seqnum != $msg['seqnum']) { return false; } /* ack($msg["seqnum"]); $seqnum++;*/ $retry_rollback_index++; //将数据加一 $total = count($retry_rollback_data); //总数据的长度 if ($retry_rollback_index < $total) { send_retry_rollback($retry_rollback_data, $retry_rollback_index); //继续发送数据 } else { $retry_rollback_index = 0; //对序号重置为0 } $exceptrollbackdata = $msg['desc']; sendToWebServer($msg); $update = []; $delete_sql = 'delete from order_progress where '; $insert_sql = "insert into order_save(type,orig_order,orig_login,dest_login,dest_order,percentage,profit,error_code,addtime) values"; foreach ($exceptrollbackdata as $key => $value) { if ($value['error_code'] == -2) { return false; } else { $time = time(); if ($value['error_code'] == 0) { $delete_sql .= "(dest_login = $value[dest_login] and orig_order = $msg[orig_order]) or "; $insert_sql .= "('$msg[type]',$msg[orig_order],$msg[orig_login],$value[dest_login],$value[dest_order],$value[percentage],$value[profit],$value[error_code],$time),"; } else { $exceptrollbackdata[$key]['type'] = $msg['type']; $exceptrollbackdata[$key]['orig_order'] = $msg['orig_order']; $exceptrollbackdata[$key]['orig_login'] = $msg['orig_login']; $exceptrollbackdata[$key]['addtime'] = time(); array_push($update, $value); } } } if (!empty($update)){ $sql = batchUpdate($exceptrollbackdata, 'dest_login', ['orig_order' => $msg['orig_order']],"order_progress"); $db->query($sql); } if ($delete_sql != 'delete from order_progress where ') { $delete_sql = trim($delete_sql, 'or '); $insert_sql = substr($insert_sql, 0, -1); $db->query($delete_sql); $db->query($insert_sql); } } if ($msg['type'] == "retry") { // if (!isset($msg['seqnum'])) { return false; } if ($seqnum != $msg['seqnum']) { return false; } /* ack($msg["seqnum"]); $seqnum++;*/ $retry_index++; //将数据加一 $total = count($retry_data); //总数据的长度 if ($retry_index < $total) { send_retry($retry_data, $retry_index); } else { $retry_index = 0; } $retrydata = $msg['desc']; sendToWebServer($msg); $succeUpdate = []; $update = []; $insert_sql = "insert into order_save(type,orig_order,orig_login,dest_login,dest_order,percentage,profit,error_code,addtime) values"; foreach ($retrydata as $key => $value) { if ($value['error_code'] == -2) { return false; } else { $time = time(); if ($value['error_code'] == 0) { $retrydata[$key]['type'] = $msg['type']; $retrydata[$key]['orig_order'] = $msg['orig_order']; $retrydata[$key]['orig_login'] = $msg['orig_login']; $retrydata[$key]['addtime'] = time(); array_push($succeUpdate, $value); $insert_sql .= "('$msg[type]',$msg[orig_order],$msg[orig_login],$value[dest_login],$value[dest_order],$value[percentage],$value[profit],$value[error_code],$time),"; } else { $retrydata[$key]['type'] = $msg['type']; $retrydata[$key]['orig_order'] = $msg['orig_order']; $retrydata[$key]['orig_login'] = $msg['orig_login']; $retrydata[$key]['addtime'] = time(); array_push($update, $value); } } } if (!empty($update)) { $sql = batchUpdate($retrydata, 'dest_login', ['orig_order' => $msg['orig_order']], "order_progress"); $db->query($sql); } if (!empty($succeUpdate)) { $sql = batchUpdate($retrydata, 'dest_login', ['orig_order' => $msg['orig_order']], "order_progress"); $db->query($sql); $insert_sql = substr($insert_sql, 0, -1); $db->query($insert_sql); } } } } //发送ack给远程服务器 function ack($seqnum){ global $connection_to_tcp; $data = ["type"=>"ack","seqnum"=>$seqnum]; $connection_to_tcp->send(json_encode($data) . "\n"); } // 发送信息到webserver服务器 function sendToWebServer($data) { global $con; log_file($data); $con->send(json_encode($data) . "\n"); } // 发送信息到接口服务器 function sendTo_tcp_Server($data){ global $connection_to_tcp; global $seqnum; $data['seqnum'] = $seqnum; log_file($data); $connection_to_tcp->send(json_encode($data) . "\n"); } // insert 数据 function send_insert($data,$index=0,$order){ $send['type'] = "insert_some"; $send['orig_order'] = (int)$order; //订单号码 $send['dest'] =[]; $send_data = $data[$index]; //将分割的数据交给全局变量 foreach ($send_data as $key => $value) { $send_child['dest_login'] = $value['LOGIN']; //登录的账号 $send_child['percentage'] = $value['ladder']; //梯度 array_push($send['dest'],$send_child); } echo "发送的数据"; var_dump($send); return false; sendTo_tcp_Server($send); } // rollback 数据 function send_rollback($data,$index=0){ $send['type'] = "rollback_some"; $send['orig_order'] = (int)$data[$index][0]['orig_order']; //订单号码 $send['orig_login'] = (int)$data[$index][0]['orig_login']; //登录者账号 $send['desc'] =[]; $send_child =[]; $send_data = $data[$index]; //将分割的数据交给全局变量 foreach ($send_data as $key => $value) { $send_child['dest_login'] = (int)$value['dest_login']; //登录的账号 $send_child['dest_order'] = (int)$value['dest_order']; //订单号 $send_child['percentage'] = (int)$value['percentage']; //手数 $send_child['profit'] = (float)$value['profit']; //利润点 $send_child['error_code'] = (int)$value['error_code']; //错误的号码 array_push($send['desc'],$send_child); } sendTo_tcp_Server($send); } // retry表中的rollback操作 function send_retry_rollback($data,$index=0){ $send['type'] = "except_rollback_some"; $send['orig_order'] = (int)$data[$index][0]['orig_order']; //订单号码 $send['orig_login'] = (int)$data[$index][0]['orig_login']; //登录者账号 $send['desc'] =[]; $send_child =[]; $send_data = $data[$index]; //将分割的数据交给全局变量 foreach ($send_data as $key => $value) { $send_child['dest_login'] = (int)$value['dest_login']; //登录的账号 $send_child['dest_order'] = (int)$value['dest_order']; //订单号 $send_child['percentage'] = (int)$value['percentage']; //手数 $send_child['profit'] = (float)$value['profit']; //利润点 $send_child['error_code'] = (int)$value['error_code']; //错误的号码 array_push($send['desc'],$send_child); } echo "发送的数据"; var_dump($send); sendTo_tcp_Server($send); } // retry表中的retry操作 function send_retry($data,$index=0){ $send['type'] = "retry_some"; $send['orig_order'] = (int)$data[$index][0]['orig_order']; //订单号码 $send['orig_login'] = (int)$data[$index][0]['orig_login']; //登录者账号 $send['desc'] =[]; $send_child =[]; $send_data = $data[$index]; //将分割的数据交给全局变量 foreach ($send_data as $key => $value) { $send_child['dest_login'] = (int)$value['dest_login']; //登录的账号 $send_child['dest_order'] = (int)$value['dest_order']; //订单号 $send_child['percentage'] = (int)$value['percentage']; //手数 $send_child['profit'] = (float)$value['profit']; //利润点 $send_child['error_code'] = (int)$value['error_code']; //错误的号码 array_push($send['desc'],$send_child); } echo "发送的数据"; var_dump($send); sendTo_tcp_Server($send); } //获取批量修改sql function batchUpdate($data,$field,$params=[],$tableName) { if (!is_array($data) || !$field || !is_array($params)) { return false; } $updates = parseUpdate($data, $field); $where = parseParams($params); $fields = array_column($data, $field); $fields = implode(',', array_map(function($value) { return "'".$value."'"; }, $fields)); $sql = sprintf("UPDATE `%s` SET %s WHERE `%s` IN (%s) %s", "$tableName", $updates, $field, $fields, $where); return $sql; } function parseUpdate($data,$field) { $sql = ""; $keys = array_keys(current($data)); foreach ($keys as $column){ $sql .= sprintf(" `%s` = CASE `%s` \n",$column,$field); foreach ($data as $line){ $sql .= sprintf("WHEN '%s' THEN '%s' \n",$line[$field],$line[$column]); } $sql .= "END,"; } return rtrim($sql,","); } function parseParams($params) { $where = []; foreach ($params as $key=>$value){ $where[] = sprintf(" `%s` = '%s' ",$key,$value); } return $where ? " AND " .implode('AND',$where) : ''; } //程序执行时间秒速 /*function curenttime() { $t = microtime(true); $micro = sprintf("%06d",($t - floor($t)) * 1000000); $d = new DateTime( date('Y-m-d H:i:s.'.$micro, $t) ); print $d->format("Y-m-d H:i:s.u"); // note at point on "u" }*/ //自定义日志类 function log_file($data){ $data = print_r($data,true); error_log(date("Y-m-d H:i:s", time()).' ===== info: '.$data."\n",3,__DIR__ . '/log/' . date("Ymd", time()) . '.log'); } Worker::runAll(); // 执行函数