Cli.php
2 weeks ago
Import.php
2 weeks ago
MailChimp.php
2 months ago
MailChimpDataMapper.php
2 months ago
index.php
3 years ago
Import.php
709 lines
| 1 | <?php // phpcs:ignore SlevomatCodingStandard.TypeHints.DeclareStrictTypes.DeclareStrictTypesMissing |
| 2 | |
| 3 | namespace MailPoet\Subscribers\ImportExport\Import; |
| 4 | |
| 5 | if (!defined('ABSPATH')) exit; |
| 6 | |
| 7 | |
| 8 | use MailPoet\CustomFields\CustomFieldsRepository; |
| 9 | use MailPoet\Entities\CustomFieldEntity; |
| 10 | use MailPoet\Entities\SubscriberCustomFieldEntity; |
| 11 | use MailPoet\Entities\SubscriberEntity; |
| 12 | use MailPoet\Entities\SubscriberSegmentEntity; |
| 13 | use MailPoet\Entities\SubscriberTagEntity; |
| 14 | use MailPoet\Newsletter\Options\NewsletterOptionsRepository; |
| 15 | use MailPoet\Segments\WP; |
| 16 | use MailPoet\Services\Validator; |
| 17 | use MailPoet\Subscribers\ImportExport\ImportExportFactory; |
| 18 | use MailPoet\Subscribers\ImportExport\ImportExportRepository; |
| 19 | use MailPoet\Subscribers\Source; |
| 20 | use MailPoet\Subscribers\SubscribersRepository; |
| 21 | use MailPoet\Tags\TagRepository; |
| 22 | use MailPoet\Util\DateConverter; |
| 23 | use MailPoet\Util\Helpers; |
| 24 | use MailPoet\Util\Security; |
| 25 | use MailPoet\WP\Functions as WPFunctions; |
| 26 | use MailPoetVendor\Carbon\Carbon; |
| 27 | |
| 28 | class Import { |
| 29 | /** @var array */ |
| 30 | public $subscribersData; |
| 31 | /** @var array */ |
| 32 | public $segmentsIds; |
| 33 | /** @var string[] */ |
| 34 | public $tags; |
| 35 | /** @var string */ |
| 36 | public $newSubscribersStatus; |
| 37 | /** @var string */ |
| 38 | public $existingSubscribersStatus; |
| 39 | /** @var bool */ |
| 40 | public $updateSubscribers; |
| 41 | /** @var array */ |
| 42 | public $subscribersFields; |
| 43 | /** @var array */ |
| 44 | public $subscribersCustomFields; |
| 45 | /** @var int */ |
| 46 | public $subscribersCount; |
| 47 | /** @var Carbon */ |
| 48 | public $createdAt; |
| 49 | /** @var Carbon */ |
| 50 | public $updatedAt; |
| 51 | /** @var array<string, mixed> */ |
| 52 | public $requiredSubscribersFields; |
| 53 | const DB_QUERY_CHUNK_SIZE = 100; |
| 54 | const STATUS_DONT_UPDATE = 'dont_update'; |
| 55 | |
| 56 | public const ACTION_CREATE = 'create'; |
| 57 | public const ACTION_UPDATE = 'update'; |
| 58 | |
| 59 | /** @var WP */ |
| 60 | private $wpSegment; |
| 61 | |
| 62 | /** @var CustomFieldsRepository */ |
| 63 | private $customFieldsRepository; |
| 64 | |
| 65 | /** @var ImportExportRepository */ |
| 66 | private $importExportRepository; |
| 67 | |
| 68 | /** @var NewsletterOptionsRepository */ |
| 69 | private $newsletterOptionsRepository; |
| 70 | |
| 71 | /** @var SubscribersRepository */ |
| 72 | private $subscriberRepository; |
| 73 | |
| 74 | /** @var TagRepository */ |
| 75 | private $tagRepository; |
| 76 | |
| 77 | /** @var Validator */ |
| 78 | private $validator; |
| 79 | |
| 80 | public function __construct( |
| 81 | WP $wpSegment, |
| 82 | CustomFieldsRepository $customFieldsRepository, |
| 83 | ImportExportRepository $importExportRepository, |
| 84 | NewsletterOptionsRepository $newsletterOptionsRepository, |
| 85 | SubscribersRepository $subscriberRepository, |
| 86 | TagRepository $tagRepository, |
| 87 | Validator $validator, |
| 88 | array $data |
| 89 | ) { |
| 90 | $this->wpSegment = $wpSegment; |
| 91 | $this->customFieldsRepository = $customFieldsRepository; |
| 92 | $this->importExportRepository = $importExportRepository; |
| 93 | $this->newsletterOptionsRepository = $newsletterOptionsRepository; |
| 94 | $this->subscriberRepository = $subscriberRepository; |
| 95 | $this->tagRepository = $tagRepository; |
| 96 | $this->validator = $validator; |
| 97 | $this->validateImportData($data); |
| 98 | $this->subscribersData = $this->transformSubscribersData( |
| 99 | $data['subscribers'], |
| 100 | $data['columns'] |
| 101 | ); |
| 102 | $this->segmentsIds = $data['segments']; |
| 103 | $this->tags = $data['tags']; |
| 104 | $this->newSubscribersStatus = $data['newSubscribersStatus']; |
| 105 | $this->existingSubscribersStatus = $data['existingSubscribersStatus']; |
| 106 | $this->updateSubscribers = $data['updateSubscribers']; |
| 107 | $this->subscribersFields = $this->getSubscribersFields( |
| 108 | array_keys($data['columns']) |
| 109 | ); |
| 110 | $this->subscribersCustomFields = $this->getCustomSubscribersFields( |
| 111 | array_keys($data['columns']) |
| 112 | ); |
| 113 | $this->subscribersCount = (reset($this->subscribersData) === false) ? 0 : count(reset($this->subscribersData)); |
| 114 | $this->createdAt = Carbon::now()->millisecond(0); |
| 115 | $this->updatedAt = Carbon::createFromTimestamp(WPFunctions::get()->currentTime('timestamp', true) + 1); |
| 116 | $this->requiredSubscribersFields = [ |
| 117 | 'status' => SubscriberEntity::STATUS_SUBSCRIBED, |
| 118 | 'first_name' => '', |
| 119 | 'last_name' => '', |
| 120 | 'created_at' => $this->createdAt, |
| 121 | ]; |
| 122 | } |
| 123 | |
| 124 | public function validateImportData(array $data): void { |
| 125 | $requiredDataFields = [ |
| 126 | 'subscribers', |
| 127 | 'columns', |
| 128 | 'segments', |
| 129 | 'timestamp', |
| 130 | 'newSubscribersStatus', |
| 131 | 'existingSubscribersStatus', |
| 132 | 'updateSubscribers', |
| 133 | 'tags', |
| 134 | ]; |
| 135 | // 1. data should contain all required fields |
| 136 | // 2. column names should only contain alphanumeric & underscore characters |
| 137 | if ( |
| 138 | count(array_intersect_key(array_flip($requiredDataFields), $data)) !== count($requiredDataFields) || |
| 139 | preg_grep('/[^a-zA-Z0-9_]/', array_keys($data['columns'])) |
| 140 | ) { |
| 141 | throw new \Exception(__('Missing or invalid import data.', 'mailpoet')); |
| 142 | } |
| 143 | } |
| 144 | |
| 145 | /** |
| 146 | * @return array{created: int, updated:int, segments: array, added_to_segment_with_welcome_notification:bool} |
| 147 | * @throws \Exception |
| 148 | */ |
| 149 | public function process(): array { |
| 150 | // validate data based on field validation rules |
| 151 | $subscribersData = $this->validateSubscribersData($this->subscribersData); |
| 152 | if (!$subscribersData) { |
| 153 | throw new \Exception(__('No valid subscribers were found.', 'mailpoet')); |
| 154 | } |
| 155 | // permanently trash deleted subscribers |
| 156 | $this->deleteExistingTrashedSubscribers($subscribersData); |
| 157 | |
| 158 | // split subscribers into "existing" and "new" and free up memory |
| 159 | $existingSubscribers = $newSubscribers = [ |
| 160 | 'data' => [], |
| 161 | 'fields' => $this->subscribersFields, |
| 162 | ]; |
| 163 | list($existingSubscribers['data'], $newSubscribers['data'], $wpUsers) = |
| 164 | $this->splitSubscribersData($subscribersData); |
| 165 | $subscribersData = null; |
| 166 | |
| 167 | // create or update subscribers |
| 168 | $createdSubscribers = $updatedSubscribers = []; |
| 169 | try { |
| 170 | if ($newSubscribers['data']) { |
| 171 | // add, if required, missing required fields to new subscribers |
| 172 | $newSubscribers = $this->addMissingRequiredFields($newSubscribers); |
| 173 | $newSubscribers = $this->setSubscriptionStatusToDefault($newSubscribers, $this->newSubscribersStatus); |
| 174 | $newSubscribers = $this->setSource($newSubscribers); |
| 175 | $newSubscribers = $this->setLinkToken($newSubscribers); |
| 176 | $createdSubscribers = |
| 177 | $this->createOrUpdateSubscribers( |
| 178 | self::ACTION_CREATE, |
| 179 | $newSubscribers, |
| 180 | $this->subscribersCustomFields |
| 181 | ); |
| 182 | } |
| 183 | |
| 184 | $updateExistingSubscribersStatus = false; |
| 185 | |
| 186 | if ($existingSubscribers['data']) { |
| 187 | $allowedStatuses = [ |
| 188 | SubscriberEntity::STATUS_SUBSCRIBED, |
| 189 | SubscriberEntity::STATUS_UNSUBSCRIBED, |
| 190 | SubscriberEntity::STATUS_INACTIVE, |
| 191 | ]; |
| 192 | if (in_array($this->existingSubscribersStatus, $allowedStatuses, true)) { |
| 193 | $updateExistingSubscribersStatus = true; |
| 194 | $existingSubscribers = $this->addField($existingSubscribers, 'status', $this->existingSubscribersStatus); |
| 195 | } |
| 196 | if ($this->updateSubscribers) { |
| 197 | // Update existing subscribers' info (first_name, last_name etc.) |
| 198 | // as well as status (optionally) if the status column was added above |
| 199 | $updatedSubscribers = |
| 200 | $this->createOrUpdateSubscribers( |
| 201 | self::ACTION_UPDATE, |
| 202 | $existingSubscribers, |
| 203 | $this->subscribersCustomFields |
| 204 | ); |
| 205 | if ($wpUsers) { |
| 206 | $this->synchronizeWPUsers($wpUsers); |
| 207 | } |
| 208 | } elseif ($updateExistingSubscribersStatus) { |
| 209 | // Only update existing subscribers' status |
| 210 | // For this we need to remove all other fields except email and status |
| 211 | $existingSubscribers['fields'] = array_intersect($existingSubscribers['fields'], ['email', 'status']); |
| 212 | $existingSubscribers['data'] = array_intersect_key($existingSubscribers['data'], array_flip(['email', 'status'])); |
| 213 | $updatedSubscribers = |
| 214 | $this->createOrUpdateSubscribers( |
| 215 | self::ACTION_UPDATE, |
| 216 | $existingSubscribers |
| 217 | ); |
| 218 | } |
| 219 | } |
| 220 | } catch (\Exception $e) { |
| 221 | throw new \Exception(__('Unable to save imported subscribers.', 'mailpoet')); |
| 222 | } |
| 223 | |
| 224 | // check if any subscribers were added to segments that have welcome notifications configured |
| 225 | $importFactory = new ImportExportFactory('import'); |
| 226 | $segments = $importFactory->getSegments(); |
| 227 | $welcomeNotificationsInSegments = |
| 228 | ($createdSubscribers || $updatedSubscribers) ? |
| 229 | $this->newsletterOptionsRepository->findWelcomeNotificationsForSegments($this->segmentsIds) : |
| 230 | false; |
| 231 | |
| 232 | return [ |
| 233 | 'created' => is_array($createdSubscribers) ? count($createdSubscribers) : 0, |
| 234 | 'updated' => is_array($updatedSubscribers) ? count($updatedSubscribers) : 0, |
| 235 | 'segments' => $segments, |
| 236 | 'added_to_segment_with_welcome_notification' => |
| 237 | ($welcomeNotificationsInSegments) ? true : false, |
| 238 | ]; |
| 239 | } |
| 240 | |
| 241 | /** |
| 242 | * @param array $subscribersData |
| 243 | * @return false|array |
| 244 | */ |
| 245 | public function validateSubscribersData(array $subscribersData) { |
| 246 | $invalidRecords = []; |
| 247 | foreach ($subscribersData as $column => &$data) { |
| 248 | if ($column === 'email') { |
| 249 | $data = array_map( |
| 250 | function($index, $email) use(&$invalidRecords) { |
| 251 | if (!$this->validator->validateNonRoleEmail($email)) { |
| 252 | $invalidRecords[] = $index; |
| 253 | } |
| 254 | return strtolower($email); |
| 255 | }, |
| 256 | array_keys($data), |
| 257 | $data |
| 258 | ); |
| 259 | } |
| 260 | if (in_array($column, ['created_at', 'confirmed_at'], true)) { |
| 261 | $data = $this->validateDateTime($data, $invalidRecords); |
| 262 | } |
| 263 | if (in_array($column, ['confirmed_ip', 'subscribed_ip'], true)) { |
| 264 | $data = array_map( |
| 265 | function($index, $ip) { |
| 266 | if (!filter_var($ip, FILTER_VALIDATE_IP)) { |
| 267 | // if invalid or empty, we allow the import but remove the IP |
| 268 | return null; |
| 269 | } |
| 270 | return $ip; |
| 271 | }, |
| 272 | array_keys($data), |
| 273 | $data |
| 274 | ); |
| 275 | } |
| 276 | // if this is a custom column |
| 277 | if (in_array($column, $this->subscribersCustomFields)) { |
| 278 | $customField = $this->customFieldsRepository->findOneById($column); |
| 279 | if (!$customField instanceof CustomFieldEntity) { |
| 280 | continue; |
| 281 | } |
| 282 | // validate date type |
| 283 | if ($customField->getType() === CustomFieldEntity::TYPE_DATE) { |
| 284 | $data = $this->validateDateTime($data, $invalidRecords); |
| 285 | } |
| 286 | } |
| 287 | } |
| 288 | if ($invalidRecords) { |
| 289 | foreach ($subscribersData as $column => &$data) { |
| 290 | $data = array_diff_key($data, array_flip($invalidRecords)); |
| 291 | $data = array_values($data); |
| 292 | } |
| 293 | } |
| 294 | if (empty($subscribersData['email'])) return false; |
| 295 | return $subscribersData; |
| 296 | } |
| 297 | |
| 298 | private function validateDateTime(array $data, array &$invalidRecords): array { |
| 299 | $siteUsesCustomFormat = WPFunctions::get()->getOption('date_format') === 'd/m/Y'; |
| 300 | if ($siteUsesCustomFormat) { |
| 301 | return $this->validateDateTimeAttemptCustomFormat($data, $invalidRecords); |
| 302 | } |
| 303 | |
| 304 | $validationRule = 'datetime'; |
| 305 | return array_map( |
| 306 | function ($index, $date) use ($validationRule, &$invalidRecords) { |
| 307 | if (empty($date)) return $date; |
| 308 | $date = (new DateConverter())->convertDateToDatetime($date, $validationRule); |
| 309 | if (!$date) { |
| 310 | $invalidRecords[] = $index; |
| 311 | } |
| 312 | return $date; |
| 313 | }, |
| 314 | array_keys($data), |
| 315 | $data |
| 316 | ); |
| 317 | } |
| 318 | |
| 319 | private function validateDateTimeAttemptCustomFormat(array $data, array &$invalidRecords): array { |
| 320 | $validationRule = 'datetime'; |
| 321 | $dateTimeDates = $data; |
| 322 | $dateTimeInvalidRecords = $invalidRecords; |
| 323 | $datetimeErrorCount = 0; |
| 324 | |
| 325 | $validationRuleCustom = 'd/m/Y'; |
| 326 | $customFormatDates = $data; |
| 327 | $customFormatInvalidRecords = $invalidRecords; |
| 328 | $customFormatErrorCount = 0; |
| 329 | |
| 330 | // We attempt converting with both date formats |
| 331 | foreach ($data as $index => $date) { |
| 332 | if (empty($date)) { |
| 333 | $dateTimeDates[$index] = $date; |
| 334 | $customFormatDates[$index] = $date; |
| 335 | continue; |
| 336 | }; |
| 337 | $dateTimeDates[$index] = (new DateConverter())->convertDateToDatetime($date, $validationRule); |
| 338 | if (!$dateTimeDates[$index]) { |
| 339 | $datetimeErrorCount ++; |
| 340 | $dateTimeInvalidRecords[] = $index; |
| 341 | } |
| 342 | $customFormatDates[$index] = (new DateConverter())->convertDateToDatetime($date, $validationRuleCustom); |
| 343 | if (!$customFormatDates[$index]) { |
| 344 | $customFormatErrorCount ++; |
| 345 | $customFormatInvalidRecords[] = $index; |
| 346 | } |
| 347 | } |
| 348 | |
| 349 | if ($customFormatErrorCount < $datetimeErrorCount) { |
| 350 | $invalidRecords = $customFormatInvalidRecords; |
| 351 | return $customFormatDates; |
| 352 | } |
| 353 | |
| 354 | $invalidRecords = $dateTimeInvalidRecords; |
| 355 | return $dateTimeDates; |
| 356 | } |
| 357 | |
| 358 | public function transformSubscribersData(array $subscribers, array $columns): array { |
| 359 | $transformedSubscribers = []; |
| 360 | foreach ($columns as $column => $data) { |
| 361 | $transformedSubscribers[$column] = array_column($subscribers, $data['index']); |
| 362 | } |
| 363 | return $transformedSubscribers; |
| 364 | } |
| 365 | |
| 366 | /** |
| 367 | * @param array $subscribersData |
| 368 | * @return array{array|false,array,array|false} |
| 369 | */ |
| 370 | public function splitSubscribersData(array $subscribersData): array { |
| 371 | // $subscribers_data is an two-dimensional associative array |
| 372 | // of all subscribers being imported: [field => [value1, value2], field => [value1, value2], ...] |
| 373 | $tempExistingSubscribers = []; |
| 374 | foreach (array_chunk($subscribersData['email'], self::DB_QUERY_CHUNK_SIZE) as $subscribersEmails) { |
| 375 | // create a two-dimensional indexed array of all existing subscribers |
| 376 | // with just wp_user_id and email fields: [[wp_user_id, email], [wp_user_id, email], ...] |
| 377 | $tempExistingSubscribers = array_merge( |
| 378 | $tempExistingSubscribers, |
| 379 | $this->subscriberRepository->findWpUserIdAndEmailByEmails($subscribersEmails) |
| 380 | ); |
| 381 | } |
| 382 | if (!$tempExistingSubscribers) { |
| 383 | return [ |
| 384 | false, // existing subscribers |
| 385 | $subscribersData, // new subscribers |
| 386 | false, // WP users |
| 387 | ]; |
| 388 | } |
| 389 | // extract WP users ids into a simple indexed array: [wp_user_id_1, wp_user_id_2, ...] |
| 390 | $wpUsers = array_filter(array_column($tempExistingSubscribers, 'wp_user_id')); |
| 391 | // create a new two-dimensional associative array with existing subscribers ($existing_subscribers) |
| 392 | // and reduce $subscribers_data to only new subscribers by removing existing subscribers |
| 393 | $existingSubscribers = []; |
| 394 | $subscribersEmails = array_flip($subscribersData['email']); |
| 395 | foreach ($tempExistingSubscribers as $tempExistingSubscriber) { |
| 396 | $existingSubscriberKey = $subscribersEmails[$tempExistingSubscriber['email']]; |
| 397 | foreach ($subscribersData as $field => &$value) { |
| 398 | $existingSubscribers[$field][] = $value[$existingSubscriberKey]; |
| 399 | unset($value[$existingSubscriberKey]); |
| 400 | } |
| 401 | } |
| 402 | $newSubscribers = $subscribersData; |
| 403 | // reindex array after unsetting elements |
| 404 | $newSubscribers = array_map('array_values', $newSubscribers); |
| 405 | // remove empty values |
| 406 | $newSubscribers = array_filter($newSubscribers); |
| 407 | return [ |
| 408 | $existingSubscribers, |
| 409 | $newSubscribers, |
| 410 | $wpUsers, |
| 411 | ]; |
| 412 | } |
| 413 | |
| 414 | public function deleteExistingTrashedSubscribers(array $subscribersData): void { |
| 415 | $existingTrashedRecords = array_filter( |
| 416 | array_map(function($subscriberEmails) { |
| 417 | return $this->subscriberRepository->findIdsOfDeletedByEmails($subscriberEmails); |
| 418 | }, array_chunk($subscribersData['email'], self::DB_QUERY_CHUNK_SIZE)) |
| 419 | ); |
| 420 | $existingTrashedRecords = Helpers::flattenArray($existingTrashedRecords); |
| 421 | if (!$existingTrashedRecords) { |
| 422 | return; |
| 423 | } |
| 424 | foreach (array_chunk($existingTrashedRecords, self::DB_QUERY_CHUNK_SIZE) as $subscriberIds) { |
| 425 | $this->subscriberRepository->bulkDelete($subscriberIds); |
| 426 | } |
| 427 | } |
| 428 | |
| 429 | public function addMissingRequiredFields(array $subscribers): array { |
| 430 | foreach (array_keys($this->requiredSubscribersFields) as $requiredField) { |
| 431 | $subscribers = $this->addField($subscribers, $requiredField, $this->requiredSubscribersFields[$requiredField]); |
| 432 | } |
| 433 | return $subscribers; |
| 434 | } |
| 435 | |
| 436 | /** |
| 437 | * @param array $subscribers |
| 438 | * @param string $fieldName |
| 439 | * @param mixed $fieldValue |
| 440 | * @return array |
| 441 | */ |
| 442 | private function addField(array $subscribers, string $fieldName, $fieldValue): array { |
| 443 | if (in_array($fieldName, $subscribers['fields'])) return $subscribers; |
| 444 | |
| 445 | $subscribersCount = count($subscribers['data'][key($subscribers['data'])]); |
| 446 | $subscribers['data'][$fieldName] = array_fill( |
| 447 | 0, |
| 448 | $subscribersCount, |
| 449 | $fieldValue |
| 450 | ); |
| 451 | $subscribers['fields'][] = $fieldName; |
| 452 | |
| 453 | return $subscribers; |
| 454 | } |
| 455 | |
| 456 | private function setSubscriptionStatusToDefault(array $subscribersData, string $defaultStatus): array { |
| 457 | if (!in_array('status', $subscribersData['fields'])) return $subscribersData; |
| 458 | $subscribersData['data']['status'] = array_map(function() use ($defaultStatus) { |
| 459 | return $defaultStatus; |
| 460 | }, $subscribersData['data']['status']); |
| 461 | |
| 462 | if ($defaultStatus === SubscriberEntity::STATUS_SUBSCRIBED) { |
| 463 | if (!in_array('last_subscribed_at', $subscribersData['fields'])) { |
| 464 | $subscribersData['fields'][] = 'last_subscribed_at'; |
| 465 | } |
| 466 | $subscribersData['data']['last_subscribed_at'] = array_map(function() { |
| 467 | return $this->createdAt; |
| 468 | }, $subscribersData['data']['status']); |
| 469 | } |
| 470 | return $subscribersData; |
| 471 | } |
| 472 | |
| 473 | private function setSource(array $subscribersData): array { |
| 474 | $subscribersCount = count($subscribersData['data'][key($subscribersData['data'])]); |
| 475 | $subscribersData['fields'][] = 'source'; |
| 476 | $subscribersData['data']['source'] = array_fill( |
| 477 | 0, |
| 478 | $subscribersCount, |
| 479 | Source::IMPORTED |
| 480 | ); |
| 481 | return $subscribersData; |
| 482 | } |
| 483 | |
| 484 | private function setLinkToken(array $subscribersData): array { |
| 485 | $subscribersCount = count($subscribersData['data'][key($subscribersData['data'])]); |
| 486 | $subscribersData['fields'][] = 'link_token'; |
| 487 | $subscribersData['data']['link_token'] = array_map( |
| 488 | function () { |
| 489 | return Security::generateRandomString(SubscriberEntity::LINK_TOKEN_LENGTH); |
| 490 | }, |
| 491 | array_fill(0, $subscribersCount, null) |
| 492 | ); |
| 493 | return $subscribersData; |
| 494 | } |
| 495 | |
| 496 | public function getSubscribersFields(array $subscribersFields): array { |
| 497 | return array_values( |
| 498 | array_filter( |
| 499 | array_map(function($field) { |
| 500 | if (!is_int($field)) return $field; |
| 501 | }, $subscribersFields) |
| 502 | ) |
| 503 | ); |
| 504 | } |
| 505 | |
| 506 | /** |
| 507 | * @param array $subscribersFields |
| 508 | * @return int[] |
| 509 | */ |
| 510 | public function getCustomSubscribersFields(array $subscribersFields): array { |
| 511 | return array_values( |
| 512 | array_filter( |
| 513 | array_map(function($field) { |
| 514 | if (is_int($field)) return $field; |
| 515 | }, $subscribersFields) |
| 516 | ) |
| 517 | ); |
| 518 | } |
| 519 | |
| 520 | public function createOrUpdateSubscribers( |
| 521 | string $action, |
| 522 | array $subscribersData, |
| 523 | array $subscribersCustomFields = [] |
| 524 | ): ?array { |
| 525 | $subscribersCount = count($subscribersData['data'][key($subscribersData['data'])]); |
| 526 | $subscribers = array_map(function($index) use ($subscribersData) { |
| 527 | return array_map(function($field) use ($index, $subscribersData) { |
| 528 | return $subscribersData['data'][$field][$index]; |
| 529 | }, $subscribersData['fields']); |
| 530 | }, range(0, $subscribersCount - 1)); |
| 531 | foreach (array_chunk($subscribers, self::DB_QUERY_CHUNK_SIZE) as $data) { |
| 532 | if ($action === self::ACTION_CREATE) { |
| 533 | $this->importExportRepository->insertMultiple( |
| 534 | SubscriberEntity::class, |
| 535 | $subscribersData['fields'], |
| 536 | $data |
| 537 | ); |
| 538 | } elseif ($action === self::ACTION_UPDATE) { |
| 539 | $this->importExportRepository->updateMultiple( |
| 540 | SubscriberEntity::class, |
| 541 | $subscribersData['fields'], |
| 542 | $data, |
| 543 | $this->updatedAt |
| 544 | ); |
| 545 | } |
| 546 | } |
| 547 | $createdOrUpdatedSubscribers = []; |
| 548 | foreach (array_chunk($subscribersData['data']['email'], self::DB_QUERY_CHUNK_SIZE) as $emails) { |
| 549 | foreach ($this->subscriberRepository->findIdAndEmailByEmails($emails) as $createdOrUpdatedSubscriber) { |
| 550 | // ensure emails loaded from the DB are lowercased (imported emails are lowercased as well) |
| 551 | $createdOrUpdatedSubscriber['email'] = mb_strtolower($createdOrUpdatedSubscriber['email']); |
| 552 | $createdOrUpdatedSubscribers[] = $createdOrUpdatedSubscriber; |
| 553 | } |
| 554 | } |
| 555 | if (empty($createdOrUpdatedSubscribers)) return null; |
| 556 | |
| 557 | $this->subscriberRepository->invalidateTotalSubscribersCache(); |
| 558 | $createdOrUpdatedSubscribersIds = array_column($createdOrUpdatedSubscribers, 'id'); |
| 559 | if ($subscribersCustomFields) { |
| 560 | $this->createOrUpdateCustomFields( |
| 561 | $action, |
| 562 | $createdOrUpdatedSubscribers, |
| 563 | $subscribersData, |
| 564 | $subscribersCustomFields |
| 565 | ); |
| 566 | } |
| 567 | $this->addSubscribersToSegments( |
| 568 | $createdOrUpdatedSubscribersIds, |
| 569 | $this->segmentsIds |
| 570 | ); |
| 571 | $this->addTagsToSubscribers( |
| 572 | $createdOrUpdatedSubscribersIds, |
| 573 | $this->tags |
| 574 | ); |
| 575 | return $createdOrUpdatedSubscribers; |
| 576 | } |
| 577 | |
| 578 | public function createOrUpdateCustomFields( |
| 579 | string $action, |
| 580 | array $createdOrUpdatedSubscribers, |
| 581 | array $subscribersData, |
| 582 | array $subscribersCustomFieldsIds |
| 583 | ): void { |
| 584 | // check if custom fields exist in the database |
| 585 | $subscribersCustomFieldsIds = array_map(function(CustomFieldEntity $customField): int { |
| 586 | return (int)$customField->getId(); |
| 587 | }, $this->customFieldsRepository->findBy(['id' => $subscribersCustomFieldsIds, 'deletedAt' => null])); |
| 588 | if (!$subscribersCustomFieldsIds) { |
| 589 | return; |
| 590 | } |
| 591 | // assemble a two-dimensional array: [[custom_field_id, subscriber_id, value], [custom_field_id, subscriber_id, value], ...] |
| 592 | $subscribersCustomFieldsData = []; |
| 593 | $subscribersEmails = array_flip($subscribersData['data']['email']); |
| 594 | foreach ($createdOrUpdatedSubscribers as $createdOrUpdatedSubscriber) { |
| 595 | $subscriberIndex = $subscribersEmails[$createdOrUpdatedSubscriber['email']]; |
| 596 | foreach ($subscribersData['data'] as $field => $values) { |
| 597 | // exclude non-custom fields |
| 598 | if (!is_int($field)) continue; |
| 599 | $subscribersCustomFieldsData[] = [ |
| 600 | (int)$field, |
| 601 | $createdOrUpdatedSubscriber['id'], |
| 602 | $values[$subscriberIndex], |
| 603 | $this->createdAt, |
| 604 | ]; |
| 605 | } |
| 606 | } |
| 607 | $columns = [ |
| 608 | 'custom_field_id', |
| 609 | 'subscriber_id', |
| 610 | 'value', |
| 611 | 'created_at', |
| 612 | ]; |
| 613 | $customFieldCount = count($subscribersCustomFieldsIds); |
| 614 | $customFieldBatchSize = (int)(round(self::DB_QUERY_CHUNK_SIZE / $customFieldCount) * $customFieldCount); |
| 615 | $customFieldBatchSize = ($customFieldBatchSize > 0) ? $customFieldBatchSize : 1; |
| 616 | foreach (array_chunk($subscribersCustomFieldsData, $customFieldBatchSize) as $subscribersCustomFieldsDataChunk) { |
| 617 | $this->importExportRepository->insertMultiple( |
| 618 | SubscriberCustomFieldEntity::class, |
| 619 | $columns, |
| 620 | $subscribersCustomFieldsDataChunk |
| 621 | ); |
| 622 | if ($action === self::ACTION_UPDATE) { |
| 623 | $this->importExportRepository->updateMultiple( |
| 624 | SubscriberCustomFieldEntity::class, |
| 625 | $columns, |
| 626 | $subscribersCustomFieldsDataChunk, |
| 627 | $this->updatedAt |
| 628 | ); |
| 629 | } |
| 630 | } |
| 631 | } |
| 632 | |
| 633 | /** |
| 634 | * @param int[] $wpUsers |
| 635 | * @return array |
| 636 | */ |
| 637 | public function synchronizeWPUsers(array $wpUsers): array { |
| 638 | $users = array_map([$this->wpSegment, 'synchronizeUser'], $wpUsers); |
| 639 | $this->subscriberRepository->invalidateTotalSubscribersCache(); |
| 640 | return $users; |
| 641 | } |
| 642 | |
| 643 | public function addSubscribersToSegments(array $subscribersIds, array $segmentsIds): void { |
| 644 | $columns = [ |
| 645 | 'subscriber_id', |
| 646 | 'segment_id', |
| 647 | 'created_at', |
| 648 | ]; |
| 649 | foreach ($segmentsIds as $segmentId) { |
| 650 | foreach (array_chunk($subscribersIds, self::DB_QUERY_CHUNK_SIZE) as $subscriberIdsChunk) { |
| 651 | $data = []; |
| 652 | $data = array_merge($data, array_map(function ($subscriberId) use ($segmentId): array { |
| 653 | return [ |
| 654 | $subscriberId, |
| 655 | $segmentId, |
| 656 | $this->createdAt, |
| 657 | ]; |
| 658 | }, $subscriberIdsChunk)); |
| 659 | |
| 660 | $this->importExportRepository->insertMultiple( |
| 661 | SubscriberSegmentEntity::class, |
| 662 | $columns, |
| 663 | $data |
| 664 | ); |
| 665 | } |
| 666 | } |
| 667 | $this->subscriberRepository->recalculateSegmentsCount($subscribersIds); |
| 668 | } |
| 669 | |
| 670 | /** |
| 671 | * @param int[] $subscribersIds |
| 672 | * @param string[] $tagNames |
| 673 | */ |
| 674 | public function addTagsToSubscribers(array $subscribersIds, array $tagNames): void { |
| 675 | $tagIds = []; |
| 676 | foreach ($tagNames as $tagName) { |
| 677 | $tag = $this->tagRepository->findOneBy(['name' => $tagName]); |
| 678 | if (!$tag) { |
| 679 | $tag = $this->tagRepository->createOrUpdate(['name' => $tagName]); |
| 680 | } |
| 681 | $tagIds[] = $tag->getId(); |
| 682 | } |
| 683 | |
| 684 | $columns = [ |
| 685 | 'subscriber_id', |
| 686 | 'tag_id', |
| 687 | 'created_at', |
| 688 | ]; |
| 689 | foreach ($tagIds as $tagId) { |
| 690 | foreach (array_chunk($subscribersIds, self::DB_QUERY_CHUNK_SIZE) as $subscriberIdsChunk) { |
| 691 | $data = []; |
| 692 | $data = array_merge($data, array_map(function ($subscriberId) use ($tagId): array { |
| 693 | return [ |
| 694 | $subscriberId, |
| 695 | $tagId, |
| 696 | $this->createdAt, |
| 697 | ]; |
| 698 | }, $subscriberIdsChunk)); |
| 699 | |
| 700 | $this->importExportRepository->insertMultiple( |
| 701 | SubscriberTagEntity::class, |
| 702 | $columns, |
| 703 | $data |
| 704 | ); |
| 705 | } |
| 706 | } |
| 707 | } |
| 708 | } |
| 709 |