| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054 |
- <?php
- require dirname(dirname(__FILE__)).'/vendor/autoload.php';
- require_once dirname(dirname(__FILE__)) . '/vendor/workerman/Autoloader.php';
- require_once dirname(dirname(__FILE__)) . '/vendor/mysql/src/Connection.php';
- require_once __DIR__ . '/config.php';
- use Workerman\Worker; //worker容器类
- use Workerman\Connection\AsyncTcpConnection; //连接异步客户端
- use Workerman\MySQL\Connection; //连接数据库的类
- use Workerman\Lib\Timer; //定时器函数
- use Monolog\Logger;
- use Monolog\Handler\StreamHandler;
- use Monolog\Handler\ErrorLogHandler;
- date_default_timezone_set("PRC");
- //初始化 (创建了一个tcp服务)
- $workerTcp = new Worker("tcp://127.0.0.1:12345");
- $workerTcp->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://47.254.202.24:10009');
- //$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) {
- 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 = "";
- echo "接收到tcp" . curenttime();
- 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);
- }
- }
- }
- };
- $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(30, function () use ($con,$connection_to_tcp) {
- // 定时发送心跳信息
- global $connection_count;
- var_dump($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) {
- 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 = "";
- echo "接收到tcp" . curenttime();
- 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);
- }
- }
- }
- };
- $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 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);
- }
- 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(); // 执行函数
|