tcpServer.php 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464
  1. <?php
  2. require dirname(dirname(__FILE__)).'/vendor/autoload.php';
  3. require_once dirname(dirname(__FILE__)) . '/vendor/workerman/Autoloader.php';
  4. require_once dirname(dirname(__FILE__)) . '/vendor/mysql/src/Connection.php';
  5. require_once __DIR__ . '/config.php';
  6. use Workerman\Worker; //worker容器类
  7. use Workerman\Connection\AsyncTcpConnection; //连接异步客户端
  8. use Workerman\MySQL\Connection; //连接数据库的类
  9. use Workerman\Lib\Timer; //定时器函数
  10. use Monolog\Logger;
  11. use Monolog\Handler\StreamHandler;
  12. use Monolog\Handler\ErrorLogHandler;
  13. //初始化 (创建了一个tcp服务)
  14. $workerTcp = new Worker("tcp://127.0.0.1:12345");
  15. $workerTcp->name = "TCP";
  16. $workerTcp::$logFile = __DIR__.'/workerman.log';
  17. $workerTcp->onWorkerStart = function ($workerTcp) {
  18. global $db; //数据库
  19. global $con; //访问本地创建的webserver服务器
  20. global $connection_to_tcp; //访问接口服务器(远程的)
  21. global $send_data; //存储insert数据
  22. global $index; //insert管理数据序号
  23. global $send_order; //insert全局的订单号
  24. global $rollback_data; //rollback数据
  25. global $rollback_index; //rollback 管理数据序号
  26. $db = new Connection(DB_HOST, DB_PORT, DB_USERNAME, DB_PASSWORD, DB_NAME); //本地数据库的配置
  27. //访问本地的websocket服务器=============================================
  28. $con = new AsyncTcpConnection("ws://127.0.0.1:12380");
  29. $con->onConnect = function ($on) {
  30. global $index;
  31. global $rollback_index;
  32. $index=0;
  33. $rollback_index=0;
  34. sendToWebServer([
  35. 'type' => 'login',
  36. 'uid' => 'tcp'
  37. ]);
  38. };
  39. // websocke服务器发送信息过来触发的函数
  40. $con->onMessage = function ($con, $data) {
  41. global $send_data;
  42. $msg = json_decode($data, true);
  43. if ($msg['type'] == "get_order_record") { //查询订单状态
  44. sendTo_tcp_Server($msg);
  45. }
  46. if ($msg['type'] == "set_current_order") { //设置当前跟单订单
  47. sendTo_tcp_Server($msg);
  48. }
  49. if ($msg['type'] == "insert") { //Insert请求
  50. global $index; //传递的index
  51. global $send_data; //全局的数据
  52. global $send_order; //全局的订单号码
  53. $insert_data = $msg['data'];
  54. $insert_data = array_chunk($insert_data,12); //每十个是一个数组
  55. $send_data = $insert_data; //传递的数据
  56. $send_order = (int)$msg['orderid'];
  57. send_insert($send_data,$index,$send_order);
  58. }
  59. if ($msg['type'] == "rollback") { //Rollback请求
  60. global $rollback_index; //传递index
  61. global $rollback_data; //全局的数据
  62. $rollback_data = array_chunk($msg['data'],5); //每几个是一个数组
  63. send_rollback($rollback_data,$rollback_index);
  64. }
  65. if ($msg['type'] == "except_rollback") { //retry表中的 Rollback请求
  66. $orllback_data = $msg['data'];
  67. foreach ($orllback_data as $key => $value) {
  68. $send['type'] = $msg['type']; //rollback类型
  69. $send['orig_order'] = (int)$value['orig_order']; //手数
  70. $send['orig_login'] = (int)$value['orig_login']; //订单号码
  71. $send['dest_login'] = (int)$value['dest_login']; //登录的账号
  72. $send['dest_order'] = (int)$value['dest_order']; //登录的账号
  73. $send['error_code'] = (int)$value['error_code']; //登录的账号
  74. $send['percentage'] = (int)$value['percentage']; //手数
  75. $send['profit'] = (float)$value['profit']; //利润点
  76. sendTo_tcp_Server($send);
  77. }
  78. }
  79. if ($msg['type'] == "retry") { //Retry请求
  80. $orllback_data = $msg['data'];
  81. foreach ($orllback_data as $key => $value) {
  82. $send['type'] = $msg['type']; //rollback类型
  83. $send['orig_order'] = (int)$value['orig_order']; //手数
  84. $send['orig_login'] = (int)$value['orig_login']; //订单号码
  85. $send['dest_login'] = (int)$value['dest_login']; //登录的账号
  86. $send['dest_order'] = (int)$value['dest_order']; //登录的账号
  87. $send['error_code'] = (int)$value['error_code']; //登录的账号
  88. $send['percentage'] = (int)$value['percentage']; //手数
  89. $send['profit'] = (float)$value['profit']; //利润点
  90. sendTo_tcp_Server($send);
  91. }
  92. }
  93. };
  94. $con->onClose = function ($con) {
  95. echo "connection closed\n";
  96. $con->reConnect(1);
  97. };
  98. $con->onError = function($con, $code, $msg)
  99. {
  100. echo "Error code:$code msg:$msg\n";
  101. };
  102. $con->connect();
  103. // 访问本地的websocket服务器=======================================================
  104. // 向远处的服务器发送信息===========================================================
  105. $connection_to_tcp = new AsyncTcpConnection('tcp://47.254.202.24:10009');
  106. //$connection_to_tcp = new AsyncTcpConnection('tcp://127.0.0.1:1235');
  107. $connection_to_tcp->onConnect = function($connection_to_tcp)
  108. {
  109. // 发送心跳信息
  110. $connection_to_tcp->send('{"type":"ping"}' . "\r\n");
  111. };
  112. // 接口服务器向我发送的信息
  113. $connection_to_tcp->onMessage = function($connection_to_tcp, $data)
  114. {
  115. //这里显示的是从接口服务器返回的数据(需要根据数据的类似判断)
  116. global $db; //操作数据的全局变量
  117. global $index; //全局的变量
  118. global $send_data; //全局的数据
  119. global $send_order; //全局的订单号码
  120. global $recv_buffer; //存储不完整的数据
  121. global $rollback_data; //rollback数据
  122. global $rollback_index; //rollback 管理数据序号
  123. $bufExplode = explode("\r\n", $data); //通过指定的分隔符,把字符串打散为数组
  124. foreach ($bufExplode as $key => $value) {
  125. if($value){
  126. $msg = json_decode($value, true); //对信息进行处理
  127. if(!$msg){
  128. $recv_buffer.=$value; //把数据组装起来
  129. $msg1 = json_decode($recv_buffer, true);
  130. if($msg1){
  131. $msg = $msg1;
  132. }else{
  133. continue; //解析的不对就跳出循环
  134. }
  135. }
  136. }else{
  137. continue; //信息不存在就跳出循环
  138. }
  139. // 接口服务器返回的信息
  140. if ($msg['type'] == "get_order_record") { //查询订单状态
  141. $recv_buffer = "";
  142. sendToWebServer($msg);
  143. }
  144. if ($msg['type'] == "set_current_order") { //设置当前跟单订单
  145. $recv_buffer = "";
  146. sendToWebServer($msg);
  147. }
  148. if ($msg['type'] == "insert") { //Insert请求 (有数据库操作,插入)
  149. $recv_buffer = "";
  150. $index++; //将数据加一
  151. $total = count($send_data); //总数据的长度
  152. if($index < $total){
  153. send_insert($send_data,$index,$send_order);
  154. }else{
  155. $index = 0;
  156. }
  157. $insert_data = $msg['dest'];
  158. sendToWebServer($msg); //发送给前端
  159. // $insert_data_length = count($insert_data);
  160. // $insert_data_index =0;
  161. foreach ($insert_data as $key => $value) {
  162. if($value['error_code']==-2){
  163. return false;
  164. }else{
  165. $msg_child['type'] = $msg['type'];
  166. $msg_child['orig_order'] = $msg['orig_order'];
  167. $msg_child['orig_login'] = $msg['orig_login'];
  168. $msg_child['dest_login'] = $value['dest_login'];
  169. $msg_child['dest_order'] = $value['dest_order'];
  170. $msg_child['percentage'] = $value['percentage'];
  171. $msg_child['profit'] = $value['profit'];
  172. $msg_child['error_code'] = $value['error_code'];
  173. $msg_child['addtime']=time(); //添加时间
  174. $result = $db->insert('order_progress')->cols($msg_child)->query(); //插入数据
  175. //插入成功,向前台发起
  176. if($result){
  177. // $insert_data_index++;
  178. }
  179. }
  180. }
  181. }
  182. if ($msg['type'] == "rollback") { //Rollback请求 (有数据库操作,插入并看看是否需要直接返回给前端)
  183. $recv_buffer = "";
  184. $rollback_index++; //将数据加一
  185. $total = count($rollback_data); //总数据的长度
  186. if($rollback_index < $total){
  187. send_rollback($rollback_data,$rollback_index);
  188. }else{
  189. $rollback_index = 0;
  190. }
  191. $rollbackdata = $msg['desc'];
  192. sendToWebServer($msg); //发送给前端
  193. foreach ($rollbackdata as $key => $value) {
  194. if($value['error_code']==-2){
  195. return false;
  196. }else{
  197. $msg_child['type'] = $msg['type'];
  198. $msg_child['orig_order'] = $msg['orig_order'];
  199. $msg_child['orig_login'] = $msg['orig_login'];
  200. $msg_child['dest_login'] = $value['dest_login'];
  201. $msg_child['dest_order'] = $value['dest_order'];
  202. $msg_child['percentage'] = $value['percentage'];
  203. $msg_child['profit'] = $value['profit'];
  204. $msg_child['error_code'] = $value['error_code'];
  205. $msg_child['addtime']=time(); //添加时间
  206. $dest_login = $value['dest_login']; //获取更新的用户号码
  207. $orig_order = $msg['orig_order'];
  208. if($value['error_code']==0){
  209. $result = $db->delete('order_progress')->where('dest_login = :dest_login AND orig_order= :orig_order')->bindValues(
  210. array("dest_login" => $value['dest_login'],"orig_order" =>$msg['orig_order'])
  211. )->query();//成功rollback 对其中的数据进行删除
  212. }else{
  213. $result = $db->update('order_progress')->cols($msg_child)->where(["dest_login = $dest_login AND orig_order = $orig_order"])->query(); //不成功就更新数据
  214. }
  215. }
  216. }
  217. }
  218. if ($msg['type'] == "except_rollback") { //retry表中的Rollback请求 (有数据库操作,插入并看看是否需要直接返回给前端)
  219. $recv_buffer = "";
  220. if($msg['error_code']==-2){
  221. log_file($msg['error_code']);
  222. }else{
  223. $msg['addtime'] = time(); //更新时间
  224. $dest_login = $msg['dest_login']; //获取更新的用户号码
  225. $orig_order = $msg['orig_order'];
  226. if($msg['error_code']==0){
  227. $result = $db->delete('order_progress')->where('dest_login = :dest_login AND orig_order= :orig_order')->bindValues(
  228. array("dest_login" => $value['dest_login'],"orig_order" =>$msg['orig_order'])
  229. )->query();//成功rollback 对其中的数据进行删除
  230. if($result){
  231. sendToWebServer($msg);
  232. }
  233. }else{
  234. $result = $db->update('order_progress')->cols($msg)->where(["dest_login = $dest_login AND orig_order = $orig_order"])->query(); //不成功就更新数据
  235. if($result){
  236. sendToWebServer($msg);
  237. }
  238. }
  239. }
  240. }
  241. if ($msg['type'] == "retry") { //Retry请求 (有数据库操作 0成功然后更新数据库 )
  242. $recv_buffer = "";
  243. if($msg['error_code']==-2){
  244. return false;
  245. }else{
  246. $msg['addtime']=time(); //更新时间
  247. $dest_login = $msg['dest_login']; //获取更新的用户号码
  248. $orig_order = $msg['orig_order'];
  249. $result = $db->update('order_progress')->cols($msg)->where(["dest_login = $dest_login AND orig_order = $orig_order"])->query(); //不成功就更新数据
  250. // 更新成功,发送请求
  251. if($result){
  252. sendToWebServer($msg);
  253. }
  254. }
  255. }
  256. if($msg['type']=="pong"){
  257. $recv_buffer = "";
  258. echo "12345\r\n";
  259. }
  260. }
  261. };
  262. $connection_to_tcp->onClose = function($connection_to_tcp)
  263. {
  264. echo "connection closed\n";
  265. log_file("to_tcp connection closed\n");
  266. $connection_to_tcp->reConnect(1); //重新连接
  267. };
  268. $connection_to_tcp->onError = function($connection_to_tcp, $code, $msg)
  269. {
  270. echo "Error code:$code msg:$msg\n";
  271. };
  272. $connection_to_tcp->connect(); //执行连接
  273. // 向远处的服务器发送信息============================================================
  274. // 发送给 webserver 服务器 心跳=====================================================
  275. Timer::add(30, function () use ($con) {
  276. // 定时发送心跳信息
  277. $con->send(json_encode([
  278. 'type' => 'heartbeat'
  279. ]));
  280. });
  281. //发送给 接口tcp服务器 心跳========================================================
  282. Timer::add(30, function () use ($connection_to_tcp) {
  283. //定时发送心跳信息
  284. $connection_to_tcp->send('{"type":"ping"}' . "\r\n");
  285. });
  286. // 心跳检测(TCP 所有客户端) 心跳===================================================
  287. Timer::add(1, function () use ($workerTcp) {
  288. $time_now = time(); //当前时间
  289. foreach ($workerTcp->connections as $connection) {
  290. // 有可能该connection还没收到过消息,则lastMessageTime设置为当前时间
  291. if (empty($connection->lastMessageTime)) {
  292. $connection->lastMessageTime = $time_now;
  293. continue;
  294. }
  295. // 上次通讯时间间隔大于心跳间隔,则认为客户端已经下线,关闭连接
  296. if ($time_now - $connection->lastMessageTime > HEARTBEAT_TIME) {
  297. $connection->close();
  298. }
  299. }
  300. });
  301. };
  302. $workerTcp->onConnect = function ($connection) {
  303. echo "connection onConnect\n";
  304. };
  305. $workerTcp->onError = function ($connection, $code, $msg) {
  306. echo "Error code:$code msg:$msg\n";
  307. };
  308. $workerTcp->onClose = function ($connection) {
  309. echo "connection closed\n";
  310. };
  311. $workerTcp->onMessage = function ($connection, $data) {
  312. echo "connection onMessage\n";
  313. };
  314. // 发送信息到webserver服务器
  315. function sendToWebServer($data)
  316. {
  317. global $con;
  318. $con->send(json_encode($data) . "\r\n");
  319. }
  320. // 发送信息到接口服务器
  321. function sendTo_tcp_Server($data){
  322. global $connection_to_tcp;
  323. $connection_to_tcp->send(json_encode($data) . "\r\n");
  324. }
  325. // insert 数据
  326. function send_insert($data,$index=0,$order){
  327. $send['type'] = "insert_some";
  328. $send['orig_order'] = (int)$order; //订单号码
  329. $send['dest'] =[];
  330. $send_data = $data[$index]; //将分割的数据交给全局变量
  331. foreach ($send_data as $key => $value) {
  332. $send_child['dest_login'] = $value['LOGIN']; //登录的账号
  333. $send_child['percentage'] = $value['ladder']; //梯度
  334. array_push($send['dest'],$send_child);
  335. }
  336. echo "发送的数据";
  337. var_dump($send);
  338. sendTo_tcp_Server($send);
  339. }
  340. // rollback 数据
  341. function send_rollback($data,$index=0){
  342. $send['type'] = "rollback_some";
  343. $send['orig_order'] = (int)$data[$index][0]['orig_order']; //订单号码
  344. $send['orig_login'] = (int)$data[$index][0]['orig_login']; //登录者账号
  345. $send['desc'] =[];
  346. $send_child =[];
  347. $send_data = $data[$index]; //将分割的数据交给全局变量
  348. foreach ($send_data as $key => $value) {
  349. $send_child['dest_login'] = (int)$value['dest_login']; //登录的账号
  350. $send_child['dest_order'] = (int)$value['dest_order']; //订单号
  351. $send_child['percentage'] = (int)$value['percentage']; //手数
  352. $send_child['profit'] = (float)$value['profit']; //利润点
  353. $send_child['error_code'] = (int)$value['error_code']; //错误的号码
  354. array_push($send['desc'],$send_child);
  355. }
  356. echo "发送的数据";
  357. var_dump($send);
  358. sendTo_tcp_Server($send);
  359. }
  360. //自定义日志类
  361. function log_file($data){
  362. error_log(date("Y-m-d H:i:s", time()).' ===== info: '.$data."\r\n",3,__DIR__ . '/log/' . date("Ymd", time()) . '.log');
  363. }
  364. Worker::runAll(); // 执行函数