Export
2 months ago
Import
2 weeks ago
PersonalDataExporters
2 months ago
ImportExportFactory.php
2 months ago
ImportExportRepository.php
2 months ago
index.php
3 years ago
ImportExportRepository.php
397 lines
| 1 | <?php // phpcs:ignore SlevomatCodingStandard.TypeHints.DeclareStrictTypes.DeclareStrictTypesMissing |
| 2 | |
| 3 | namespace MailPoet\Subscribers\ImportExport; |
| 4 | |
| 5 | if (!defined('ABSPATH')) exit; |
| 6 | |
| 7 | |
| 8 | use DateTime; |
| 9 | use MailPoet\Config\SubscriberChangesNotifier; |
| 10 | use MailPoet\Entities\CustomFieldEntity; |
| 11 | use MailPoet\Entities\SegmentEntity; |
| 12 | use MailPoet\Entities\SubscriberCustomFieldEntity; |
| 13 | use MailPoet\Entities\SubscriberEntity; |
| 14 | use MailPoet\Entities\SubscriberSegmentEntity; |
| 15 | use MailPoet\Segments\DynamicSegments\FilterHandler; |
| 16 | use MailPoet\Subscribers\SubscriberCustomFieldRepository; |
| 17 | use MailPoet\Subscribers\SubscribersRepository; |
| 18 | use MailPoetVendor\Doctrine\DBAL\ArrayParameterType; |
| 19 | use MailPoetVendor\Doctrine\DBAL\Query\QueryBuilder; |
| 20 | use MailPoetVendor\Doctrine\DBAL\Result; |
| 21 | use MailPoetVendor\Doctrine\ORM\EntityManager; |
| 22 | use MailPoetVendor\Doctrine\ORM\Mapping\ClassMetadata; |
| 23 | |
| 24 | class ImportExportRepository { |
| 25 | private const IGNORED_COLUMNS_FOR_BULK_UPDATE = [ |
| 26 | SubscriberEntity::class => [ |
| 27 | 'wp_user_id', |
| 28 | 'is_woocommerce_user', |
| 29 | 'email', |
| 30 | 'created_at', |
| 31 | 'last_subscribed_at', |
| 32 | ], |
| 33 | SubscriberCustomFieldEntity::class => [ |
| 34 | 'created_at', |
| 35 | ], |
| 36 | SubscriberSegmentEntity::class => [ |
| 37 | 'created_at', |
| 38 | ], |
| 39 | ]; |
| 40 | |
| 41 | private const KEY_COLUMNS_FOR_BULK_UPDATE = [ |
| 42 | SubscriberEntity::class => [ |
| 43 | 'email', |
| 44 | ], |
| 45 | SubscriberCustomFieldEntity::class => [ |
| 46 | 'subscriber_id', |
| 47 | 'custom_field_id', |
| 48 | ], |
| 49 | ]; |
| 50 | |
| 51 | /** @var EntityManager */ |
| 52 | protected $entityManager; |
| 53 | |
| 54 | /** @var SubscriberChangesNotifier */ |
| 55 | private $subscriberChangesNotifier; |
| 56 | |
| 57 | /** @var FilterHandler */ |
| 58 | private $filterHandler; |
| 59 | |
| 60 | /** @var SubscribersRepository */ |
| 61 | private $subscribersRepository; |
| 62 | |
| 63 | /** @var SubscriberCustomFieldRepository */ |
| 64 | private $subscriberCustomFieldRepository; |
| 65 | |
| 66 | public function __construct( |
| 67 | EntityManager $entityManager, |
| 68 | SubscriberChangesNotifier $changesNotifier, |
| 69 | FilterHandler $filterHandler, |
| 70 | SubscribersRepository $subscribersRepository, |
| 71 | SubscriberCustomFieldRepository $subscriberCustomFieldRepository |
| 72 | ) { |
| 73 | $this->entityManager = $entityManager; |
| 74 | $this->subscriberChangesNotifier = $changesNotifier; |
| 75 | $this->filterHandler = $filterHandler; |
| 76 | $this->subscribersRepository = $subscribersRepository; |
| 77 | $this->subscriberCustomFieldRepository = $subscriberCustomFieldRepository; |
| 78 | } |
| 79 | |
| 80 | /** |
| 81 | * @param class-string<object> $className |
| 82 | * @return ClassMetadata<object> |
| 83 | */ |
| 84 | protected function getClassMetadata(string $className): ClassMetadata { |
| 85 | return $this->entityManager->getClassMetadata($className); |
| 86 | } |
| 87 | |
| 88 | /** |
| 89 | * @param class-string<object> $className |
| 90 | */ |
| 91 | protected function getTableName(string $className): string { |
| 92 | return $this->getClassMetadata($className)->getTableName(); |
| 93 | } |
| 94 | |
| 95 | /** |
| 96 | * @param class-string<object> $className |
| 97 | */ |
| 98 | protected function getTableColumns(string $className): array { |
| 99 | return $this->getClassMetadata($className)->getColumnNames(); |
| 100 | } |
| 101 | |
| 102 | /** |
| 103 | * @param class-string<object> $className |
| 104 | */ |
| 105 | public function insertMultiple( |
| 106 | string $className, |
| 107 | array $columns, |
| 108 | array $data |
| 109 | ): int { |
| 110 | $tableName = $this->getTableName($className); |
| 111 | |
| 112 | if (!$columns || !$data) { |
| 113 | return 0; |
| 114 | } |
| 115 | |
| 116 | $rows = []; |
| 117 | $parameters = []; |
| 118 | foreach ($data as $key => $item) { |
| 119 | $paramNames = array_map(function (string $parameter) use ($key): string { |
| 120 | return ":{$parameter}_{$key}"; |
| 121 | }, $columns); |
| 122 | |
| 123 | foreach ($item as $columnKey => $column) { |
| 124 | // We need to remove the colon character from the query parameter name that is passed to the query builder |
| 125 | $parameters[substr($paramNames[$columnKey], 1)] = $column; |
| 126 | } |
| 127 | $rows[] = "(" . implode(', ', $paramNames) . ")"; |
| 128 | } |
| 129 | |
| 130 | $count = (int)$this->entityManager->getConnection()->executeStatement(" |
| 131 | INSERT IGNORE INTO {$tableName} (`" . implode("`, `", $columns) . "`) VALUES |
| 132 | " . implode(", \n", $rows) . " |
| 133 | ", $parameters); |
| 134 | $this->notifyCreations($className, $columns, $data); |
| 135 | return $count; |
| 136 | } |
| 137 | |
| 138 | /** |
| 139 | * @param class-string<object> $className |
| 140 | */ |
| 141 | public function updateMultiple( |
| 142 | string $className, |
| 143 | array $columns, |
| 144 | array $data, |
| 145 | ?DateTime $updatedAt = null |
| 146 | ): int { |
| 147 | $tableName = $this->getTableName($className); |
| 148 | $entityColumns = $this->getTableColumns($className); |
| 149 | |
| 150 | if (!$columns || !$data) { |
| 151 | return 0; |
| 152 | } |
| 153 | |
| 154 | $parameters = []; |
| 155 | $parameterTypes = []; |
| 156 | $keyColumns = self::KEY_COLUMNS_FOR_BULK_UPDATE[$className] ?? []; |
| 157 | if (!$keyColumns) { |
| 158 | return 0; |
| 159 | } |
| 160 | |
| 161 | $keyColumnsConditions = []; |
| 162 | foreach ($keyColumns as $keyColumn) { |
| 163 | $columnIndex = array_search($keyColumn, $columns); |
| 164 | $parameters[$keyColumn] = array_map(function(array $row) use ($columnIndex) { |
| 165 | return $row[$columnIndex]; |
| 166 | }, $data); |
| 167 | $parameterTypes[$keyColumn] = ArrayParameterType::STRING; |
| 168 | $keyColumnsConditions[] = "{$keyColumn} IN (:{$keyColumn})"; |
| 169 | } |
| 170 | |
| 171 | $restoredSubscriberIds = $className === SubscriberEntity::class |
| 172 | ? $this->getDeletedSubscriberIdsByEmail($columns, $data) |
| 173 | : []; |
| 174 | |
| 175 | $ignoredColumns = self::IGNORED_COLUMNS_FOR_BULK_UPDATE[$className] ?? ['created_at']; |
| 176 | $updateColumns = array_map(function($columnName) use ($keyColumns, $columns, $data, &$parameters): string { |
| 177 | $values = []; |
| 178 | foreach ($data as $index => $row) { |
| 179 | $keyCondition = array_map(function($keyColumn) use ($index, $row, $columns, &$parameters): string { |
| 180 | $parameters["{$keyColumn}_{$index}"] = $row[array_search($keyColumn, $columns)]; |
| 181 | return "{$keyColumn} = :{$keyColumn}_{$index}"; |
| 182 | }, $keyColumns); |
| 183 | $values[] = "WHEN " . implode(' AND ', $keyCondition) . " THEN :{$columnName}_{$index}"; |
| 184 | $parameters["{$columnName}_{$index}"] = $row[array_search($columnName, $columns)]; |
| 185 | } |
| 186 | return "{$columnName} = (CASE " . implode("\n", $values) . " END)"; |
| 187 | }, array_diff($columns, $ignoredColumns)); |
| 188 | |
| 189 | if ($updatedAt && in_array('updated_at', $entityColumns, true)) { |
| 190 | $parameters['updated_at'] = $updatedAt; |
| 191 | $updateColumns[] = "updated_at = :updated_at"; |
| 192 | } |
| 193 | |
| 194 | // we want to reset deleted_at for updated rows |
| 195 | if (in_array('deleted_at', $entityColumns, true)) { |
| 196 | $updateColumns[] = 'deleted_at = NULL'; |
| 197 | } |
| 198 | |
| 199 | $count = (int)$this->entityManager->getConnection()->executeStatement(" |
| 200 | UPDATE {$tableName} SET |
| 201 | " . implode(", \n", $updateColumns) . " |
| 202 | WHERE |
| 203 | " . implode(' AND ', $keyColumnsConditions) . " |
| 204 | ", $parameters, $parameterTypes); |
| 205 | $updatedSubscriberIds = $this->notifyUpdates($className, $columns, $data); |
| 206 | $countChangedSubscriberIds = $this->getCountChangedSubscriberIdsForBulkUpdate($className, $columns, $updatedSubscriberIds, $restoredSubscriberIds); |
| 207 | if ($countChangedSubscriberIds) { |
| 208 | $this->subscriberChangesNotifier->subscribersCountChanged($countChangedSubscriberIds); |
| 209 | } |
| 210 | if ($className === SubscriberEntity::class) { |
| 211 | $this->subscribersRepository->refreshAll(); |
| 212 | } |
| 213 | if ($className === SubscriberCustomFieldEntity::class) { |
| 214 | $this->subscriberCustomFieldRepository->refreshAll(); |
| 215 | } |
| 216 | return $count; |
| 217 | } |
| 218 | |
| 219 | public function getSubscribersBatchBySegment(?SegmentEntity $segment, int $limit, int $offset = 0): array { |
| 220 | $subscriberSegmentTable = $this->getTableName(SubscriberSegmentEntity::class); |
| 221 | $subscriberTable = $this->getTableName(SubscriberEntity::class); |
| 222 | $segmentTable = $this->getTableName(SegmentEntity::class); |
| 223 | |
| 224 | $qb = $this->createSubscribersQueryBuilder($limit, $offset); |
| 225 | $qb = $this->addSubscriberCustomFieldsToQueryBuilder($qb); |
| 226 | |
| 227 | if (!$segment || $segment->isStatic()) { |
| 228 | // joining with the segments table is used only when there is no segment or for static segments. |
| 229 | // this because dynamic segments don't have a corresponding entry in the segments table. |
| 230 | $qb->leftJoin($subscriberSegmentTable, $segmentTable, $segmentTable, "{$segmentTable}.id = {$subscriberSegmentTable}.segment_id") |
| 231 | ->groupBy("{$subscriberTable}.id, {$segmentTable}.id"); |
| 232 | } |
| 233 | |
| 234 | if (!$segment) { |
| 235 | // if there are subscribers who do not belong to any segment, use |
| 236 | // a CASE function to group them under "Not In Segment" |
| 237 | $qb->addSelect("'" . __('Not In Segment', 'mailpoet') . "' AS segment_name") |
| 238 | ->leftJoin($subscriberTable, $subscriberTable, 's2', "{$subscriberTable}.id = s2.id") |
| 239 | ->leftJoin('s2', $subscriberSegmentTable, 'ssg2', "s2.id = ssg2.subscriber_id AND ssg2.status = :statusSubscribed AND {$segmentTable}.id <> ssg2.segment_id") |
| 240 | ->leftJoin('ssg2', $segmentTable, 'sg2', 'ssg2.segment_id = sg2.id AND sg2.deleted_at IS NULL') |
| 241 | ->andWhere("({$subscriberSegmentTable}.status != :statusSubscribed OR {$subscriberSegmentTable}.id IS NULL OR {$segmentTable}.deleted_at IS NOT NULL)") |
| 242 | ->andWhere('sg2.id IS NULL') |
| 243 | ->setParameter('statusSubscribed', SubscriberEntity::STATUS_SUBSCRIBED); |
| 244 | } elseif ($segment->isStatic()) { |
| 245 | $qb->addSelect("{$segmentTable}.name AS segment_name") |
| 246 | ->andWhere("{$subscriberSegmentTable}.segment_id = :segmentId") |
| 247 | ->setParameter('segmentId', $segment->getId()); |
| 248 | } else { |
| 249 | // Dynamic segments don't have a relation to the segment table, |
| 250 | // So we need to use a placeholder |
| 251 | $qb->addSelect(":segmentName AS segment_name") |
| 252 | ->setParameter('segmentName', $segment->getName()) |
| 253 | ->groupBy("{$subscriberTable}.id"); |
| 254 | $qb = $this->filterHandler->apply($qb, $segment); |
| 255 | } |
| 256 | |
| 257 | $statement = $qb->execute(); |
| 258 | return $statement instanceof Result ? $statement->fetchAll() : []; |
| 259 | } |
| 260 | |
| 261 | private function createSubscribersQueryBuilder(int $limit, int $offset): QueryBuilder { |
| 262 | $subscriberSegmentTable = $this->getTableName(SubscriberSegmentEntity::class); |
| 263 | $subscriberTable = $this->getTableName(SubscriberEntity::class); |
| 264 | |
| 265 | return $this->entityManager->getConnection()->createQueryBuilder() |
| 266 | ->select(" |
| 267 | {$subscriberTable}.first_name, |
| 268 | {$subscriberTable}.last_name, |
| 269 | {$subscriberTable}.email, |
| 270 | {$subscriberTable}.subscribed_ip, |
| 271 | {$subscriberTable}.confirmed_at, |
| 272 | {$subscriberTable}.confirmed_ip, |
| 273 | {$subscriberTable}.created_at, |
| 274 | {$subscriberTable}.last_subscribed_at, |
| 275 | {$subscriberTable}.status AS global_status, |
| 276 | {$subscriberSegmentTable}.status AS list_status |
| 277 | ") |
| 278 | ->from($subscriberTable) |
| 279 | ->leftJoin($subscriberTable, $subscriberSegmentTable, $subscriberSegmentTable, "{$subscriberTable}.id = {$subscriberSegmentTable}.subscriber_id") |
| 280 | ->andWhere("{$subscriberTable}.deleted_at IS NULL") |
| 281 | ->orderBy("{$subscriberTable}.id") |
| 282 | ->setFirstResult($offset) |
| 283 | ->setMaxResults($limit); |
| 284 | } |
| 285 | |
| 286 | private function addSubscriberCustomFieldsToQueryBuilder(QueryBuilder $qb): QueryBuilder { |
| 287 | $segmentsTable = $this->getTableName(SubscriberEntity::class); |
| 288 | $customFieldsTable = $this->getTableName(CustomFieldEntity::class); |
| 289 | $subscriberCustomFieldTable = $this->getTableName(SubscriberCustomFieldEntity::class); |
| 290 | |
| 291 | $customFields = $this->entityManager->getConnection()->createQueryBuilder() |
| 292 | ->select("{$customFieldsTable}.*") |
| 293 | ->from($customFieldsTable) |
| 294 | ->execute(); |
| 295 | |
| 296 | $customFields = $customFields->fetchAll(); |
| 297 | |
| 298 | foreach ($customFields as $customField) { |
| 299 | if (!is_array($customField) || !is_numeric($customField['id'] ?? null)) { |
| 300 | continue; |
| 301 | } |
| 302 | $rawId = (int)$customField['id']; |
| 303 | $customFieldId = "customFieldId{$rawId}export"; |
| 304 | $qb->addSelect("MAX(CASE WHEN {$customFieldsTable}.id = :{$customFieldId} THEN {$subscriberCustomFieldTable}.value END) AS :{$customFieldId}") |
| 305 | ->setParameter($customFieldId, $rawId); |
| 306 | } |
| 307 | |
| 308 | $qb->leftJoin($segmentsTable, $subscriberCustomFieldTable, $subscriberCustomFieldTable, "{$segmentsTable}.id = {$subscriberCustomFieldTable}.subscriber_id") |
| 309 | ->leftJoin($subscriberCustomFieldTable, $customFieldsTable, $customFieldsTable, "{$customFieldsTable}.id = {$subscriberCustomFieldTable}.custom_field_id"); |
| 310 | |
| 311 | return $qb; |
| 312 | } |
| 313 | |
| 314 | private function notifyCreations(string $className, array $columns, array $data): void { |
| 315 | if ($className === SubscriberEntity::class) { |
| 316 | $ids = $this->getIdsByEmail($className, $columns, $data); |
| 317 | $this->subscriberChangesNotifier->subscribersCreated($ids); |
| 318 | } |
| 319 | } |
| 320 | |
| 321 | private function notifyUpdates(string $className, array $columns, array $data): array { |
| 322 | if ($className === SubscriberEntity::class) { |
| 323 | $ids = $this->getIdsByEmail($className, $columns, $data); |
| 324 | $this->subscriberChangesNotifier->subscribersUpdated($ids); |
| 325 | return $ids; |
| 326 | } |
| 327 | return []; |
| 328 | } |
| 329 | |
| 330 | private function getCountChangedSubscriberIdsForBulkUpdate(string $className, array $columns, array $updatedSubscriberIds, array $restoredSubscriberIds): array { |
| 331 | if ($className !== SubscriberEntity::class) { |
| 332 | return []; |
| 333 | } |
| 334 | |
| 335 | $subscriberIds = $restoredSubscriberIds; |
| 336 | if (in_array('status', $columns, true)) { |
| 337 | $subscriberIds = array_merge($subscriberIds, $updatedSubscriberIds); |
| 338 | } |
| 339 | |
| 340 | return array_values(array_unique(array_map('intval', $subscriberIds))); |
| 341 | } |
| 342 | |
| 343 | /** |
| 344 | * @param class-string<object> $className |
| 345 | */ |
| 346 | private function getIdsByEmail(string $className, array $columns, array $data): array { |
| 347 | $tableName = $this->getTableName($className); |
| 348 | $emailIndex = array_search('email', $columns); |
| 349 | if ($emailIndex === false) { |
| 350 | return []; |
| 351 | } |
| 352 | $emails = []; |
| 353 | foreach ($data as $item) { |
| 354 | $emails[] = $item[$emailIndex]; |
| 355 | } |
| 356 | // get ids for updated/created rows |
| 357 | return $this->entityManager->getConnection()->executeQuery(" |
| 358 | SELECT id |
| 359 | FROM {$tableName} |
| 360 | WHERE email IN (:emails) |
| 361 | ", ['emails' => $emails], ['emails' => ArrayParameterType::STRING])->fetchFirstColumn(); |
| 362 | } |
| 363 | |
| 364 | private function getDeletedSubscriberIdsByEmail(array $columns, array $data): array { |
| 365 | $emailColumnIndex = array_search('email', $columns, true); |
| 366 | if ($emailColumnIndex === false) { |
| 367 | return []; |
| 368 | } |
| 369 | |
| 370 | $emails = array_map(function(array $row) use ($emailColumnIndex): string { |
| 371 | return (string)$row[$emailColumnIndex]; |
| 372 | }, $data); |
| 373 | |
| 374 | if (!$emails) { |
| 375 | return []; |
| 376 | } |
| 377 | |
| 378 | $subscriberTable = $this->getTableName(SubscriberEntity::class); |
| 379 | $subscriberIds = $this->entityManager->getConnection()->executeQuery( |
| 380 | "SELECT `id` |
| 381 | FROM {$subscriberTable} |
| 382 | WHERE `email` IN (:emails) |
| 383 | AND `deleted_at` IS NOT NULL", |
| 384 | [ |
| 385 | 'emails' => $emails, |
| 386 | ], |
| 387 | [ |
| 388 | 'emails' => ArrayParameterType::STRING, |
| 389 | ] |
| 390 | )->fetchFirstColumn(); |
| 391 | |
| 392 | return array_values(array_map('intval', array_filter($subscriberIds, static function($id): bool { |
| 393 | return is_int($id) || (is_string($id) && ctype_digit($id)); |
| 394 | }))); |
| 395 | } |
| 396 | } |
| 397 |