| 1 |
<?php |
| 2 |
|
| 3 |
require_once dirname( __FILE__ ) . '/class.jetpack-sync-queue.php'; |
| 4 |
require_once dirname( __FILE__ ) . '/class.jetpack-sync-defaults.php'; |
| 5 |
require_once dirname( __FILE__ ) . '/class.jetpack-sync-json-deflate-array-codec.php'; |
| 6 |
require_once dirname( __FILE__ ) . '/class.jetpack-sync-simple-codec.php'; |
| 7 |
require_once dirname( __FILE__ ) . '/class.jetpack-sync-modules.php'; |
| 8 |
require_once dirname( __FILE__ ) . '/class.jetpack-sync-settings.php'; |
| 9 |
|
| 10 |
/** |
| 11 |
* This class grabs pending actions from the queue and sends them |
| 12 |
*/ |
| 13 |
class Jetpack_Sync_Sender { |
| 14 |
|
| 15 |
const NEXT_SYNC_TIME_OPTION_NAME = 'jetpack_next_sync_time'; |
| 16 |
const WPCOM_ERROR_SYNC_DELAY = 60; |
| 17 |
const QUEUE_LOCKED_SYNC_DELAY = 10; |
| 18 |
|
| 19 |
private $dequeue_max_bytes; |
| 20 |
private $upload_max_bytes; |
| 21 |
private $upload_max_rows; |
| 22 |
private $max_dequeue_time; |
| 23 |
private $sync_wait_time; |
| 24 |
private $sync_wait_threshold; |
| 25 |
private $enqueue_wait_time; |
| 26 |
private $sync_queue; |
| 27 |
private $full_sync_queue; |
| 28 |
private $codec; |
| 29 |
private $old_user; |
| 30 |
|
| 31 |
// singleton functions |
| 32 |
private static $instance; |
| 33 |
|
| 34 |
public static function get_instance() { |
| 35 |
if ( null === self::$instance ) { |
| 36 |
self::$instance = new self(); |
| 37 |
} |
| 38 |
|
| 39 |
return self::$instance; |
| 40 |
} |
| 41 |
|
| 42 |
// this is necessary because you can't use "new" when you declare instance properties >:( |
| 43 |
protected function __construct() { |
| 44 |
$this->set_defaults(); |
| 45 |
$this->init(); |
| 46 |
} |
| 47 |
|
| 48 |
private function init() { |
| 49 |
add_action( 'jetpack_sync_before_send_queue_sync', array( $this, 'maybe_set_user_from_token' ), 1 ); |
| 50 |
add_action( 'jetpack_sync_before_send_queue_sync', array( $this, 'maybe_clear_user_from_token' ), 20 ); |
| 51 |
foreach ( Jetpack_Sync_Modules::get_modules() as $module ) { |
| 52 |
$module->init_before_send(); |
| 53 |
} |
| 54 |
} |
| 55 |
|
| 56 |
public function maybe_set_user_from_token( ) { |
| 57 |
$jetpack = Jetpack::init(); |
| 58 |
$verified_user = $jetpack->verify_xml_rpc_signature(); |
| 59 |
if ( Jetpack_Constants::is_true( 'XMLRPC_REQUEST' ) && |
| 60 |
! is_wp_error( $verified_user ) |
| 61 |
&& $verified_user |
| 62 |
) { |
| 63 |
$old_user = wp_get_current_user(); |
| 64 |
$this->old_user = isset( $old_user->ID ) ? $old_user->ID : 0; |
| 65 |
wp_set_current_user( $verified_user['user_id'] ); |
| 66 |
} |
| 67 |
} |
| 68 |
|
| 69 |
public function maybe_clear_user_from_token() { |
| 70 |
if ( isset( $this->old_user ) ) { |
| 71 |
wp_set_current_user( $this->old_user ); |
| 72 |
} |
| 73 |
} |
| 74 |
|
| 75 |
public function get_next_sync_time( $queue_name ) { |
| 76 |
return (double) get_option( self::NEXT_SYNC_TIME_OPTION_NAME . '_' . $queue_name, 0 ); |
| 77 |
} |
| 78 |
|
| 79 |
public function set_next_sync_time( $time, $queue_name ) { |
| 80 |
return update_option( self::NEXT_SYNC_TIME_OPTION_NAME . '_' . $queue_name, $time, true ); |
| 81 |
} |
| 82 |
|
| 83 |
public function do_full_sync() { |
| 84 |
if ( ! Jetpack_Sync_Modules::get_module( 'full-sync' ) ) { |
| 85 |
return; |
| 86 |
} |
| 87 |
$this->continue_full_sync_enqueue(); |
| 88 |
return $this->do_sync_and_set_delays( $this->full_sync_queue ); |
| 89 |
} |
| 90 |
|
| 91 |
private function continue_full_sync_enqueue() { |
| 92 |
if ( defined( 'WP_IMPORTING' ) && WP_IMPORTING ) { |
| 93 |
return false; |
| 94 |
} |
| 95 |
|
| 96 |
if ( $this->get_next_sync_time( 'full-sync-enqueue' ) > microtime( true ) ) { |
| 97 |
return false; |
| 98 |
} |
| 99 |
|
| 100 |
Jetpack_Sync_Modules::get_module( 'full-sync' )->continue_enqueuing(); |
| 101 |
|
| 102 |
$this->set_next_sync_time( time() + $this->get_enqueue_wait_time(), 'full-sync-enqueue' ); |
| 103 |
} |
| 104 |
|
| 105 |
public function do_sync() { |
| 106 |
return $this->do_sync_and_set_delays( $this->sync_queue ); |
| 107 |
} |
| 108 |
|
| 109 |
public function do_sync_and_set_delays( $queue ) { |
| 110 |
// don't sync if importing |
| 111 |
if ( defined( 'WP_IMPORTING' ) && WP_IMPORTING ) { |
| 112 |
return new WP_Error( 'is_importing' ); |
| 113 |
} |
| 114 |
|
| 115 |
// don't sync if we are throttled |
| 116 |
if ( $this->get_next_sync_time( $queue->id ) > microtime( true ) ) { |
| 117 |
return new WP_Error( 'sync_throttled' ); |
| 118 |
} |
| 119 |
|
| 120 |
$start_time = microtime( true ); |
| 121 |
|
| 122 |
Jetpack_Sync_Settings::set_is_syncing( true ); |
| 123 |
|
| 124 |
$sync_result = $this->do_sync_for_queue( $queue ); |
| 125 |
|
| 126 |
Jetpack_Sync_Settings::set_is_syncing( false ); |
| 127 |
|
| 128 |
$exceeded_sync_wait_threshold = ( microtime( true ) - $start_time ) > (double) $this->get_sync_wait_threshold(); |
| 129 |
|
| 130 |
if ( is_wp_error( $sync_result ) ) { |
| 131 |
if ( 'unclosed_buffer' === $sync_result->get_error_code() ) { |
| 132 |
$this->set_next_sync_time( time() + self::QUEUE_LOCKED_SYNC_DELAY, $queue->id ); |
| 133 |
} |
| 134 |
if ( 'wpcom_error' === $sync_result->get_error_code() ) { |
| 135 |
$this->set_next_sync_time( time() + self::WPCOM_ERROR_SYNC_DELAY, $queue->id ); |
| 136 |
} |
| 137 |
} elseif ( $exceeded_sync_wait_threshold ) { |
| 138 |
// if we actually sent data and it took a while, wait before sending again |
| 139 |
$this->set_next_sync_time( time() + $this->get_sync_wait_time(), $queue->id ); |
| 140 |
} |
| 141 |
|
| 142 |
return $sync_result; |
| 143 |
} |
| 144 |
|
| 145 |
public function get_items_to_send( $buffer, $encode = true ) { |
| 146 |
// track how long we've been processing so we can avoid request timeouts |
| 147 |
$start_time = microtime( true ); |
| 148 |
$upload_size = 0; |
| 149 |
$items_to_send = array(); |
| 150 |
$items = $buffer->get_items(); |
| 151 |
// set up current screen to avoid errors rendering content |
| 152 |
require_once( ABSPATH . 'wp-admin/includes/class-wp-screen.php' ); |
| 153 |
require_once( ABSPATH . 'wp-admin/includes/screen.php' ); |
| 154 |
set_current_screen( 'sync' ); |
| 155 |
$skipped_items_ids = array(); |
| 156 |
// we estimate the total encoded size as we go by encoding each item individually |
| 157 |
// this is expensive, but the only way to really know :/ |
| 158 |
foreach ( $items as $key => $item ) { |
| 159 |
// Suspending cache addition help prevent overloading in memory cache of large sites. |
| 160 |
wp_suspend_cache_addition( true ); |
| 161 |
/** |
| 162 |
* Modify the data within an action before it is serialized and sent to the server |
| 163 |
* For example, during full sync this expands Post ID's into full Post objects, |
| 164 |
* so that we don't have to serialize the whole object into the queue. |
| 165 |
* |
| 166 |
* @since 4.2.0 |
| 167 |
* |
| 168 |
* @param array The action parameters |
| 169 |
* @param int The ID of the user who triggered the action |
| 170 |
*/ |
| 171 |
$item[1] = apply_filters( 'jetpack_sync_before_send_' . $item[0], $item[1], $item[2] ); |
| 172 |
wp_suspend_cache_addition( false ); |
| 173 |
if ( $item[1] === false ) { |
| 174 |
$skipped_items_ids[] = $key; |
| 175 |
continue; |
| 176 |
} |
| 177 |
$encoded_item = $encode ? $this->codec->encode( $item ) : $item; |
| 178 |
$upload_size += strlen( $encoded_item ); |
| 179 |
if ( $upload_size > $this->upload_max_bytes && count( $items_to_send ) > 0 ) { |
| 180 |
break; |
| 181 |
} |
| 182 |
$items_to_send[ $key ] = $encoded_item; |
| 183 |
if ( microtime(true) - $start_time > $this->max_dequeue_time ) { |
| 184 |
break; |
| 185 |
} |
| 186 |
} |
| 187 |
|
| 188 |
return array( $items_to_send, $skipped_items_ids, $items, microtime( true ) - $start_time ); |
| 189 |
} |
| 190 |
|
| 191 |
public function do_sync_for_queue( $queue ) { |
| 192 |
|
| 193 |
do_action( 'jetpack_sync_before_send_queue_' . $queue->id ); |
| 194 |
if ( $queue->size() === 0 ) { |
| 195 |
return new WP_Error( 'empty_queue_' . $queue->id ); |
| 196 |
} |
| 197 |
// now that we're sure we are about to sync, try to |
| 198 |
// ignore user abort so we can avoid getting into a |
| 199 |
// bad state |
| 200 |
if ( function_exists( 'ignore_user_abort' ) ) { |
| 201 |
ignore_user_abort( true ); |
| 202 |
} |
| 203 |
|
| 204 |
$checkout_start_time = microtime( true ); |
| 205 |
|
| 206 |
$buffer = $queue->checkout_with_memory_limit( $this->dequeue_max_bytes, $this->upload_max_rows ); |
| 207 |
|
| 208 |
if ( ! $buffer ) { |
| 209 |
// buffer has no items |
| 210 |
return new WP_Error( 'empty_buffer' ); |
| 211 |
} |
| 212 |
|
| 213 |
if ( is_wp_error( $buffer ) ) { |
| 214 |
return $buffer; |
| 215 |
} |
| 216 |
|
| 217 |
$checkout_duration = microtime( true ) - $checkout_start_time; |
| 218 |
|
| 219 |
list( $items_to_send, $skipped_items_ids, $items, $preprocess_duration ) = $this->get_items_to_send( $buffer, true ); |
| 220 |
if ( ! empty( $items_to_send ) ) { |
| 221 |
/** |
| 222 |
* Fires when data is ready to send to the server. |
| 223 |
* Return false or WP_Error to abort the sync (e.g. if there's an error) |
| 224 |
* The items will be automatically re-sent later |
| 225 |
* |
| 226 |
* @since 4.2.0 |
| 227 |
* |
| 228 |
* @param array $data The action buffer |
| 229 |
* @param string $codec The codec name used to encode the data |
| 230 |
* @param double $time The current time |
| 231 |
* @param string $queue The queue used to send ('sync' or 'full_sync') |
| 232 |
*/ |
| 233 |
Jetpack_Sync_Settings::set_is_sending( true ); |
| 234 |
$processed_item_ids = apply_filters( 'jetpack_sync_send_data', $items_to_send, $this->codec->name(), microtime( true ), $queue->id, $checkout_duration, $preprocess_duration ); |
| 235 |
Jetpack_Sync_Settings::set_is_sending( false ); |
| 236 |
} else { |
| 237 |
$processed_item_ids = $skipped_items_ids; |
| 238 |
$skipped_items_ids = array(); |
| 239 |
} |
| 240 |
|
| 241 |
if ( ! $processed_item_ids || is_wp_error( $processed_item_ids ) ) { |
| 242 |
$checked_in_item_ids = $queue->checkin( $buffer ); |
| 243 |
if ( is_wp_error( $checked_in_item_ids ) ) { |
| 244 |
error_log( 'Error checking in buffer: ' . $checked_in_item_ids->get_error_message() ); |
| 245 |
$queue->force_checkin(); |
| 246 |
} |
| 247 |
if ( is_wp_error( $processed_item_ids ) ) { |
| 248 |
return new WP_Error( 'wpcom_error', $processed_item_ids->get_error_code() ); |
| 249 |
} |
| 250 |
// returning a WP_Error('wpcom_error') is a sign to the caller that we should wait a while |
| 251 |
// before syncing again |
| 252 |
return new WP_Error( 'wpcom_error', 'jetpack_sync_send_data_false' ); |
| 253 |
} else { |
| 254 |
// detect if the last item ID was an error |
| 255 |
$had_wp_error = is_wp_error( end( $processed_item_ids ) ); |
| 256 |
if ( $had_wp_error ) { |
| 257 |
$wp_error = array_pop( $processed_item_ids ); |
| 258 |
} |
| 259 |
// also checkin any items that were skipped |
| 260 |
if ( count( $skipped_items_ids ) > 0 ) { |
| 261 |
$processed_item_ids = array_merge( $processed_item_ids, $skipped_items_ids ); |
| 262 |
} |
| 263 |
$processed_items = array_intersect_key( $items, array_flip( $processed_item_ids ) ); |
| 264 |
/** |
| 265 |
* Allows us to keep track of all the actions that have been sent. |
| 266 |
* Allows us to calculate the progress of specific actions. |
| 267 |
* |
| 268 |
* @since 4.2.0 |
| 269 |
* |
| 270 |
* @param array $processed_actions The actions that we send successfully. |
| 271 |
*/ |
| 272 |
do_action( 'jetpack_sync_processed_actions', $processed_items ); |
| 273 |
$queue->close( $buffer, $processed_item_ids ); |
| 274 |
// returning a WP_Error is a sign to the caller that we should wait a while |
| 275 |
// before syncing again |
| 276 |
if ( $had_wp_error ) { |
| 277 |
return new WP_Error( 'wpcom_error', $wp_error->get_error_code() ); |
| 278 |
} |
| 279 |
} |
| 280 |
return true; |
| 281 |
} |
| 282 |
|
| 283 |
function get_sync_queue() { |
| 284 |
return $this->sync_queue; |
| 285 |
} |
| 286 |
|
| 287 |
function get_full_sync_queue() { |
| 288 |
return $this->full_sync_queue; |
| 289 |
} |
| 290 |
|
| 291 |
function get_codec() { |
| 292 |
return $this->codec; |
| 293 |
} |
| 294 |
function set_codec() { |
| 295 |
if ( function_exists( 'gzinflate' ) ) { |
| 296 |
$this->codec = new Jetpack_Sync_JSON_Deflate_Array_Codec(); |
| 297 |
} else { |
| 298 |
$this->codec = new Jetpack_Sync_Simple_Codec(); |
| 299 |
} |
| 300 |
} |
| 301 |
|
| 302 |
function send_checksum() { |
| 303 |
require_once 'class.jetpack-sync-wp-replicastore.php'; |
| 304 |
$store = new Jetpack_Sync_WP_Replicastore(); |
| 305 |
do_action( 'jetpack_sync_checksum', $store->checksum_all() ); |
| 306 |
} |
| 307 |
|
| 308 |
function reset_sync_queue() { |
| 309 |
$this->sync_queue->reset(); |
| 310 |
} |
| 311 |
|
| 312 |
function reset_full_sync_queue() { |
| 313 |
$this->full_sync_queue->reset(); |
| 314 |
} |
| 315 |
|
| 316 |
function set_dequeue_max_bytes( $size ) { |
| 317 |
$this->dequeue_max_bytes = $size; |
| 318 |
} |
| 319 |
|
| 320 |
// in bytes |
| 321 |
function set_upload_max_bytes( $max_bytes ) { |
| 322 |
$this->upload_max_bytes = $max_bytes; |
| 323 |
} |
| 324 |
|
| 325 |
// in rows |
| 326 |
function set_upload_max_rows( $max_rows ) { |
| 327 |
$this->upload_max_rows = $max_rows; |
| 328 |
} |
| 329 |
|
| 330 |
// in seconds |
| 331 |
function set_sync_wait_time( $seconds ) { |
| 332 |
$this->sync_wait_time = $seconds; |
| 333 |
} |
| 334 |
|
| 335 |
function get_sync_wait_time() { |
| 336 |
return $this->sync_wait_time; |
| 337 |
} |
| 338 |
|
| 339 |
function set_enqueue_wait_time( $seconds ) { |
| 340 |
$this->enqueue_wait_time = $seconds; |
| 341 |
} |
| 342 |
|
| 343 |
function get_enqueue_wait_time() { |
| 344 |
return $this->enqueue_wait_time; |
| 345 |
} |
| 346 |
|
| 347 |
// in seconds |
| 348 |
function set_sync_wait_threshold( $seconds ) { |
| 349 |
$this->sync_wait_threshold = $seconds; |
| 350 |
} |
| 351 |
|
| 352 |
function get_sync_wait_threshold() { |
| 353 |
return $this->sync_wait_threshold; |
| 354 |
} |
| 355 |
|
| 356 |
// in seconds |
| 357 |
function set_max_dequeue_time( $seconds ) { |
| 358 |
$this->max_dequeue_time = $seconds; |
| 359 |
} |
| 360 |
|
| 361 |
|
| 362 |
|
| 363 |
function set_defaults() { |
| 364 |
$this->sync_queue = new Jetpack_Sync_Queue( 'sync' ); |
| 365 |
$this->full_sync_queue = new Jetpack_Sync_Queue( 'full_sync' ); |
| 366 |
$this->set_codec(); |
| 367 |
|
| 368 |
// saved settings |
| 369 |
Jetpack_Sync_Settings::set_importing( null ); |
| 370 |
$settings = Jetpack_Sync_Settings::get_settings(); |
| 371 |
$this->set_dequeue_max_bytes( $settings['dequeue_max_bytes'] ); |
| 372 |
$this->set_upload_max_bytes( $settings['upload_max_bytes'] ); |
| 373 |
$this->set_upload_max_rows( $settings['upload_max_rows'] ); |
| 374 |
$this->set_sync_wait_time( $settings['sync_wait_time'] ); |
| 375 |
$this->set_enqueue_wait_time( $settings['enqueue_wait_time'] ); |
| 376 |
$this->set_sync_wait_threshold( $settings['sync_wait_threshold'] ); |
| 377 |
$this->set_max_dequeue_time( Jetpack_Sync_Defaults::get_max_sync_execution_time() ); |
| 378 |
} |
| 379 |
|
| 380 |
function reset_data() { |
| 381 |
$this->reset_sync_queue(); |
| 382 |
$this->reset_full_sync_queue(); |
| 383 |
|
| 384 |
foreach ( Jetpack_Sync_Modules::get_modules() as $module ) { |
| 385 |
$module->reset_data(); |
| 386 |
} |
| 387 |
|
| 388 |
foreach ( array( 'sync', 'full_sync', 'full-sync-enqueue' ) as $queue_name ) { |
| 389 |
delete_option( self::NEXT_SYNC_TIME_OPTION_NAME . '_' . $queue_name ); |
| 390 |
} |
| 391 |
|
| 392 |
Jetpack_Sync_Settings::reset_data(); |
| 393 |
} |
| 394 |
|
| 395 |
function uninstall() { |
| 396 |
// Lets delete all the other fun stuff like transient and option and the sync queue |
| 397 |
$this->reset_data(); |
| 398 |
|
| 399 |
// delete the full sync status |
| 400 |
delete_option( 'jetpack_full_sync_status' ); |
| 401 |
|
| 402 |
// clear the sync cron. |
| 403 |
wp_clear_scheduled_hook( 'jetpack_sync_cron' ); |
| 404 |
wp_clear_scheduled_hook( 'jetpack_sync_full_cron' ); |
| 405 |
} |
| 406 |
} |
| 407 |
|