| 1 |
<?php |
| 2 |
/* |
| 3 |
* This file is part of the ManageWP Worker plugin. |
| 4 |
* |
| 5 |
* (c) ManageWP LLC <[email protected]> |
| 6 |
* |
| 7 |
* For the full copyright and license information, please view the LICENSE |
| 8 |
* file that was distributed with this source code. |
| 9 |
*/ |
| 10 |
|
| 11 |
class MWP_IncrementalBackup_Database_StreamableQuerySequenceDump |
| 12 |
{ |
| 13 |
|
| 14 |
/** |
| 15 |
* @var MWP_IncrementalBackup_Database_ConnectionInterface |
| 16 |
*/ |
| 17 |
private $connection; |
| 18 |
|
| 19 |
/** |
| 20 |
* @var MWP_IncrementalBackup_Database_DumpOptions |
| 21 |
*/ |
| 22 |
private $options; |
| 23 |
|
| 24 |
public function __construct(MWP_IncrementalBackup_Database_ConnectionInterface $connection, MWP_IncrementalBackup_Database_DumpOptions $options) |
| 25 |
{ |
| 26 |
$this->connection = $connection; |
| 27 |
$this->options = $options; |
| 28 |
} |
| 29 |
|
| 30 |
/** |
| 31 |
* @return MWP_IncrementalBackup_Database_ConnectionInterface |
| 32 |
*/ |
| 33 |
protected function getConnection() |
| 34 |
{ |
| 35 |
return $this->connection; |
| 36 |
} |
| 37 |
|
| 38 |
/** |
| 39 |
* @inherit |
| 40 |
*/ |
| 41 |
public function createStream() |
| 42 |
{ |
| 43 |
$stream = new MWP_Stream_Append(); |
| 44 |
|
| 45 |
$stream->addStream(MWP_Stream_Stream::factory(" |
| 46 |
/*!40101 SET @OLD_CHARACTER_SET_CLIENT=@@CHARACTER_SET_CLIENT */; |
| 47 |
/*!40101 SET @OLD_CHARACTER_SET_RESULTS=@@CHARACTER_SET_RESULTS */; |
| 48 |
/*!40101 SET @OLD_COLLATION_CONNECTION=@@COLLATION_CONNECTION */; |
| 49 |
/*!40101 SET NAMES utf8 */; |
| 50 |
/*!40103 SET @OLD_TIME_ZONE=@@TIME_ZONE */; |
| 51 |
/*!40103 SET TIME_ZONE='+00:00' */; |
| 52 |
/*!40014 SET @OLD_UNIQUE_CHECKS=@@UNIQUE_CHECKS, UNIQUE_CHECKS=0 */; |
| 53 |
/*!40014 SET @OLD_FOREIGN_KEY_CHECKS=@@FOREIGN_KEY_CHECKS, FOREIGN_KEY_CHECKS=0 */; |
| 54 |
/*!40101 SET @OLD_SQL_MODE=@@SQL_MODE, SQL_MODE='NO_AUTO_VALUE_ON_ZERO' */; |
| 55 |
/*!40111 SET @OLD_SQL_NOTES=@@SQL_NOTES, SQL_NOTES=0 */;\n\n" |
| 56 |
)); |
| 57 |
|
| 58 |
$allTables = self::arrayColumn($this->getConnection()->query('SHOW TABLES')->fetchAll()); |
| 59 |
$tables = array_intersect($allTables, $this->options->getTables() ? $this->options->getTables() : $allTables); |
| 60 |
|
| 61 |
foreach ($tables as $tableName) { |
| 62 |
$stream->addStream( |
| 63 |
new MWP_Stream_Callable(array($this, 'streamCreateTable'), array($tableName)) |
| 64 |
); |
| 65 |
} |
| 66 |
|
| 67 |
$stream->addStream(MWP_Stream_Stream::factory(" |
| 68 |
/*!40103 SET TIME_ZONE=@OLD_TIME_ZONE */; |
| 69 |
/*!40101 SET SQL_MODE=@OLD_SQL_MODE */; |
| 70 |
/*!40014 SET FOREIGN_KEY_CHECKS=@OLD_FOREIGN_KEY_CHECKS */; |
| 71 |
/*!40014 SET UNIQUE_CHECKS=@OLD_UNIQUE_CHECKS */; |
| 72 |
/*!40101 SET CHARACTER_SET_CLIENT=@OLD_CHARACTER_SET_CLIENT */; |
| 73 |
/*!40101 SET CHARACTER_SET_RESULTS=@OLD_CHARACTER_SET_RESULTS */; |
| 74 |
/*!40101 SET COLLATION_CONNECTION=@OLD_COLLATION_CONNECTION */; |
| 75 |
/*!40111 SET SQL_NOTES=@OLD_SQL_NOTES */;\n" |
| 76 |
)); |
| 77 |
|
| 78 |
return $stream; |
| 79 |
} |
| 80 |
|
| 81 |
public function streamCreateTable($length, $tableName) |
| 82 |
{ |
| 83 |
// Get the SHOW CREATE TABLE part |
| 84 |
$content = $this->getConnection() |
| 85 |
->query("SHOW CREATE TABLE `{$tableName}`;") |
| 86 |
->fetchAll(); |
| 87 |
|
| 88 |
if (!is_array($content)) { |
| 89 |
return new MWP_Stream_Buffer(); |
| 90 |
} |
| 91 |
|
| 92 |
$stream = new MWP_Stream_Append(); |
| 93 |
|
| 94 |
foreach ($content as $entry) { |
| 95 |
// Add drop table query |
| 96 |
if ($this->options->isDropTables()) { |
| 97 |
$stream->addStream(MWP_Stream_Stream::factory("DROP TABLE IF EXISTS `$tableName`;\n")); |
| 98 |
} |
| 99 |
|
| 100 |
// Add create table query |
| 101 |
$stream->addStream(MWP_Stream_Stream::factory(" |
| 102 |
/*!40101 SET @saved_cs_client = @@character_set_client */; |
| 103 |
/*!40101 SET character_set_client = utf8 */;\n" |
| 104 |
)); |
| 105 |
$stream->addStream(MWP_Stream_Stream::factory($entry['Create Table'].";\n")); |
| 106 |
$stream->addStream(MWP_Stream_Stream::factory("/*!40101 SET character_set_client = @saved_cs_client */;\n\n")); |
| 107 |
} |
| 108 |
|
| 109 |
// Export content |
| 110 |
$stream->addStream( |
| 111 |
new MWP_Stream_Callable(array($this, 'createExportTableStream'), array($tableName)) |
| 112 |
); |
| 113 |
|
| 114 |
return $stream; |
| 115 |
} |
| 116 |
|
| 117 |
public function createExportTableStream($length, $tableName) |
| 118 |
{ |
| 119 |
$stream = new MWP_Stream_Append(); |
| 120 |
|
| 121 |
$columns = $this->getConnection() |
| 122 |
->query("SHOW COLUMNS IN `{$tableName}`;") |
| 123 |
->fetchAll(); |
| 124 |
|
| 125 |
if (is_array($columns)) { |
| 126 |
$columns = $this->repack($columns, 'Field'); |
| 127 |
} |
| 128 |
|
| 129 |
$query = $this->selectAllDataQuery($tableName, $columns); |
| 130 |
$statement = $this->getConnection()->query($query, true); |
| 131 |
|
| 132 |
// Go through row by row |
| 133 |
if (!$this->options->isSkipLockTables()) { |
| 134 |
$stream->addStream(MWP_Stream_Stream::factory("LOCK TABLES `$tableName` WRITE;\n")); |
| 135 |
} |
| 136 |
|
| 137 |
$stream->addStream(MWP_Stream_Stream::factory("/*!40000 ALTER TABLE `$tableName` DISABLE KEYS */;\n")); |
| 138 |
|
| 139 |
$stream->addStream( |
| 140 |
new MWP_Stream_Callable(array($this, 'createExportRowStream'), array($statement, $tableName, $columns)) |
| 141 |
); |
| 142 |
|
| 143 |
$stream->addStream(MWP_Stream_Stream::factory("\n")); |
| 144 |
$stream->addStream(MWP_Stream_Stream::factory("/*!40000 ALTER TABLE `$tableName` ENABLE KEYS */;\n")); |
| 145 |
|
| 146 |
if (!$this->options->isSkipLockTables()) { |
| 147 |
$stream->addStream(MWP_Stream_Stream::factory("UNLOCK TABLES;\n")); |
| 148 |
} |
| 149 |
|
| 150 |
return $stream; |
| 151 |
} |
| 152 |
|
| 153 |
public function createExportRowStream($length, MWP_IncrementalBackup_Database_StatementInterface $statement, $tableName, $columns) |
| 154 |
{ |
| 155 |
$row = $statement->fetch(); |
| 156 |
if (!$row) { |
| 157 |
// This statement is using unbuffered queries and MUST be closed explicitly. |
| 158 |
$statement->close(); |
| 159 |
|
| 160 |
return false; |
| 161 |
} |
| 162 |
|
| 163 |
return $this->createRowInsertStatement($tableName, $row, $columns)."\n"; |
| 164 |
} |
| 165 |
|
| 166 |
/** |
| 167 |
* Repacks an array by making a key of a particular column |
| 168 |
* |
| 169 |
* @param array $array |
| 170 |
* @param $column |
| 171 |
* |
| 172 |
* @return array |
| 173 |
*/ |
| 174 |
protected function repack(array $array, $column) |
| 175 |
{ |
| 176 |
$repacked = array(); |
| 177 |
foreach ($array as $element) { |
| 178 |
$repacked[$element[$column]] = $element; |
| 179 |
} |
| 180 |
|
| 181 |
return $repacked; |
| 182 |
} |
| 183 |
|
| 184 |
/** |
| 185 |
* Creates an SQL statement for fetching all data from a particular table |
| 186 |
* |
| 187 |
* @param $tableName |
| 188 |
* @param $columnData |
| 189 |
* |
| 190 |
* @return string |
| 191 |
*/ |
| 192 |
protected function selectAllDataQuery($tableName, $columnData) |
| 193 |
{ |
| 194 |
$columns = array(); |
| 195 |
foreach ($columnData as $columnName => $metadata) { |
| 196 |
if (strpos($metadata['Type'], 'blob') !== false) { |
| 197 |
$fullColumnName = "`{$tableName}`.`{$columnName}`"; |
| 198 |
$columns[] = "HEX($fullColumnName) as `{$columnName}`"; |
| 199 |
} else { |
| 200 |
$columns[] = "`{$tableName}`.`{$columnName}`"; |
| 201 |
} |
| 202 |
} |
| 203 |
$cols = join(', ', $columns); |
| 204 |
$sql = "SELECT $cols FROM `$tableName`;"; |
| 205 |
|
| 206 |
return $sql; |
| 207 |
} |
| 208 |
|
| 209 |
/** |
| 210 |
* Creates an sql statement for row insertion |
| 211 |
* |
| 212 |
* @param string $tableName |
| 213 |
* @param array $row |
| 214 |
* @param array $columns |
| 215 |
* |
| 216 |
* @return string |
| 217 |
*/ |
| 218 |
protected function createRowInsertStatement($tableName, array $row, array $columns = array()) |
| 219 |
{ |
| 220 |
$values = $this->createRowInsertValues($row, $columns); |
| 221 |
$joined = join(', ', $values); |
| 222 |
$sql = "INSERT INTO `$tableName` VALUES($joined);"; |
| 223 |
|
| 224 |
return $sql; |
| 225 |
} |
| 226 |
|
| 227 |
protected function createRowInsertValues($row, $columns) |
| 228 |
{ |
| 229 |
$values = array(); |
| 230 |
|
| 231 |
foreach ($row as $columnName => $value) { |
| 232 |
$type = $columns[$columnName]['Type']; |
| 233 |
|
| 234 |
// Used to determine if the column is enum in case some of the allowed values contain reserved type identifiers |
| 235 |
$trimmedType = strtolower(trim($type)); |
| 236 |
|
| 237 |
// If it should not be enclosed |
| 238 |
if ($value === null) { |
| 239 |
$values[] = 'null'; |
| 240 |
} elseif (strpos($trimmedType, 'enum') !== 0 && |
| 241 |
(strpos($type, 'int') !== false |
| 242 |
|| strpos($type, 'float') !== false |
| 243 |
|| strpos($type, 'double') !== false |
| 244 |
|| strpos($type, 'decimal') !== false |
| 245 |
|| strpos($type, 'bool') !== false) |
| 246 |
) { |
| 247 |
$values[] = $value; |
| 248 |
} elseif (strpos($type, 'blob') !== false) { |
| 249 |
$values[] = strlen($value) ? ('0x'.$value) : "''"; |
| 250 |
} else { |
| 251 |
$values[] = $this->getConnection()->quote($value); |
| 252 |
} |
| 253 |
} |
| 254 |
|
| 255 |
return $values; |
| 256 |
} |
| 257 |
|
| 258 |
private static function arrayColumn($array, $columnIndex = 0) |
| 259 |
{ |
| 260 |
$result = array(); |
| 261 |
foreach ($array as $arr) { |
| 262 |
if (!is_array($arr)) { |
| 263 |
continue; |
| 264 |
} |
| 265 |
$arr = array_values($arr); |
| 266 |
$result[] = $arr[$columnIndex]; |
| 267 |
} |
| 268 |
return $result; |
| 269 |
} |
| 270 |
} |
| 271 |
|