← All changes
|
jetpack_vendor/automattic/jetpack-sync/src/class-sender.php
+69
-33
12.1.3
→
16.3-beta
View file →
| @@ -24,8 +24,26 @@ | ||
| 24 | 24 | */ |
| 25 | 25 | const NEXT_SYNC_TIME_OPTION_NAME = 'jetpack_next_sync_time'; |
| 26 | 26 | |
| 27 | 27 | /** |
| 28 | + * Name of the transient responsible for temprorarily disabling Sync sending during Pulls. | |
| 29 | + * | |
| 30 | + * @access public | |
| 31 | + * | |
| 32 | + * @var string | |
| 33 | + */ | |
| 34 | + const TEMP_SYNC_DISABLE_TRANSIENT_NAME = 'jetpack_disable_sync_sending'; | |
| 35 | + | |
| 36 | + /** | |
| 37 | + * Expiry of the transient responsible for temprorarily disabling Sync sending during Pulls. | |
| 38 | + * | |
| 39 | + * @access public | |
| 40 | + * | |
| 41 | + * @var int | |
| 42 | + */ | |
| 43 | + const TEMP_SYNC_DISABLE_TRANSIENT_EXPIRY = MINUTE_IN_SECONDS; | |
| 44 | + | |
| 45 | + /** | |
| 28 | 46 | * Sync timeout after a WPCOM error. |
| 29 | 47 | * |
| 30 | 48 | * @access public |
| 31 | 49 | * |
| @@ -109,9 +127,9 @@ | ||
| 109 | 127 | * Incremental sync queue object. |
| 110 | 128 | * |
| 111 | 129 | * @access private |
| 112 | 130 | * |
| 113 | - * @var Automattic\Jetpack\Sync\Queue | |
| 131 | + * @var \Automattic\Jetpack\Sync\Queue | |
| 114 | 132 | */ |
| 115 | 133 | private $sync_queue; |
| 116 | 134 | |
| 117 | 135 | /** |
| @@ -118,9 +136,9 @@ | ||
| 118 | 136 | * Full sync queue object. |
| 119 | 137 | * |
| 120 | 138 | * @access private |
| 121 | 139 | * |
| 122 | - * @var Automattic\Jetpack\Sync\Queue | |
| 140 | + * @var \Automattic\Jetpack\Sync\Queue | |
| 123 | 141 | */ |
| 124 | 142 | private $full_sync_queue; |
| 125 | 143 | |
| 126 | 144 | /** |
| @@ -127,9 +145,9 @@ | ||
| 127 | 145 | * Codec object for encoding and decoding sync items. |
| 128 | 146 | * |
| 129 | 147 | * @access private |
| 130 | 148 | * |
| 131 | - * @var Automattic\Jetpack\Sync\Codec_Interface | |
| 149 | + * @var \Automattic\Jetpack\Sync\Codec_Interface | |
| 132 | 150 | */ |
| 133 | 151 | private $codec; |
| 134 | 152 | |
| 135 | 153 | /** |
| @@ -146,9 +164,9 @@ | ||
| 146 | 164 | * |
| 147 | 165 | * @access private |
| 148 | 166 | * @static |
| 149 | 167 | * |
| 150 | - * @var Automattic\Jetpack\Sync\Sender | |
| 168 | + * @var \Automattic\Jetpack\Sync\Sender | |
| 151 | 169 | */ |
| 152 | 170 | private static $instance; |
| 153 | 171 | |
| 154 | 172 | /** |
| @@ -207,9 +225,9 @@ | ||
| 207 | 225 | ! is_wp_error( $verified_user ) |
| 208 | 226 | && $verified_user |
| 209 | 227 | ) { |
| 210 | 228 | $old_user = wp_get_current_user(); |
| 211 | - $this->old_user = isset( $old_user->ID ) ? $old_user->ID : 0; | |
| 229 | + $this->old_user = $old_user->ID ?? 0; | |
| 212 | 230 | wp_set_current_user( $verified_user['user_id'] ); |
| 213 | 231 | } |
| 214 | 232 | } |
| 215 | 233 | |
| @@ -273,8 +291,9 @@ | ||
| 273 | 291 | * @return boolean|WP_Error True if this sync sending was successful, error object otherwise. |
| 274 | 292 | */ |
| 275 | 293 | public function do_full_sync() { |
| 276 | 294 | $sync_module = Modules::get_module( 'full-sync' ); |
| 295 | + '@phan-var Modules\Full_Sync_Immediately|Modules\Full_Sync $sync_module'; | |
| 277 | 296 | if ( ! $sync_module ) { |
| 278 | 297 | return; |
| 279 | 298 | } |
| 280 | 299 | // Full Sync Disabled. |
| @@ -294,9 +313,9 @@ | ||
| 294 | 313 | } |
| 295 | 314 | |
| 296 | 315 | $this->continue_full_sync_enqueue(); |
| 297 | 316 | // immediate full sync sends data in continue_full_sync_enqueue. |
| 298 | - if ( false === strpos( get_class( $sync_module ), 'Full_Sync_Immediately' ) ) { | |
| 317 | + if ( ! $sync_module instanceof Modules\Full_Sync_Immediately ) { | |
| 299 | 318 | return $this->do_sync_and_set_delays( $this->full_sync_queue ); |
| 300 | 319 | } else { |
| 301 | 320 | $status = $sync_module->get_status(); |
| 302 | 321 | // Sync not started or Sync finished. |
| @@ -323,9 +342,11 @@ | ||
| 323 | 342 | if ( $this->get_next_sync_time( 'full-sync-enqueue' ) > microtime( true ) ) { |
| 324 | 343 | return false; |
| 325 | 344 | } |
| 326 | 345 | |
| 327 | - Modules::get_module( 'full-sync' )->continue_enqueuing(); | |
| 346 | + $full_sync_module = Modules::get_module( 'full-sync' ); | |
| 347 | + '@phan-var Modules\Full_Sync_Immediately|Modules\Full_Sync $full_sync_module'; | |
| 348 | + $full_sync_module->continue_enqueuing(); | |
| 328 | 349 | |
| 329 | 350 | $this->set_next_sync_time( time() + $this->get_enqueue_wait_time(), 'full-sync-enqueue' ); |
| 330 | 351 | } |
| 331 | 352 | |
| @@ -336,9 +357,12 @@ | ||
| 336 | 357 | * |
| 337 | 358 | * @return boolean|WP_Error True if this sync sending was successful, error object otherwise. |
| 338 | 359 | */ |
| 339 | 360 | public function do_sync() { |
| 340 | - if ( ! Settings::is_dedicated_sync_enabled() ) { | |
| 361 | + // Sync directly during cron. We are doing this because otherwise | |
| 362 | + // the dedicated sync flow would be spawning HTTP requests during cron shutdown, | |
| 363 | + // which can be unreliable and cause sync lag for time-sensitive events like updates. | |
| 364 | + if ( ! Settings::is_dedicated_sync_enabled() || Settings::is_doing_cron() ) { | |
| 341 | 365 | $result = $this->do_sync_and_set_delays( $this->sync_queue ); |
| 342 | 366 | } else { |
| 343 | 367 | $result = Dedicated_Sender::spawn_sync( $this->sync_queue ); |
| 344 | 368 | } |
| @@ -372,9 +396,9 @@ | ||
| 372 | 396 | * This is used to test the feature is working. |
| 373 | 397 | * |
| 374 | 398 | * @see \Automattic\Jetpack\Sync\Dedicated_Sender::can_spawn_dedicated_sync_request |
| 375 | 399 | */ |
| 376 | - // phpcs:ignore WordPress.Security.EscapeOutput.OutputNotEscaped | |
| 400 | + // phpcs:ignore WordPress.Security.EscapeOutput.OutputNotEscaped -- This is just a constant string used for Validation. | |
| 377 | 401 | echo Dedicated_Sender::DEDICATED_SYNC_VALIDATION_STRING; |
| 378 | 402 | |
| 379 | 403 | // Try to disconnect the request as quickly as possible and process things in the background. |
| 380 | 404 | $this->fastcgi_finish_request(); |
| @@ -396,14 +420,14 @@ | ||
| 396 | 420 | if ( session_status() === PHP_SESSION_ACTIVE ) { |
| 397 | 421 | session_write_close(); |
| 398 | 422 | } |
| 399 | 423 | |
| 424 | + // Actually try to send Sync events. | |
| 425 | + $result = $this->do_sync_and_set_delays( $this->sync_queue ); | |
| 426 | + | |
| 400 | 427 | // Output not used right now. Try to release dedicated sync lock |
| 401 | 428 | Dedicated_Sender::try_release_lock_spawn_request(); |
| 402 | 429 | |
| 403 | - // Actually try to send Sync events. | |
| 404 | - $result = $this->do_sync_and_set_delays( $this->sync_queue ); | |
| 405 | - | |
| 406 | 430 | // If no errors occurred, re-spawn a dedicated Sync request. |
| 407 | 431 | if ( true === $result ) { |
| 408 | 432 | Dedicated_Sender::spawn_sync( $this->sync_queue ); |
| 409 | 433 | } |
| @@ -408,9 +432,9 @@ | ||
| 408 | 432 | Dedicated_Sender::spawn_sync( $this->sync_queue ); |
| 409 | 433 | } |
| 410 | 434 | |
| 411 | 435 | if ( $do_real_exit ) { |
| 412 | - exit; | |
| 436 | + exit( 0 ); | |
| 413 | 437 | } |
| 414 | 438 | } |
| 415 | 439 | |
| 416 | 440 | /** |
| @@ -420,9 +444,9 @@ | ||
| 420 | 444 | * Will be delayed until the next sync time comes. |
| 421 | 445 | * |
| 422 | 446 | * @access public |
| 423 | 447 | * |
| 424 | - * @param Automattic\Jetpack\Sync\Queue $queue Queue object. | |
| 448 | + * @param \Automattic\Jetpack\Sync\Queue $queue Queue object. | |
| 425 | 449 | * |
| 426 | 450 | * @return boolean|WP_Error True if this sync sending was successful, error object otherwise. |
| 427 | 451 | */ |
| 428 | 452 | public function do_sync_and_set_delays( $queue ) { |
| @@ -439,8 +463,12 @@ | ||
| 439 | 463 | if ( ! Settings::is_sender_enabled( $queue->id ) ) { |
| 440 | 464 | return new WP_Error( 'sender_disabled_for_queue_' . $queue->id ); |
| 441 | 465 | } |
| 442 | 466 | |
| 467 | + if ( get_transient( self::TEMP_SYNC_DISABLE_TRANSIENT_NAME ) ) { | |
| 468 | + return new WP_Error( 'sender_temporarily_disabled_while_pulling' ); | |
| 469 | + } | |
| 470 | + | |
| 443 | 471 | // Return early if we've gotten a retry-after header response. |
| 444 | 472 | $retry_time = get_option( Actions::RETRY_AFTER_PREFIX . $queue->id ); |
| 445 | 473 | if ( $retry_time ) { |
| 446 | 474 | // If expired update to false but don't send. Send will occurr in new request to avoid race conditions. |
| @@ -471,10 +499,12 @@ | ||
| 471 | 499 | } |
| 472 | 500 | if ( 'wpcom_error' === $sync_result->get_error_code() ) { |
| 473 | 501 | $this->set_next_sync_time( time() + self::WPCOM_ERROR_SYNC_DELAY, $queue->id ); |
| 474 | 502 | } |
| 475 | - } elseif ( $exceeded_sync_wait_threshold ) { | |
| 476 | - // If we actually sent data and it took a while, wait before sending again. | |
| 503 | + } elseif ( $exceeded_sync_wait_threshold && ! Settings::is_doing_cron() ) { | |
| 504 | + // If a send was slow, briefly pause before the next one. | |
| 505 | + // Applies only to Dedicated/Normal Sync to avoid impacting user traffic; | |
| 506 | + // cron jobs are exempt. | |
| 477 | 507 | $this->set_next_sync_time( time() + $this->get_sync_wait_time(), $queue->id ); |
| 478 | 508 | } |
| 479 | 509 | |
| 480 | 510 | return $sync_result; |
| @@ -484,10 +514,10 @@ | ||
| 484 | 514 | * Retrieve the next sync items to send. |
| 485 | 515 | * |
| 486 | 516 | * @access public |
| 487 | 517 | * |
| 488 | - * @param (array|Automattic\Jetpack\Sync\Queue_Buffer) $buffer_or_items Queue buffer or array of objects. | |
| 489 | - * @param boolean $encode Whether to encode the items. | |
| 518 | + * @param (array|\Automattic\Jetpack\Sync\Queue_Buffer) $buffer_or_items Queue buffer or array of objects. | |
| 519 | + * @param boolean $encode Whether to encode the items. | |
| 490 | 520 | * @return array Sync items to send. |
| 491 | 521 | */ |
| 492 | 522 | public function get_items_to_send( $buffer_or_items, $encode = true ) { |
| 493 | 523 | // Track how long we've been processing so we can avoid request timeouts. |
| @@ -508,8 +538,13 @@ | ||
| 508 | 538 | * We estimate the total encoded size as we go by encoding each item individually. |
| 509 | 539 | * This is expensive, but the only way to really know :/ |
| 510 | 540 | */ |
| 511 | 541 | foreach ( $items as $key => $item ) { |
| 542 | + if ( ! is_array( $item ) ) { | |
| 543 | + $skipped_items_ids[] = $key; | |
| 544 | + continue; | |
| 545 | + } | |
| 546 | + | |
| 512 | 547 | // Suspending cache addition help prevent overloading in memory cache of large sites. |
| 513 | 548 | wp_suspend_cache_addition( true ); |
| 514 | 549 | /** |
| 515 | 550 | * Modify the data within an action before it is serialized and sent to the server |
| @@ -530,9 +565,9 @@ | ||
| 530 | 565 | continue; |
| 531 | 566 | } |
| 532 | 567 | $encoded_item = $this->codec->encode( $item ); |
| 533 | 568 | $upload_size += strlen( $encoded_item ); |
| 534 | - if ( $upload_size > $this->upload_max_bytes && count( $items_to_send ) > 0 ) { | |
| 569 | + if ( $upload_size > $this->upload_max_bytes && array() !== $items_to_send ) { | |
| 535 | 570 | break; |
| 536 | 571 | } |
| 537 | 572 | $items_to_send[ $key ] = $encode ? $encoded_item : $item; |
| 538 | 573 | if ( microtime( true ) - $start_time > $this->max_dequeue_time ) { |
| @@ -549,9 +584,9 @@ | ||
| 549 | 584 | * |
| 550 | 585 | * @access private |
| 551 | 586 | */ |
| 552 | 587 | private function fastcgi_finish_request() { |
| 553 | - if ( function_exists( 'fastcgi_finish_request' ) && version_compare( phpversion(), '7.0.16', '>=' ) ) { | |
| 588 | + if ( function_exists( 'fastcgi_finish_request' ) ) { | |
| 554 | 589 | fastcgi_finish_request(); |
| 555 | 590 | } |
| 556 | 591 | } |
| 557 | 592 | |
| @@ -559,9 +594,9 @@ | ||
| 559 | 594 | * Perform sync for a certain sync queue. |
| 560 | 595 | * |
| 561 | 596 | * @access public |
| 562 | 597 | * |
| 563 | - * @param Automattic\Jetpack\Sync\Queue $queue Queue object. | |
| 598 | + * @param \Automattic\Jetpack\Sync\Queue $queue Queue object. | |
| 564 | 599 | * |
| 565 | 600 | * @return boolean|WP_Error True if this sync sending was successful, error object otherwise. |
| 566 | 601 | */ |
| 567 | 602 | public function do_sync_for_queue( $queue ) { |
| @@ -573,8 +608,9 @@ | ||
| 573 | 608 | /** |
| 574 | 609 | * Now that we're sure we are about to sync, try to ignore user abort |
| 575 | 610 | * so we can avoid getting into a bad state. |
| 576 | 611 | */ |
| 612 | + // https://plugins.trac.wordpress.org/ticket/2041 | |
| 577 | 613 | if ( function_exists( 'ignore_user_abort' ) ) { |
| 578 | 614 | ignore_user_abort( true ); |
| 579 | 615 | } |
| 580 | 616 | |
| @@ -640,13 +676,11 @@ | ||
| 640 | 676 | return new WP_Error( 'wpcom_error', 'jetpack_sync_send_data_false' ); |
| 641 | 677 | } else { |
| 642 | 678 | // Detect if the last item ID was an error. |
| 643 | 679 | $had_wp_error = is_wp_error( end( $processed_item_ids ) ); |
| 644 | - if ( $had_wp_error ) { | |
| 645 | - $wp_error = array_pop( $processed_item_ids ); | |
| 646 | - } | |
| 680 | + $wp_error = $had_wp_error ? array_pop( $processed_item_ids ) : null; | |
| 647 | 681 | // Also checkin any items that were skipped. |
| 648 | - if ( count( $skipped_items_ids ) > 0 ) { | |
| 682 | + if ( array() !== $skipped_items_ids ) { | |
| 649 | 683 | $processed_item_ids = array_merge( $processed_item_ids, $skipped_items_ids ); |
| 650 | 684 | } |
| 651 | 685 | $processed_items = array_intersect_key( $items, array_flip( $processed_item_ids ) ); |
| 652 | 686 | /** |
| @@ -674,18 +708,19 @@ | ||
| 674 | 708 | * Immediately sends a single item without firing or enqueuing it |
| 675 | 709 | * |
| 676 | 710 | * @param string $action_name The action. |
| 677 | 711 | * @param array $data The data associated with the action. |
| 712 | + * @param string $key The key to use for the action. | |
| 678 | 713 | * |
| 679 | - * @return Items processed. TODO: this doesn't make much sense anymore, it should probably be just a bool. | |
| 714 | + * @return array Items processed. TODO: this doesn't make much sense anymore, it should probably be just a bool. | |
| 680 | 715 | */ |
| 681 | - public function send_action( $action_name, $data = null ) { | |
| 716 | + public function send_action( $action_name, $data = null, $key = null ) { | |
| 682 | 717 | if ( ! Settings::is_sender_enabled( 'full_sync' ) ) { |
| 683 | 718 | return array(); |
| 684 | 719 | } |
| 685 | 720 | |
| 686 | 721 | // Compose the data to be sent. |
| 687 | - $action_to_send = $this->create_action_to_send( $action_name, $data ); | |
| 722 | + $action_to_send = $this->create_action_to_send( $action_name, $data, $key ); | |
| 688 | 723 | |
| 689 | 724 | list( $items_to_send, $skipped_items_ids, $items, $preprocess_duration ) = $this->get_items_to_send( $action_to_send, true ); // phpcs:ignore VariableAnalysis.CodeAnalysis.VariableAnalysis.UnusedVariable |
| 690 | 725 | Settings::set_is_sending( true ); |
| 691 | 726 | $processed_item_ids = apply_filters( 'jetpack_sync_send_data', $items_to_send, $this->get_codec()->name(), microtime( true ), 'immediate-send', 0, $preprocess_duration ); |
| @@ -711,13 +746,14 @@ | ||
| 711 | 746 | * @access private |
| 712 | 747 | * |
| 713 | 748 | * @param string $action_name The action. |
| 714 | 749 | * @param array $data The data associated with the action. |
| 750 | + * @param string $key The key to use for the action. | |
| 715 | 751 | * @return array An array of synthetic sync actions keyed by current microtime(true) |
| 716 | 752 | */ |
| 717 | - private function create_action_to_send( $action_name, $data ) { | |
| 753 | + private function create_action_to_send( $action_name, $data, $key = null ) { | |
| 718 | 754 | return array( |
| 719 | - (string) microtime( true ) => array( | |
| 755 | + $key ?? (string) microtime( true ) => array( | |
| 720 | 756 | $action_name, |
| 721 | 757 | $data, |
| 722 | 758 | get_current_user_id(), |
| 723 | 759 | microtime( true ), |
| @@ -763,9 +799,9 @@ | ||
| 763 | 799 | * Get the incremental sync queue object. |
| 764 | 800 | * |
| 765 | 801 | * @access public |
| 766 | 802 | * |
| 767 | - * @return Automattic\Jetpack\Sync\Queue Queue object. | |
| 803 | + * @return \Automattic\Jetpack\Sync\Queue Queue object. | |
| 768 | 804 | */ |
| 769 | 805 | public function get_sync_queue() { |
| 770 | 806 | return $this->sync_queue; |
| 771 | 807 | } |
| @@ -774,9 +810,9 @@ | ||
| 774 | 810 | * Get the full sync queue object. |
| 775 | 811 | * |
| 776 | 812 | * @access public |
| 777 | 813 | * |
| 778 | - * @return Automattic\Jetpack\Sync\Queue Queue object. | |
| 814 | + * @return \Automattic\Jetpack\Sync\Queue Queue object. | |
| 779 | 815 | */ |
| 780 | 816 | public function get_full_sync_queue() { |
| 781 | 817 | return $this->full_sync_queue; |
| 782 | 818 | } |
| @@ -785,9 +821,9 @@ | ||
| 785 | 821 | * Get the codec object. |
| 786 | 822 | * |
| 787 | 823 | * @access public |
| 788 | 824 | * |
| 789 | - * @return Automattic\Jetpack\Sync\Codec_Interface Codec object. | |
| 825 | + * @return \Automattic\Jetpack\Sync\Codec_Interface Codec object. | |
| 790 | 826 | */ |
| 791 | 827 | public function get_codec() { |
| 792 | 828 | return $this->codec; |
| 793 | 829 | } |