stanley-king 3 年 前
コミット
852dccaec5
2 ファイル変更200 行追加0 行削除
  1. 15 0
      docker/compose/stanley/docker-compose.yml
  2. 185 0
      rdispatcher/coall.php

+ 15 - 0
docker/compose/stanley/docker-compose.yml

@@ -106,6 +106,21 @@ services:
     depends_on:
       - "redisrv"
 
+  coallsrv:
+    image: php-swool-redis:latest
+    volumes:
+      - ../../conf/etc/localtime:/etc/localtime:ro
+      - ../../../:/var/www/html
+      - /Volumes/Transcend/upload:/var/www/html/data/upload
+      - /Users/stanley-king/work/PHPProject/shoplog:/var/www/html/data/log
+      - ../../conf/php/php-swoole-debug.ini:/usr/local/etc/php/php.ini
+    links:
+      - redisrv
+    container_name: "panda-coall"
+    command: [php,"/var/www/html/rdispatcher/coall.php","1"]
+    depends_on:
+      - "redisrv"
+
   sendsrv:
     image: php-swool-redis:latest
     volumes:

+ 185 - 0
rdispatcher/coall.php

@@ -0,0 +1,185 @@
+<?php
+declare(strict_types=0);
+
+define('APP_ID', 'coall');
+define('MOBILE_SERVER',true);
+define('USE_COROUTINE',true);
+define('SUPPORT_PTHREAD',false);
+
+
+define('BASE_ROOT_PATH',str_replace('/rdispatcher','',dirname(__FILE__)));
+define('BASE_PATH',BASE_ROOT_PATH . '/rdispatcher');
+
+require_once(BASE_ROOT_PATH . '/global.php');
+require_once(BASE_ROOT_PATH . '/fooder.php');
+
+require_once(BASE_HELPER_PATH . '/event_looper.php');
+require_once(BASE_HELPER_PATH . '/queue/rdispatcher.php');
+require_once(BASE_HELPER_PATH . '/algorithm.php');
+require_once(BASE_HELPER_PATH . '/refill/RefillFactory.php');
+require_once(BASE_PATH . '/processor.php');
+require_once(BASE_PATH . '/proxy.php');
+
+Co::set(['hook_flags' => SWOOLE_HOOK_NATIVE_CURL | SWOOLE_HOOK_SLEEP | SWOOLE_HOOK_TCP]);
+if (empty($_SERVER['argv'][1])) exit('parameter error');
+$process_count = intval($_SERVER['argv'][1]);
+
+function all_channels() {
+    return ['refill'];
+}
+
+function strbool($value) {
+    return $value ? 'true' : 'false';
+}
+
+function handle_error($level, $message, $file, $line)
+{
+    if($level == E_NOTICE) return;
+    $trace = "handle_error: level={$level},msg={$message} file={$file},line={$line}\n";
+    $backtrace = debug_backtrace();
+    foreach ($backtrace as $item) {
+        $trace .= "{$item['file']}\t{$item['line']}\t{$item['function']}\n";
+    }
+
+    Log::record($trace,Log::ERR);
+}
+
+function subscribe_message(&$quit, &$redis, $channels)
+{
+    $redis = new Swoole\Coroutine\Redis();
+    while (!$quit)
+    {
+        try
+        {
+            Log::record("subscribe_message start quit=" . strbool($quit),Log::DEBUG);
+            if(!$redis->connected) {
+                $ret = $redis->connect(C('coroutine.redis_host'), C('coroutine.redis_port'));
+            }
+            else {
+                $ret = true;
+            }
+
+            if(!$ret) {
+                Log::record("subscribe_message cannot connet redis.",Log::DEBUG);
+                $redis->close();
+            }
+            elseif($redis->subscribe($channels))
+            {
+                while ($msg = $redis->recv())
+                {
+                    [$sub_type, $channel, $content] = $msg;
+                    $content = json_decode($content,true);
+                    $type = $content['type'];
+
+                    if($channel != 'refill') continue;
+                    if($quit) break;
+
+                    if($type == 'channel' || $type == 'merchant') {
+                        refill\RefillFactory::instance()->load();
+                    }
+                    elseif($type == 'ratio') {
+                        $start = microtime(true);
+                        $ins = Cache::getInstance('cacheredis');
+                        $val = $ins->get_org('channel_ratios');
+
+                        if(empty($val)) continue;
+                        $val = json_decode($val,true);
+                        if(empty($val)) continue;
+                        $ratios = $val['ratios'];
+                        if(empty($ratios)) continue;
+
+                        refill\RefillFactory::instance()->UpdateRatio($ratios);
+                        $use_time = microtime(true) - $start;
+                        $msg = sprintf("subscribe_message UpdateRatio use_time=%.6f",$use_time);
+                        Log::record($msg,Log::DEBUG);
+                    }
+                    else {
+                        Log::record("subscribe_message dont not handle mgs:{$sub_type}-{$channel}-{$type}",Log::DEBUG);
+                    }
+                }
+            }
+            else {
+                Log::record("subscribe_message subscribe error",Log::ERR);
+            }
+        }
+        catch (Exception $ex)
+        {
+            Log::record($ex->getMessage(),Log::ERR);
+        }
+    }
+
+    Log::record("subscribe_message quit =" . strbool($quit),Log::DEBUG);
+    $quit = true;
+}
+
+$workers = [];
+for ($i = 0; $i < $process_count;$i++)
+{
+    $process = new Swoole\Process(function(Swoole\Process $worker)
+    {
+        Base::run_util();
+        set_error_handler('handle_error');
+
+        $sub_quit = false;
+        $sub_redis = null;
+        go(function () use (&$sub_quit,&$sub_redis) {
+            subscribe_message($sub_quit,$sub_redis,['refill']);
+        });
+
+        $looper = new processor(false);
+        go(function () use ($looper) {
+            $looper->run();
+        });
+
+        Swoole\Process::signal(SIGTERM, function($signal_num) use ($worker,&$sub_quit,$sub_redis,$looper)
+        {
+            Log::record("signal call SIGTERM begin = $signal_num, #{$worker->pid}",Log::DEBUG);
+
+            set_error_handler(null);
+            try {
+                $sub_quit = true;
+                $sub_redis->close();
+            }
+            catch(Exception $ex) {
+                Log::record($ex->getMessage(),Log::DEBUG);
+            }
+
+            $looper->stop();
+            do {
+                $res = Swoole\Coroutine::stats();
+                $num = $res['coroutine_num'];
+                if($num > 1) {
+                    Swoole\Coroutine::sleep(0.1);
+                }
+                Log::record("coroutine_num = {$num}",Log::DEBUG);
+            } while($num > 1);
+        });
+
+    }, false, false, true);
+
+    $pid = $process->start();
+    $workers[$pid] = $process;
+}
+
+
+Log::record("main process start wait sub process....",Log::DEBUG);
+while (true)
+{
+    if($status = Swoole\Process::wait(true)) {
+        Log::record("Sub process #{$status['pid']} quit, code={$status['code']}, signal={$status['signal']}",Log::DEBUG);
+    }
+    else
+    {
+        foreach ($workers as $pid => $worker)
+        {
+            Swoole\Process::kill($pid,SIGTERM);
+            if($status = Swoole\Process::wait(true)) {
+                Log::record("Graceful Recycled #{$status['pid']}, code={$status['code']}, signal={$status['signal']}",Log::DEBUG);
+            }
+        }
+        break;
+    }
+}
+Log::record("Quit all",Log::DEBUG);
+
+