Repository navigation
Expand file tree
/
Copy pathRedis.php
More file actions
141 lines (125 loc) · 4.2 KB
/
Copy pathRedis.php
File metadata and controls
141 lines (125 loc) · 4.2 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
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
<?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\Redis;
use Hyperf\Context\Context;
use Hyperf\Redis\Exception\InvalidRedisConnectionException;
use Hyperf\Redis\Pool\PoolFactory;
use Throwable;
use function Hyperf\Coroutine\defer;
/**
* @mixin \Redis
*/
class Redis
{
use Traits\ScanCaller;
use Traits\MultiExec;
protected string $poolName = 'default';
public function __construct(protected PoolFactory $factory)
{
}
public function __call($name, $arguments)
{
// Get a connection from coroutine context or connection pool.
$hasContextConnection = Context::has($this->getContextKey());
$connection = $this->getConnection($hasContextConnection);
// Record the start time of the command.
$start = (float) microtime(true);
try {
/** @var RedisConnection $connection */
$connection = $connection->getConnection();
// Execute the command with the arguments.
$result = $connection->{$name}(...$arguments);
} catch (Throwable $exception) {
throw $exception;
} finally {
$time = round((microtime(true) - $start) * 1000, 2);
// Dispatch the command executed event.
$connection->getEventDispatcher()?->dispatch(
new Event\CommandExecuted(
$name,
$arguments,
$time,
$connection,
$this->poolName,
$result ?? null,
$exception ?? null,
)
);
// Release connection.
if (! $hasContextConnection) {
if ($this->shouldUseSameConnection($name)) {
if ($name === 'select' && $db = $arguments[0]) {
$connection->setDatabase((int) $db);
}
// Should storage the connection to coroutine context, then use defer() to release the connection.
Context::set($this->getContextKey(), $connection);
defer(function () {
$this->releaseContextConnection();
});
} else {
// Release the connection after command executed.
$connection->release();
}
}
}
return $result;
}
/**
* Release the connection stored in coroutine context.
*/
private function releaseContextConnection(): void
{
$contextKey = $this->getContextKey();
$connection = Context::get($contextKey);
if ($connection) {
Context::set($contextKey, null);
$connection->release();
}
}
/**
* Define the commands that need same connection to execute.
* When these commands executed, the connection will storage to coroutine context.
*/
private function shouldUseSameConnection(string $methodName): bool
{
return in_array($methodName, [
'multi',
'pipeline',
'select',
]);
}
/**
* Get a connection from coroutine context, or from redis connection pool.
* @param mixed $hasContextConnection
*/
private function getConnection($hasContextConnection): RedisConnection
{
$connection = null;
if ($hasContextConnection) {
$connection = Context::get($this->getContextKey());
}
if (! $connection instanceof RedisConnection) {
$pool = $this->factory->getPool($this->poolName);
$connection = $pool->get();
}
if (! $connection instanceof RedisConnection) {
throw new InvalidRedisConnectionException('The connection is not a valid RedisConnection.');
}
return $connection;
}
/**
* The key to identify the connection object in coroutine context.
*/
private function getContextKey(): string
{
return sprintf('redis.connection.%s', $this->poolName);
}
}