Repository navigation
Expand file tree
/
Copy pathConnection.php
More file actions
147 lines (120 loc) · 4.58 KB
/
Copy pathConnection.php
File metadata and controls
147 lines (120 loc) · 4.58 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
142
143
144
145
146
147
<?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\DbConnection;
use Hyperf\Contract\ConnectionInterface;
use Hyperf\Contract\StdoutLoggerInterface;
use Hyperf\Database\ConnectionInterface as DbConnectionInterface;
use Hyperf\Database\Connectors\ConnectionFactory;
use Hyperf\DbConnection\Pool\DbPool;
use Hyperf\DbConnection\Traits\DbConnection;
use Hyperf\Pool\Connection as BaseConnection;
use Hyperf\Pool\Exception\ConnectionException;
use Psr\Container\ContainerInterface;
use Psr\EventDispatcher\EventDispatcherInterface;
use Psr\Log\LoggerInterface;
use Throwable;
class Connection extends BaseConnection implements ConnectionInterface, DbConnectionInterface
{
use DbConnection;
protected ?DbConnectionInterface $connection = null;
protected ConnectionFactory $factory;
protected LoggerInterface $logger;
public function __construct(ContainerInterface $container, DbPool $pool, protected array $config)
{
parent::__construct($container, $pool);
$this->factory = $container->get(ConnectionFactory::class);
$this->logger = $container->get(StdoutLoggerInterface::class);
$this->reconnect();
}
public function __call($name, $arguments)
{
return $this->connection->{$name}(...$arguments);
}
public function getActiveConnection(): DbConnectionInterface
{
if ($this->check()) {
return $this;
}
if (! $this->reconnect()) {
throw new ConnectionException('Connection reconnect failed.');
}
return $this;
}
public function reconnect(): bool
{
$this->close();
$this->connection = $this->factory->make($this->config);
if ($this->connection instanceof \Hyperf\Database\Connection) {
// Reset event dispatcher after db reconnect.
if ($this->container->has(EventDispatcherInterface::class)) {
$dispatcher = $this->container->get(EventDispatcherInterface::class);
$this->connection->setEventDispatcher($dispatcher);
}
// Reset reconnector after db reconnect.
$this->connection->setReconnector(function ($connection) {
$this->logger->warning('Database connection refreshing.');
if ($connection instanceof \Hyperf\Database\Connection) {
$this->refresh($connection);
}
});
}
$this->lastUseTime = microtime(true);
return true;
}
public function close(): bool
{
if ($this->connection instanceof \Hyperf\Database\Connection) {
$this->connection->disconnect();
}
unset($this->connection);
return true;
}
public function isTransaction(): bool
{
return $this->transactionLevel() > 0;
}
public function release(): void
{
try {
if ($this->connection instanceof \Hyperf\Database\Connection) {
// Reset $recordsModified property of connection to false before the connection release into the pool.
$this->connection->resetRecordsModified();
if ($this->connection->getErrorCount() > 100) {
// If the error count of connection is more than 100, we think it is a bad connection,
// So we'll reset it at the next time
$this->lastUseTime = 0.0;
}
}
if ($this->transactionLevel() > 0) {
$this->rollBack(0);
$this->logger->error('Maybe you\'ve forgotten to commit or rollback the MySQL transaction.');
}
} catch (Throwable $exception) {
$this->logger->error('Rollback connection failed, caused by ' . $exception);
// Ensure that the connection must be reset the next time after broken.
$this->lastUseTime = 0.0;
}
parent::release();
}
/**
* Refresh pdo and readPdo for current connection.
*/
protected function refresh(\Hyperf\Database\Connection $connection)
{
$refresh = $this->factory->make($this->config);
if ($refresh instanceof \Hyperf\Database\Connection) {
$connection->disconnect();
$connection->setPdo($refresh->getPdo());
$connection->setReadPdo($refresh->getReadPdo());
}
$this->logger->warning('Database connection refreshed.');
}
}