From 1653056386e1d8bd8ee147f99c7aafb8ec5592bb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9C=B1=E9=87=91=E5=8B=87?= Date: Mon, 20 Nov 2023 07:51:42 +0000 Subject: [PATCH] * Merge code for cron. --- framework/base/router.class.php | 9 +- module/cron/control.php | 160 ++++++++++++++++++++++++++------ module/cron/model.php | 89 ++++++++++++++++++ roadrunner/.rr.yaml | 2 + 4 files changed, 231 insertions(+), 29 deletions(-) diff --git a/framework/base/router.class.php b/framework/base/router.class.php index 88fb2f4e51..6049016360 100644 --- a/framework/base/router.class.php +++ b/framework/base/router.class.php @@ -2164,9 +2164,13 @@ class baseRouter } /* 合并Hook文件。Cycle the hook methods and merge hook codes. */ - $hookedMethods = array_keys($hookCodes); + $hookedMethods = array_keys($hookCodes); $mainTargetCodes = file($mainTargetFile); $mergedTargetCodes = file($tmpTargetFile); + + /* 如果已经合并,不需要再合并了。 If merged by other thread, return directly. */ + if(!$mergedTargetCodes) return; + foreach($hookedMethods as $method) { /* 通过反射获得hook脚本对应的方法所在的文件和起止行数。Reflection the hooked method to get it's defined position. */ @@ -3202,7 +3206,8 @@ class baseRouter if(!$this->config->debug) return true; if(!class_exists('dao')) return; - $sqlLog = $this->getLogRoot() . 'sql.' . date('Ymd') . '.log.php'; + $runMode = PHP_SAPI == 'cli' ? '_cli' : ''; + $sqlLog = $this->getLogRoot() . "sql$runMode." . date('Ymd') . '.log.php'; if(!is_file($sqlLog)) file_put_contents($sqlLog, "\n"); if(!is_writable($sqlLog)) return false; diff --git a/module/cron/control.php b/module/cron/control.php index f0b0924877..da2a63ed35 100644 --- a/module/cron/control.php +++ b/module/cron/control.php @@ -185,38 +185,95 @@ class cron extends control { if(empty($this->config->global->cron)) return; - if('cli' !== PHP_SAPI) + /* Run as daemon. */ + ignore_user_abort(true); + set_time_limit(0); + session_write_close(); + + $execId = mt_rand(); + if($restart) $this->cron->restartCron($execId); + + while(true) { - ignore_user_abort(true); - set_time_limit(0); - session_write_close(); + /* Only one scheduler and max 4 consumers. */ + $roles = $this->applyExecRoles($execId); + if(empty($roles)) + { + ignore_user_abort(false); + return; + } + + if(in_array('scheduler', $roles)) $this->schedule($execId); + if(in_array('consumer', $roles)) $this->consumeTasks($execId); + + sleep(20); } + } + + /** + * Schedule cron task by RoadRunner. + * + * @access public + * @return void + */ + public function rrSchedule() + { + if('cli' !== PHP_SAPI) return; + + set_time_limit(0); + session_write_close(); + + $execId = mt_rand(); + $this->cron->restartCron($execId); + + $this->loadModel('common'); + + while(true) + { + if(empty($this->config->global->cron)) + { + sleep(60); + continue; + } + + if($this->canSchedule($execId)) $this->schedule($execId); + + sleep(20); + } + } + + /** + * Consume cron task by RoadRunner. + * + * @access public + * @return void + */ + public function rrConsume() + { + if('cli' !== PHP_SAPI) return; + + set_time_limit(0); + session_write_close(); $this->loadModel('common'); $execId = mt_rand(); while(true) { - dao::$cache = array(); - if($restart || $this->canSchedule($execId)) + if(empty($this->config->global->cron)) { - $this->schedule($execId); - $this->execTasks($execId); + sleep(60); + continue; + } - $restart = false; - sleep(20); - } - else - { - $this->execTasks($execId); - return; // If no task in queue, executor will exit. - } + $this->consumeTasks($execId); + sleep(20); } } /** * 检查该execId是否可以进行调度(最近1分钟没有其他进程在调度). - * Check the execId can schedule(No other execId scheduled 10 minutes ago). + * Check the execId can schedule(No other execId scheduled 1 minutes ago). * * @param int $execId * @access protected @@ -224,13 +281,56 @@ class cron extends control */ protected function canSchedule($execId) { - $settings = $this->dao->select('`key`,`value`')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->andWhere('section')->eq('run')->fetchPairs(); + $settings = $this->dao->select('`key`,`value`')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->andWhere('section')->eq('scheduler')->fetchPairs(); if(!isset($settings['execId']) || $settings['execId'] == $execId) return true; if(!isset($settings['lastTime']) || $settings['lastTime'] < date('Y-m-d H:i:s', strtotime('-1 minute'))) return true; return false; } + /** + * 检查该execId是否可以执行任务(最多允许1个调度进程,4个执行进程). + * Check the execId can exec(1 scheduler, 4 consumers). + * + * @param int $execId + * @access protected + * @return array + */ + protected function applyExecRoles($execId) + { + $roles = array(); + + $settings = $this->dao->select('*')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->fetchAll(); + + $scheduler = array('execId' => 0, 'lastTime' => ''); + $consumerCount = 0; + + $expirDate = date('Y-m-d H:i:s', strtotime('-1 minute')); + foreach($settings as $setting) + { + if($setting->section == 'scheduler' && $setting->key == 'lastTime') $scheduler['lastTime'] = $setting->value; + if($setting->section == 'scheduler' && $setting->key == 'execId') $scheduler['execId'] = $setting->value; + + if($setting->section == 'consumer') + { + if($setting->value > $expirDate) + { + if($setting->key == strval($execId)) $roles[] = 'consumer'; + $consumerCount ++; + } + else + { + $this->dao->delete()->from(TABLE_CONFIG)->where('id')->eq($setting->id)->exec(); + } + } + } + + if(in_array($scheduler['execId'], array(0, $execId)) || $scheduler['lastTime'] < $expirDate) $roles[] = 'scheduler'; + if($consumerCount < 4) $roles[] = 'consumer'; + + return $roles; + } + /** * 调度生成队列任务. * Schedule, push tasks to queue. @@ -241,10 +341,11 @@ class cron extends control */ public function schedule($execId) { - $now = date(DT_DATETIME1); + $this->loadModel('common'); - $this->loadModel('setting')->setItem('system.cron.run.execId', $execId); - $this->setting->setItem('system.cron.run.lastTime', $now); + dao::$cache = array(); + + $this->cron->updateTime('scheduler', $execId); /* Get and parse crons. */ $tasks = $this->dao->select('cron,MAX(`createdDate`) `datetime`')->from(TABLE_QUEUE)->groupBy('cron')->fetchAll('cron'); @@ -255,6 +356,7 @@ class cron extends control } $parsedCrons = $this->cron->parseCron($crons); + $now = date(DT_DATETIME1); foreach($parsedCrons as $id => $cron) { @@ -289,17 +391,19 @@ class cron extends control * @access public * @return bool */ - public function execTasks($execId) + public function consumeTasks($execId) { while(true) { dao::$cache = array(); + + $this->cron->updateTime('consumer', $execId); + + /* Consume. */ $task = $this->dao->select('*')->from(TABLE_QUEUE)->where('status')->eq('wait')->andWhere('command')->ne('')->orderBy('createdDate')->fetch(); if(!$task) break; - $this->cron->logCron(strval($task->id) . "\n"); - - $this->execTask($execId, $task); + $this->consumeTask($execId, $task); } } @@ -312,7 +416,7 @@ class cron extends control * @access public * @return bool */ - public function execTask($execId, $task) + public function consumeTask($execId, $task) { /* Other executor may execute the task at the same time,so we mark execId and wait 500ms to check whether we own it. */ $this->dao->update(TABLE_QUEUE)->set('status')->eq('doing')->set('execId')->eq($execId)->where('id')->eq($task->id)->exec(); @@ -327,6 +431,8 @@ class cron extends control unset($_SESSION['company']); unset($this->app->company); + + $this->loadModel('common'); $this->common->setCompany(); $this->common->loadConfigFromDB(); @@ -359,7 +465,7 @@ class cron extends control $this->dao->update(TABLE_QUEUE)->set('status')->eq('done')->where('id')->eq($task->id)->exec(); $this->dao->update(TABLE_CRON)->set('lastTime')->eq(date(DT_DATETIME1))->where('id')->eq($task->cron)->exec(); - $log = date('G:i:s') . " execute\ncronId: {$task->cron}\nexecId: $execId\noutput: taskId:{$task->id}.\ncommand: {$task->command}.\nreturn : $return.\noutput : $output\n\n"; + $log = date('G:i:s') . " execute\ncronId: {$task->cron}\nexecId: $execId\ntaskId: {$task->id}\ncommand: {$task->command}\nreturn : $return\noutput : $output\n\n"; $this->cron->logCron($log); return true; diff --git a/module/cron/model.php b/module/cron/model.php index c9048d1094..272132256f 100644 --- a/module/cron/model.php +++ b/module/cron/model.php @@ -302,4 +302,93 @@ class cronModel extends model ->andWhere('`key`')->eq('cron') ->fetch('value'); } + + /** + * Restart cron. + * + * @access public + * @return void + */ + public function restartCron($execId) + { + $this->dao->update(TABLE_CONFIG)->set('value')->eq($execId) + ->where('owner')->eq('system') + ->andWhere('module')->eq('cron') + ->andWhere('section')->eq('scheduler') + ->andWhere('`key`')->eq($execId) + ->exec(); + $this->dao->delete()->from(TABLE_QUEUE)->where('createdDate')->le(date("Y-m-d H:i:s", strtotime("-1 week")))->exec(); + + $this->logCron(date('G:i:s') . " restart\n\n"); + } + + /** + * Update last time of cron. + * + * @param string $role + * @param int $execId + * @access public + * @return void + */ + public function updateTime($role, $execId) + { + $now = date(DT_DATETIME1); + + $settings = $this->dao->select('*')->from(TABLE_CONFIG)->where('owner')->eq('system')->andWhere('module')->eq('cron')->andWhere('section')->eq($role)->fetchAll('key'); + if($role == 'scheduler') + { + if(isset($settings['execId'])) + { + $setting = $settings['execId']; + if($setting->value != strval($execId)) $this->dao->update(TABLE_CONFIG)->set('value')->eq($execId)->where('id')->eq($setting->id)->exec(); + } + else + { + $data = new stdclass(); + $data->owner = 'system'; + $data->module = 'cron'; + $data->section = 'scheduler'; + $data->key = 'execId'; + $data->value = $execId; + + $this->dao->insert(TABLE_CONFIG)->data($data)->exec(); + } + + if(isset($settings['lastTime'])) + { + $setting = $settings['lastTime']; + $this->dao->update(TABLE_CONFIG)->set('value')->eq($now)->where('id')->eq($setting->id)->exec(); + } + else + { + $data = new stdclass(); + $data->owner = 'system'; + $data->module = 'cron'; + $data->section = 'scheduler'; + $data->key = 'lastTime'; + $data->value = $now; + + $this->dao->insert(TABLE_CONFIG)->data($data)->exec(); + } + } + else + { + if(isset($settings[strval($execId)])) + { + $setting = $settings[strval($execId)]; + $this->dao->update(TABLE_CONFIG)->set('value')->eq($now)->where('id')->eq($setting->id)->exec(); + } + else + { + $data = new stdclass(); + $data->owner = 'system'; + $data->module = 'cron'; + $data->section = 'consumer'; + $data->key = $execId; + $data->value = $now; + + $this->dao->insert(TABLE_CONFIG)->data($data)->exec(); + } + } + } } diff --git a/roadrunner/.rr.yaml b/roadrunner/.rr.yaml index cbe89792b3..e9fbe14f59 100644 --- a/roadrunner/.rr.yaml +++ b/roadrunner/.rr.yaml @@ -7,6 +7,7 @@ service: exec_timeout: 0 restart_sec: 1 relay: pipes + user: "nobody" cron_consumer: command: php consumer.php @@ -14,3 +15,4 @@ service: exec_timeout: 0 restart_sec: 1 relay: pipes + user: "nobody"