codispatcher.php 7.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229
  1. <?php
  2. declare(strict_types=0);
  3. define('APP_ID', 'cordispatcher');
  4. define('MOBILE_SERVER',true);
  5. define('USE_COROUTINE',true);
  6. define('SUPPORT_PTHREAD',false);
  7. define('BASE_ROOT_PATH',str_replace('/rdispatcher','',dirname(__FILE__)));
  8. define('BASE_PATH',BASE_ROOT_PATH . '/rdispatcher');
  9. require_once(BASE_ROOT_PATH . '/global.php');
  10. require_once(BASE_ROOT_PATH . '/fooder.php');
  11. require_once(BASE_HELPER_PATH . '/event_looper.php');
  12. require_once(BASE_HELPER_PATH . '/queue/rdispatcher.php');
  13. require_once(BASE_HELPER_PATH . '/algorithm.php');
  14. require_once(BASE_HELPER_PATH . '/refill/RefillFactory.php');
  15. require_once(BASE_PATH . '/processor.php');
  16. require_once(BASE_PATH . '/proxy.php');
  17. Co::set(['hook_flags' => SWOOLE_HOOK_NATIVE_CURL|SWOOLE_HOOK_SLEEP]);
  18. Co::set(['socket_connect_timeout' => 60,'socket_read_timeout' => 900,'socket_write_timeout' => 900,'socket_timeout' => 900]);
  19. if (empty($_SERVER['argv'][1])) exit('parameter error');
  20. $process_count = intval($_SERVER['argv'][1]);
  21. function all_channels() {
  22. return ['refill'];
  23. }
  24. function strbool($value) {
  25. return $value ? 'true' : 'false';
  26. }
  27. function handle_error($level, $message, $file, $line)
  28. {
  29. if($level == E_NOTICE) return;
  30. $trace = "handle_error: level={$level},msg={$message} file={$file},line={$line}\n";
  31. $backtrace = debug_backtrace();
  32. foreach ($backtrace as $item) {
  33. $trace .= "{$item['file']}\t{$item['line']}\t{$item['function']}\n";
  34. }
  35. Log::record($trace,Log::ERR);
  36. }
  37. function subscribe_message(&$quit, &$redis, $channels)
  38. {
  39. $redis = new Swoole\Coroutine\Redis();
  40. while (!$quit)
  41. {
  42. try
  43. {
  44. Log::record("subscribe_message start quit=" . strbool($quit),Log::DEBUG);
  45. if(!$redis->connected) {
  46. $ret = $redis->connect(C('coroutine.redis_host'), C('coroutine.redis_port'));
  47. }
  48. else {
  49. $ret = true;
  50. }
  51. if(!$ret) {
  52. Log::record("subscribe_message cannot connet redis.",Log::DEBUG);
  53. $redis->close();
  54. }
  55. elseif($redis->subscribe($channels))
  56. {
  57. while ($msg = $redis->recv())
  58. {
  59. [$sub_type, $channel, $content] = $msg;
  60. $content = json_decode($content,true);
  61. $type = $content['type'];
  62. if($channel != 'refill') continue;
  63. if($quit) break;
  64. if($type == 'channel' || $type == 'merchant') {
  65. refill\RefillFactory::instance()->load();
  66. refill\transfer::instance()->load();
  67. }
  68. elseif($type == 'ratio') {
  69. $ins = Cache::getInstance('cacheredis');
  70. $val = $ins->get_org('channel_ratios');
  71. if(empty($val)) continue;
  72. $val = json_decode($val,true);
  73. if(empty($val)) continue;
  74. $ratios = $val['ratios'];
  75. if(empty($ratios)) continue;
  76. refill\RefillFactory::instance()->UpdateRatio($ratios);
  77. }
  78. elseif($type == 'mch_counts') {
  79. $ins = Cache::getInstance('cacheredis');
  80. $content = $ins->get_org('merchant_refill_counts');
  81. if(empty($content)) continue;
  82. $counts = json_decode($content,true);
  83. if(empty($counts)) continue;
  84. $gross = $counts['gross'];
  85. $detail = $counts['detail'];
  86. refill\RefillFactory::instance()->UpdateMchRatios($gross,$detail);
  87. }
  88. elseif($type == 'channels_speed') {
  89. $ins = Cache::getInstance('cacheredis');
  90. $content = $ins->get_org('channels_speed');
  91. if(empty($content)) continue;
  92. $speeds = json_decode($content,true);
  93. if(empty($speeds)) continue;
  94. refill\RefillFactory::instance()->UpdateSpeeds($speeds);
  95. }
  96. else {
  97. //Log::record("subscribe_message dont not handle mgs:{$sub_type}-{$channel}-{$type}",Log::DEBUG);
  98. }
  99. }
  100. }
  101. else {
  102. Log::record("subscribe_message subscribe error",Log::ERR);
  103. }
  104. }
  105. catch (Exception $ex)
  106. {
  107. Log::record($ex->getMessage(),Log::ERR);
  108. }
  109. }
  110. Log::record("subscribe_message quit =" . strbool($quit),Log::DEBUG);
  111. $quit = true;
  112. }
  113. $workers = [];
  114. for ($i = 0; $i < $process_count;$i++)
  115. {
  116. $process = new Swoole\Process(function(Swoole\Process $worker)
  117. {
  118. Base::run_util();
  119. refill\RefillFactory::instance();
  120. refill\transfer::instance();
  121. $sub_quit = false;
  122. $sub_redis = null;
  123. $looper = new processor(false);
  124. set_error_handler('handle_error');
  125. register_shutdown_function(function () use ($looper)
  126. {
  127. $error = error_get_last();
  128. if(!empty($error) && $error['type'] != E_NOTICE) {
  129. $msg = "register_shutdown_function type:{$error['type']}\n message:{$error['message']}\n file={$error['file']} line={$error['line']}";
  130. Log::record("$msg",Log::ERR);
  131. }
  132. });
  133. go(function () use (&$sub_quit,&$sub_redis) {
  134. subscribe_message($sub_quit,$sub_redis,['refill']);
  135. });
  136. go(function () use ($looper) {
  137. $looper->run();
  138. });
  139. Swoole\Process::signal(SIGTERM, function($signal_num) use ($worker,&$sub_quit,$sub_redis,$looper)
  140. {
  141. Log::record("signal call SIGTERM begin signum={$signal_num}, pid={$worker->pid}",Log::DEBUG);
  142. set_error_handler(null);
  143. try {
  144. $sub_quit = true;
  145. $sub_redis->close();
  146. }
  147. catch(Exception $ex) {
  148. Log::record($ex->getMessage(),Log::DEBUG);
  149. }
  150. $looper->stop();
  151. });
  152. }, false, false, true);
  153. $pid = $process->start();
  154. $workers[$pid] = $process;
  155. }
  156. Log::record("main process start wait sub process....",Log::DEBUG);
  157. while (true)
  158. {
  159. if($status = Swoole\Process::wait(true))
  160. {
  161. $quit_pid = $status['pid'];
  162. Log::record("Sub process pid={$status['pid']} quit, code={$status['code']}, signum={$status['signal']}",Log::DEBUG);
  163. foreach ($workers as $pid => $worker)
  164. {
  165. if($pid != $quit_pid) {
  166. Swoole\Process::kill($pid, SIGTERM);
  167. }
  168. }
  169. foreach ($workers as $pid => $worker)
  170. {
  171. if($pid != $quit_pid && ($status = Swoole\Process::wait(true))) {
  172. Log::record("Graceful Recycled pid={$status['pid']}, code={$status['code']}, signum={$status['signal']}",Log::DEBUG);
  173. }
  174. }
  175. break;
  176. }
  177. else
  178. {
  179. foreach ($workers as $pid => $worker) {
  180. Swoole\Process::kill($pid, SIGTERM);
  181. }
  182. foreach ($workers as $pid => $worker)
  183. {
  184. if($status = Swoole\Process::wait(true)) {
  185. Log::record("Graceful Recycled pid={$status['pid']}, code={$status['code']}, signal={$status['signal']}",Log::DEBUG);
  186. }
  187. }
  188. break;
  189. }
  190. }
  191. Log::record("Quit all",Log::DEBUG);