Compare commits

...

50 Commits

Author SHA1 Message Date
as2252258 2939d58585 1 2022-01-09 15:44:54 +08:00
as2252258 146da1d00c 1 2022-01-09 15:41:37 +08:00
as2252258 51babeddc1 1 2022-01-09 03:55:58 +08:00
as2252258 550197de96 e 2022-01-09 03:49:02 +08:00
as2252258 b5ad99332d 1 2022-01-09 03:46:41 +08:00
as2252258 e9b51ffb1e 1 2022-01-09 02:44:05 +08:00
as2252258 8c9de79027 Revert "改名"
This reverts commit fdf58326
2022-01-08 18:49:06 +08:00
as2252258 85c973469f Revert "改名"
This reverts commit fdf58326
2022-01-08 10:07:19 +08:00
as2252258 db8d9a8376 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:58:56 +08:00
as2252258 18b4e79f5f Revert "改名"
This reverts commit fdf58326
2022-01-07 18:58:16 +08:00
as2252258 0b7f380298 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:56:58 +08:00
as2252258 ea400c4ea5 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:54:46 +08:00
as2252258 07cfb62bb8 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:54:18 +08:00
as2252258 a5efa0813f Revert "改名"
This reverts commit fdf58326
2022-01-07 18:51:46 +08:00
as2252258 d1f2912324 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:51:12 +08:00
as2252258 ba0acf8f00 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:50:24 +08:00
as2252258 2788008fa1 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:45:44 +08:00
as2252258 5beee8352d Revert "改名"
This reverts commit fdf58326
2022-01-07 18:41:28 +08:00
as2252258 35859fead1 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:40:45 +08:00
as2252258 0031d73eff Revert "改名"
This reverts commit fdf58326
2022-01-07 18:40:05 +08:00
as2252258 89807fe548 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:38:45 +08:00
as2252258 9c28ce91ae Revert "改名"
This reverts commit fdf58326
2022-01-07 18:35:35 +08:00
as2252258 7ed01ba893 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:33:31 +08:00
as2252258 926a8c0719 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:29:08 +08:00
as2252258 688a30531a Revert "改名"
This reverts commit fdf58326
2022-01-07 18:27:01 +08:00
as2252258 1ff9d7f210 Revert "改名"
This reverts commit fdf58326
2022-01-07 18:16:29 +08:00
as2252258 da038f4fd8 Revert "改名"
This reverts commit fdf58326
2022-01-07 14:38:36 +08:00
as2252258 de20fa9582 Revert "改名"
This reverts commit fdf58326
2022-01-06 19:05:15 +08:00
as2252258 80cf72adc9 Revert "改名"
This reverts commit fdf58326
2022-01-04 16:06:29 +08:00
as2252258 dabcea452a 改名 2021-12-23 18:25:31 +08:00
as2252258 2feda943af 改名 2021-12-23 18:22:26 +08:00
as2252258 bffe375230 1 2021-12-17 04:20:20 +08:00
as2252258 efb69d02f3 改名 2021-12-08 11:46:08 +08:00
as2252258 79e5fa0b08 改名 2021-12-08 11:45:19 +08:00
as2252258 de7e005b46 改名 2021-12-08 11:32:31 +08:00
as2252258 a64fa8f493 改名 2021-12-08 11:16:35 +08:00
as2252258 c19e141f68 改名 2021-12-06 17:58:11 +08:00
as2252258 94a522588d 改名 2021-12-06 13:51:54 +08:00
as2252258 308c992fd3 改名 2021-12-06 11:47:12 +08:00
as2252258 458c6cbd22 改名 2021-12-03 16:05:49 +08:00
as2252258 b717a7e9f7 改名 2021-12-03 15:42:04 +08:00
as2252258 6786b313b5 改名 2021-12-03 14:46:38 +08:00
as2252258 14db451313 改名 2021-12-02 18:24:23 +08:00
as2252258 ced92a10e1 改名 2021-12-02 14:16:57 +08:00
as2252258 0d12baf49f 改名 2021-12-01 19:04:59 +08:00
as2252258 b798270bb7 改名 2021-12-01 17:21:15 +08:00
as2252258 955f358577 改名 2021-12-01 16:09:49 +08:00
as2252258 97a5deda23 改名 2021-12-01 15:28:59 +08:00
as2252258 8ae395e9cd 改名 2021-12-01 15:18:34 +08:00
as2252258 44a9a507f6 改名 2021-12-01 15:16:08 +08:00
45 changed files with 2114 additions and 2065 deletions
+34 -6
View File
@@ -3,7 +3,7 @@
namespace Server\Abstracts;
use JetBrains\PhpStorm\Pure;
use Kiri\Context;
use Server\Contract\OnProcessInterface;
use Swoole\Coroutine;
use Swoole\Process;
@@ -17,7 +17,7 @@ abstract class BaseProcess implements OnProcessInterface
protected bool $isStop = false;
protected mixed $redirect_stdin_and_stdout = null;
protected bool $redirect_stdin_and_stdout = FALSE;
protected int $pipe_type = SOCK_DGRAM;
@@ -47,10 +47,10 @@ abstract class BaseProcess implements OnProcessInterface
}
/**
* @return mixed
*/
public function getRedirectStdinAndStdout(): mixed
{
* @return bool
*/
public function getRedirectStdinAndStdout(): bool
{
return $this->redirect_stdin_and_stdout;
}
@@ -80,4 +80,32 @@ abstract class BaseProcess implements OnProcessInterface
}
/**
*
*/
public function onSigterm(): static
{
if (!Context::inCoroutine()) {
Process::signal(SIGTERM, fn($data) => $this->onShutdown($data));
} else {
Coroutine::create(function () {
$data = Coroutine::waitSignal(SIGTERM, -1);
if ($data) {
$this->onShutdown($data);
}
});
}
return $this;
}
/**
* @param $data
*/
protected function onShutdown($data): void
{
$this->isStop = true;
}
}
+1 -1
View File
@@ -4,7 +4,7 @@
namespace Server\Abstracts;
use Note\Inject;
use Kiri\Annotation\Inject;
use Exception;
use Kiri\Abstracts\Config;
use Kiri\Exception\ConfigException;
+1 -2
View File
@@ -13,10 +13,9 @@ interface OnCloseInterface
/**
* @param Server $server
* @param int $fd
*/
public function onClose(Server $server, int $fd): void;
public function onClose(int $fd): void;
}
+1 -2
View File
@@ -10,10 +10,9 @@ interface OnDisconnectInterface
/**
* @param Server $server
* @param int $fd
*/
public function onDisconnect(Server $server, int $fd): void;
public function onDisconnect(int $fd): void;
}
+2 -3
View File
@@ -2,7 +2,7 @@
namespace Server\Contract;
use Swoole\Server;
use Kiri\Websocket\WebSocketInterface;
use Swoole\WebSocket\Frame;
interface OnMessageInterface
@@ -10,10 +10,9 @@ interface OnMessageInterface
/**
* @param Server $server
* @param Frame $frame
* @return void
*/
public function onMessage(Server $server, Frame $frame): void;
public function onMessage(Frame $frame): void;
}
+4 -6
View File
@@ -3,16 +3,14 @@
namespace Server\Contract;
use Swoole\Http\Request;
use Swoole\WebSocket\Server;
interface OnOpenInterface
{
/**
* @param Server $server
* @param Request $request
*/
public function onOpen(Server $server, Request $request): void;
/**
* @param Request $request
*/
public function onOpen(Request $request): void;
}
+1 -1
View File
@@ -2,7 +2,7 @@
namespace Server\Handler;
use Note\Inject;
use Kiri\Annotation\Inject;
use Server\Abstracts\Server;
use Exception;
use Server\Contract\OnPipeMessageInterface;
+5 -1
View File
@@ -2,9 +2,10 @@
namespace Server\Handler;
use Note\Inject;
use Kiri\Annotation\Inject;
use Kiri\Events\EventDispatch;
use Kiri\Exception\ConfigException;
use ReflectionException;
use Server\Abstracts\Server;
use Server\Events\OnBeforeShutdown;
use Server\Events\OnShutdown;
@@ -28,6 +29,7 @@ class OnServer extends Server
/**
* @param \Swoole\Server $server
* @throws ConfigException
* @throws ReflectionException
*/
public function onStart(\Swoole\Server $server)
{
@@ -39,6 +41,7 @@ class OnServer extends Server
/**
* @param \Swoole\Server $server
* @throws ReflectionException
*/
public function onBeforeShutdown(\Swoole\Server $server)
{
@@ -48,6 +51,7 @@ class OnServer extends Server
/**
* @param \Swoole\Server $server
* @throws ReflectionException
*/
public function onShutdown(\Swoole\Server $server)
{
+5 -3
View File
@@ -2,8 +2,9 @@
namespace Server\Handler;
use Note\Inject;
use Kiri\Annotation\Inject;
use Kiri\Events\EventDispatch;
use ReflectionException;
use Server\Abstracts\Server;
use Kiri\Exception\ConfigException;
use Server\Events\OnManagerStart;
@@ -26,8 +27,8 @@ class OnServerManager extends Server
/**
* @param \Swoole\Server $server
* @throws ConfigException
*/
* @throws ConfigException|ReflectionException
*/
public function onManagerStart(\Swoole\Server $server)
{
$this->setProcessName(sprintf('manger[%d].0', $server->manager_pid));
@@ -38,6 +39,7 @@ class OnServerManager extends Server
/**
* @param \Swoole\Server $server
* @throws ReflectionException
*/
public function onManagerStop(\Swoole\Server $server)
{
+3 -1
View File
@@ -2,7 +2,7 @@
namespace Server\Handler;
use Note\Inject;
use Kiri\Annotation\Inject;
use Kiri\Events\EventDispatch;
use Server\Events\OnAfterReload;
use Server\Events\OnBeforeReload;
@@ -25,6 +25,7 @@ class OnServerReload
/**
* @param Server $server
* @throws \ReflectionException
*/
public function onBeforeReload(Server $server)
{
@@ -34,6 +35,7 @@ class OnServerReload
/**
* @param Server $server
* @throws \ReflectionException
*/
public function onAfterReload(Server $server)
{
+7 -7
View File
@@ -2,12 +2,12 @@
namespace Server\Handler;
use Note\Inject;
use Exception;
use Kiri\Abstracts\Config;
use Kiri\Core\Help;
use Kiri\Events\EventDispatch;
use Kiri\Kiri;
use Kiri\Annotation\Inject;
use Server\Events\OnAfterWorkerStart;
use Server\Events\OnBeforeWorkerStart;
use Server\Events\OnTaskerStart as OnTaskStart;
@@ -62,8 +62,8 @@ class OnServerWorker extends \Server\Abstracts\Server
*/
public function onWorkerStop(Server $server, int $workerId)
{
Timer::clearAll();
$this->eventDispatch->dispatch(new OnWorkerStop($server, $workerId));
Timer::clearAll();
$this->eventDispatch->dispatch(new OnWorkerStop($server, $workerId));
}
@@ -76,8 +76,8 @@ class OnServerWorker extends \Server\Abstracts\Server
{
set_env('state', 'exit');
$this->eventDispatch->dispatch(new OnWorkerExit($server, $workerId));
}
$this->eventDispatch->dispatch(new OnWorkerExit($server, $workerId));
}
/**
@@ -98,7 +98,7 @@ class OnServerWorker extends \Server\Abstracts\Server
$this->logger->error($message);
// $this->system_mail($message);
$this->system_mail($message);
}
@@ -109,7 +109,7 @@ class OnServerWorker extends \Server\Abstracts\Server
protected function system_mail($messageContent)
{
try {
$email = Config::get('email');
$email = Config::get('email', ['enable' => false]);
if (!empty($email) && ($email['enable'] ?? false) == true) {
Help::sendEmail($email, 'Service Error', $messageContent);
}
+38 -27
View File
@@ -3,13 +3,15 @@
namespace Server;
use Note\Inject;
use Exception;
use Http\Handler\Abstracts\HttpService;
use JetBrains\PhpStorm\Pure;
use Kiri\Abstracts\Config;
use Kiri\Events\EventDispatch;
use Kiri\Exception\ConfigException;
use Kiri\Kiri;
use Kiri\Annotation\Inject;
use Psr\Container\ContainerExceptionInterface;
use Psr\Container\NotFoundExceptionInterface;
use Server\Events\OnShutdown;
@@ -26,19 +28,14 @@ class Server extends HttpService
];
/**
* @Inject ServerManager
* @var null|ServerManager
*/
#[Inject(ServerManager::class)]
public ?ServerManager $manager = null;
private mixed $daemon = 0;
/** @var EventDispatch */
#[Inject(EventDispatch::class)]
public EventDispatch $eventDispatch;
/**
* @var State
*/
#[Inject(State::class)]
public State $state;
/**
@@ -59,54 +56,58 @@ class Server extends HttpService
/**
* @return string start server
*
* start server
* @return string
* @throws ConfigException
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws \ReflectionException
* @throws Exception
*/
public function start(): string
{
$this->manager->initBaseServer(Config::get('server', [], true), $this->daemon);
$this->manager()->initBaseServer(Config::get('server', [], true), $this->daemon);
$rpcService = Config::get('rpc', []);
if (!empty($rpcService)) {
$this->manager->addListener($rpcService['type'], $rpcService['host'], $rpcService['port'],
$this->manager()->addListener($rpcService['type'], $rpcService['host'], $rpcService['port'],
$rpcService['mode'], $rpcService);
}
$processes = array_merge($this->process, Config::get('processes', []));
foreach ($processes as $process) {
$this->manager->addProcess($process);
$this->manager()->addProcess($process);
}
return $this->manager->getServer()->start();
return $this->manager()->getServer()->start();
}
/**
* @return void
*
* start server
* @throws ConfigException
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws \ReflectionException
* @throws Exception
*/
public function shutdown()
{
$configs = Config::get('server', [], true);
foreach ($this->manager->sortService($configs['ports'] ?? []) as $config) {
$this->manager->stopServer($config['port']);
foreach ($this->manager()->sortService($configs['ports'] ?? []) as $config) {
$this->state->exit($config['port']);
}
$this->eventDispatch->dispatch(new OnShutdown());
$this->container->get(EventDispatch::class)->dispatch(new OnShutdown());
}
/**
* @return bool
* @throws ConfigException
* @throws Exception
*/
public function isRunner(): bool
{
return $this->manager->isRunner();
return $this->state->isRunner();
}
@@ -126,10 +127,20 @@ class Server extends HttpService
/**
* @return \Swoole\Http\Server|\Swoole\Server|\Swoole\WebSocket\Server|null
* @throws \ReflectionException
*/
#[Pure] public function getServer(): \Swoole\Http\Server|\Swoole\Server|\Swoole\WebSocket\Server|null
public function getServer(): \Swoole\Http\Server|\Swoole\Server|\Swoole\WebSocket\Server|null
{
return $this->manager->getServer();
return $this->manager()->getServer();
}
/**
* @return ServerManager
*/
private function manager(): ServerManager
{
return Kiri::getDi()->get(ServerManager::class);
}
}
+15 -11
View File
@@ -4,13 +4,14 @@ declare(strict_types=1);
namespace Server;
use Note\Inject;
use Exception;
use Kiri\Abstracts\Config;
use Kiri\Events\EventDispatch;
use Kiri\Exception\ConfigException;
use Kiri\Kiri;
use Server\Events\OnServerBeforeStart;
use Kiri\Annotation\Inject;
use Psr\Container\ContainerExceptionInterface;
use Psr\Container\NotFoundExceptionInterface;
use Swoole\Coroutine;
use Symfony\Component\Console\Command\Command;
use Symfony\Component\Console\Input\InputArgument;
@@ -52,6 +53,10 @@ class ServerCommand extends Command
* @param InputInterface $input
* @param OutputInterface $output
* @return int
* @throws ConfigException
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws \ReflectionException
* @throws Exception
*/
public function execute(InputInterface $input, OutputInterface $output): int
@@ -59,19 +64,19 @@ class ServerCommand extends Command
$manager = Kiri::app()->getServer();
$manager->setDaemon((int)!is_null($input->getOption('daemon')));
if (is_null($input->getArgument('action'))) {
$input->setArgument('action', 'restart');
throw new Exception('I don\'t know what I want to do.');
}
if (!in_array($input->getArgument('action'), self::ACTIONS)) {
throw new Exception('I don\'t know what I want to do.');
}
if ($manager->isRunner() && $input->getArgument('action') == 'start') {
throw new Exception('Service is running. Please use restart.');
if ($input->getArgument('action') == 'restart') {
$manager->shutdown();
}
if ($input->getArgument('action') != 'stop') {
return $this->generate_runtime_builder($manager);
}
$manager->shutdown();
if ($input->getArgument('action') == 'stop') {
return 0;
}
return $this->generate_runtime_builder($manager);
return 0;
}
@@ -81,6 +86,7 @@ class ServerCommand extends Command
private function configure_set()
{
$enable_coroutine = Config::get('servers.settings.enable_coroutine', false);
Config::set('servers.settings.enable_coroutine', true);
if ($enable_coroutine != true) {
return;
}
@@ -106,8 +112,6 @@ class ServerCommand extends Command
Kiri::app()->getRouter()->read_files();
$this->eventProvider->dispatch(new OnServerBeforeStart());
$manager->start();
return 1;
+63 -107
View File
@@ -2,14 +2,14 @@
namespace Server;
use Note\Inject;
use Exception;
use Kiri\Abstracts\Component;
use Kiri\Abstracts\Config;
use Kiri\Di\Container;
use Psr\Container\ContainerInterface;
use Kiri\Error\Logger;
use Kiri\Events\EventDispatch;
use Kiri\Exception\ConfigException;
use Kiri\Kiri;
use Kiri\Annotation\Inject;
use Psr\Container\ContainerExceptionInterface;
use Psr\Container\NotFoundExceptionInterface;
use ReflectionException;
@@ -22,25 +22,29 @@ use Server\Contract\OnMessageInterface;
use Server\Contract\OnPacketInterface;
use Server\Contract\OnProcessInterface;
use Server\Contract\OnReceiveInterface;
use Server\Contract\OnTaskInterface;
use Server\Events\OnServerBeforeStart;
use Server\Handler\OnPipeMessage;
use Server\Handler\OnServer;
use Server\Handler\OnServerManager;
use Server\Handler\OnServerReload;
use Server\Tasker\OnServerTask;
use Server\Handler\OnServerWorker;
use Server\Tasker\OnServerTask;
use Swoole\Http\Server as HServer;
use Swoole\Process;
use Swoole\Server;
use Swoole\Server\Port;
use Swoole\WebSocket\Server as WServer;
use Symfony\Component\Console\Helper\Table;
use Symfony\Component\Console\Helper\TableSeparator;
use Symfony\Component\Console\Output\ConsoleOutput;
use Symfony\Component\Console\Output\OutputInterface;
/**
* Class OnServerManager
* @package Http\Service
*/
class ServerManager
class ServerManager extends Component
{
@@ -52,10 +56,12 @@ class ServerManager
public int $port = 0;
/**
* @var Logger
*/
#[Inject(Logger::class)]
public Logger $logger;
/** @var array<string,Port> */
public array $ports = [];
@@ -65,11 +71,7 @@ class ServerManager
private Server|null $server = null;
/**
* @var Container
*/
#[Inject(ContainerInterface::class)]
public ContainerInterface $container;
protected array $initProcesses = [];
const DEFAULT_EVENT = [
@@ -100,9 +102,11 @@ class ServerManager
/**
* @return Server|WServer|HServer|null
* @throws ReflectionException
*/
public function getServer(): Server|WServer|HServer|null
{
di(EventDispatch::class)->dispatch(new OnServerBeforeStart());
return $this->server;
}
@@ -113,9 +117,9 @@ class ServerManager
* @param int $port
* @param int $mode
* @param array $settings
* @throws ReflectionException
* @throws ConfigException
* @throws Exception
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
*/
public function addListener(string $type, string $host, int $port, int $mode, array $settings = [])
{
@@ -131,8 +135,11 @@ class ServerManager
/**
* @throws ReflectionException
* @param $configs
* @param int $daemon
* @throws ConfigException
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
*/
public function initBaseServer($configs, int $daemon = 0): void
{
@@ -144,23 +151,6 @@ class ServerManager
}
/**
* @return bool
* @throws ConfigException
* @throws Exception
*/
public function isRunner(): bool
{
$configs = Config::get('server', [], true);
foreach ($this->sortService($configs['ports']) as $config) {
if (checkPortIsAlready($config['port'])) {
return true;
}
}
return false;
}
/**
* @param string|OnProcessInterface|BaseProcess $customProcess
* @throws Exception
@@ -171,15 +161,26 @@ class ServerManager
$customProcess = Kiri::getDi()->get($customProcess);
}
$system = sprintf('[%s].process', Config::get('id', 'system-service'));
$process = new Process(function (Process $process) use ($customProcess, $system) {
$this->logger->debug($system . ' ' . $customProcess->getName() . ' start.');
$this->server->addProcess(new Process(function (Process $process) use ($customProcess, $system) {
if (Kiri::getPlatform()->isLinux()) {
$process->name($system . '(' . $customProcess->getName() . ')');
}
$customProcess->process($process);
}, $customProcess->getRedirectStdinAndStdout(), $customProcess->getPipeType(), $customProcess->isEnableCoroutine());
$this->logger->debug($system . ' ' . $customProcess->getName() . ' start.');
$this->container->setBindings($customProcess->getName(), $process);
$this->server->addProcess($process);
$customProcess->onSigterm()->process($process);
},
$customProcess->getRedirectStdinAndStdout(),
$customProcess->getPipeType(),
$customProcess->isEnableCoroutine()
));
}
/**
* @return array<string,Process>
*/
public function getProcesses(): array
{
return $this->initProcesses;
}
@@ -192,22 +193,13 @@ class ServerManager
}
/**
* @param string $key
* @param string|int $value
*/
public static function setEnv(string $key, string|int $value): void
{
putenv(sprintf('%s=%s', $key, (string)$value));
}
/**
* @param ServerManager $context
* @param array $config
* @param int $daemon
* @throws ConfigException
* @throws ReflectionException
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws Exception
*/
private function startListenerHandler(ServerManager $context, array $config, int $daemon = 0)
@@ -256,7 +248,7 @@ class ServerManager
{
$id = Config::get('id', 'system-service');
$this->logger->debug(sprintf('[%s]' . $type . ' service %s::%d start', $id, $host, $port));
$this->logger->debug(sprintf('[%s].' . $type . ' service %s::%d start', $id, $host, $port));
/** @var Server\Port $service */
$this->ports[$port] = $this->server->addlistener($host, $port, $mode);
@@ -286,7 +278,6 @@ class ServerManager
* @throws ConfigException
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws ReflectionException
* @throws Exception
*/
private function createBaseServer(string $type, string $host, int $port, int $mode, array $settings = [])
@@ -300,35 +291,33 @@ class ServerManager
$this->server = new $match($host, $port, SWOOLE_PROCESS, $mode);
$this->server->set(array_merge(Config::get('server.settings', []), $settings['settings']));
$data = new Table($this->container->get(OutputInterface::class));
$data->setHeaders(['key', 'value']);
$array = [];
foreach ($this->server->setting as $key => $value) {
$array[] = [$key, $value];
$array[] = new TableSeparator();
}
array_pop($array);
$data->setStyle('box-double');
$data->setRows($array);
$data->render();
$id = Config::get('id', 'system-service');
$this->logger->debug(sprintf('[%s]' . $type . ' service %s::%d start', $id, $host, $port));
$this->logger->debug(sprintf('[%s].' . $type . ' service %s::%d start', $id, $host, $port));
$this->addDefaultListener($settings);
}
/**
* @param int $port
* @throws Exception
*/
public function stopServer(int $port)
{
if (!($pid = checkPortIsAlready($port))) {
return;
}
while (checkPortIsAlready($port)) {
Process::kill($pid, SIGTERM);
usleep(300);
}
}
/**
* @param array $settings
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws ReflectionException
*/
private function addDefaultListener(array $settings): void
{
@@ -348,7 +337,6 @@ class ServerManager
* @param Server|Port $server
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws ReflectionException
*/
private function addServiceEvents(array $events, Server|Port $server)
{
@@ -370,40 +358,6 @@ class ServerManager
}
/**
* @param string $class
* @return object
*/
private function getNewInstance(string $class): object
{
return $this->container->create($class);
}
/**
* @param OnTaskInterface|string $handler
* @param array $params
* @param int|null $workerId
* @throws ReflectionException
* @throws Exception
*/
public function task(OnTaskInterface|string $handler, array $params = [], int $workerId = null)
{
if ($workerId === null || $workerId <= $this->server->setting['worker_num']) {
$workerId = random_int($this->server->setting['worker_num'] + 1,
$this->server->setting['worker_num'] + 1 + $this->server->setting['task_worker_num']);
}
if (is_string($handler)) {
$implements = $this->container->getReflect($handler);
if (!in_array(OnTaskInterface::class, $implements->getInterfaceNames())) {
throw new Exception('Task must instance ' . OnTaskInterface::class);
}
$handler = $implements->newInstanceArgs($params);
}
$this->server->task(serialize($handler), $workerId);
}
/**
* @param mixed $message
* @param int $workerId
@@ -417,12 +371,14 @@ class ServerManager
/**
* @param array $events
* @throws ReflectionException
* @return void
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
*/
private function addTaskListener(array $events = []): void
{
$task_use_object = $this->server->setting['task_object'] ?? $this->server->setting['task_use_object'] ?? false;
$reflect = $this->container->getReflect(OnServerTask::class)?->newInstance();
$reflect = $this->container->get(OnServerTask::class);
$this->server->on('finish', $events[Constant::FINISH] ?? [$reflect, 'onFinish']);
if ($task_use_object || $this->server->setting['task_enable_coroutine']) {
$this->server->on('task', $events[Constant::TASK] ?? [$reflect, 'onCoroutineTask']);
@@ -434,9 +390,9 @@ class ServerManager
/**
* @param array|null $settings
* @return void
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws ReflectionException
*/
public function bindCallback(?array $settings = [])
{
+56
View File
@@ -0,0 +1,56 @@
<?php
namespace Server;
use Exception;
use Kiri\Abstracts\Component;
use Kiri\Abstracts\Config;
use Swoole\Process;
class State extends Component
{
use TraitServer;
public array $servers = [];
public function init()
{
$this->servers = Config::get('server.ports');
}
/**
* @return bool
* @throws Exception
*/
public function isRunner(): bool
{
$ports = $this->sortService($this->servers);
foreach ($ports as $config) {
if (checkPortIsAlready($config['port'])) {
return true;
}
}
return false;
}
/**
* @param $port
* @throws Exception
*/
public function exit($port)
{
if (!($pid = checkPortIsAlready($port))) {
return;
}
while (checkPortIsAlready($port)) {
Process::kill($pid, SIGTERM);
usleep(300);
}
}
}
+6 -14
View File
@@ -2,13 +2,10 @@
namespace Server\Tasker;
use Kiri\Kiri;
use Note\Inject;
use Exception;
use Kiri\Abstracts\BaseObject;
use Kiri\Abstracts\Component;
use Kiri\Core\HashMap;
use Kiri\Di\Container;
use Psr\Container\ContainerInterface;
use Kiri\Kiri;
use ReflectionException;
use Server\Contract\OnTaskInterface;
use Server\SwooleServerInterface;
@@ -17,7 +14,7 @@ use Server\SwooleServerInterface;
/**
*
*/
class AsyncTaskExecute extends BaseObject
class AsyncTaskExecute extends Component
{
@@ -27,13 +24,6 @@ class AsyncTaskExecute extends BaseObject
public ?SwooleServerInterface $server = null;
/**
* @var Container
*/
#[Inject(ContainerInterface::class)]
public ContainerInterface $container;
private HashMap $hashMap;
@@ -68,7 +58,9 @@ class AsyncTaskExecute extends BaseObject
$this->server = Kiri::getDi()->get(SwooleServerInterface::class);
}
if ($workerId === null || $workerId <= $this->server->setting['worker_num']) {
$workerId = random_int($this->server->setting['worker_num'] + 1, $this->server->setting['task_worker_num']);
$workerNum = $this->server->setting['worker_num'];
$taskerNum = $workerNum + $this->server->setting['task_worker_num'];
$workerId = random_int($workerNum, $taskerNum - 1);
}
if (is_string($handler)) {
$handler = $this->handle($handler, $params);
+1 -2
View File
@@ -4,10 +4,9 @@
namespace Server\Tasker;
use Note\Inject;
use Kiri\Annotation\Inject;
use Kiri\Abstracts\Logger;
use Kiri\Exception\ConfigException;
use ReflectionException;
use Server\Contract\OnTaskInterface;
use Swoole\Server;
+1 -1
View File
@@ -15,7 +15,7 @@
"psr/http-server-middleware": "^1.0",
"psr/http-message": "^1.0",
"psr/event-dispatcher": "^1.0",
"game-worker/kiri-http-message": "^v1.2"
"game-worker/kiri-http-message": "~v2.0"
},
"autoload": {
"psr-4": {