Repository navigation
Expand file tree
/
Copy pathSocketFactory.php
More file actions
122 lines (103 loc) · 3.54 KB
/
Copy pathSocketFactory.php
File metadata and controls
122 lines (103 loc) · 3.54 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
<?php
declare(strict_types=1);
/**
* This file is part of Hyperf.
*
* @link https://www.hyperf.io
* @document https://hyperf.wiki
* @contact [email protected]
* @license https://github.com/hyperf/hyperf/blob/master/LICENSE
*/
namespace Hyperf\RpcMultiplex;
use Hyperf\Collection\Arr;
use Hyperf\Contract\StdoutLoggerInterface;
use Hyperf\LoadBalancer\LoadBalancerInterface;
use Hyperf\LoadBalancer\Node;
use Hyperf\RpcMultiplex\Exception\NoAvailableNodesException;
use Psr\Container\ContainerInterface;
use function Hyperf\Support\make;
class SocketFactory
{
protected ?LoadBalancerInterface $loadBalancer = null;
/**
* @var Socket[]
*/
protected array $clients = [];
public function __construct(protected ContainerInterface $container, protected array $config)
{
}
public function getLoadBalancer(): ?LoadBalancerInterface
{
return $this->loadBalancer;
}
public function setLoadBalancer(LoadBalancerInterface $loadBalancer)
{
if ($loadBalancer->isAutoRefresh()) {
$this->bindAfterRefreshed($loadBalancer);
}
$this->loadBalancer = $loadBalancer;
}
public function refresh(): void
{
$nodes = $this->getNodes();
$nodeCount = count($nodes);
$count = $this->getCount();
for ($i = 0; $i < $count; ++$i) {
if (! isset($this->clients[$i])) {
$this->clients[$i] = make(Socket::class);
}
$client = $this->clients[$i];
$node = $nodes[$i % $nodeCount];
$client->setName($node->host)->setPort($node->port)->set([
'package_max_length' => $this->config['settings']['package_max_length'] ?? 1024 * 1024 * 2,
'recv_timeout' => $this->config['recv_timeout'] ?? 10,
'connect_timeout' => $this->config['connect_timeout'] ?? 0.5,
'heartbeat' => $this->config['heartbeat'] ?? null,
'max_requests' => $this->config['max_requests'] ?? 0,
'max_wait_close_seconds' => $this->config['max_wait_close_seconds'] ?? 0.5,
]);
if ($this->container->has(StdoutLoggerInterface::class)) {
$client->setLogger($this->container->get(StdoutLoggerInterface::class));
}
}
}
public function get(): Socket
{
if (count($this->clients) === 0) {
$this->refresh();
}
return Arr::random($this->clients);
}
protected function bindAfterRefreshed(LoadBalancerInterface $loadBalancer): void
{
$loadBalancer->afterRefreshed(static::class, function ($beforeNodes, $nodes) {
$items = [];
/** @var Node $node */
foreach ($beforeNodes as $node) {
$key = $node->host . $node->port . $node->weight . $node->pathPrefix;
$items[$key] = true;
}
foreach ($nodes as $node) {
$key = $node->host . $node->port . $node->weight . $node->pathPrefix;
if (array_key_exists($key, $items)) {
unset($items[$key]);
}
}
if (! empty($items)) {
$this->refresh();
}
});
}
protected function getNodes(): array
{
$nodes = $this->getLoadBalancer()->getNodes();
if (empty($nodes)) {
throw new NoAvailableNodesException();
}
return $nodes;
}
protected function getCount(): int
{
return (int) ($this->config['client_count'] ?? 4);
}
}