class.jetpack-sync-sender.php
| 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 (float) 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 ) > (float) $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 | private function fastcgi_finish_request() { |
| 192 | if ( function_exists( 'fastcgi_finish_request' ) && version_compare( phpversion(), '7.0.16', '>=' ) ) { |
| 193 | fastcgi_finish_request(); |
| 194 | } |
| 195 | } |
| 196 | |
| 197 | public function do_sync_for_queue( $queue ) { |
| 198 | |
| 199 | do_action( 'jetpack_sync_before_send_queue_' . $queue->id ); |
| 200 | if ( $queue->size() === 0 ) { |
| 201 | return new WP_Error( 'empty_queue_' . $queue->id ); |
| 202 | } |
| 203 | // now that we're sure we are about to sync, try to |
| 204 | // ignore user abort so we can avoid getting into a |
| 205 | // bad state |
| 206 | if ( function_exists( 'ignore_user_abort' ) ) { |
| 207 | ignore_user_abort( true ); |
| 208 | } |
| 209 | |
| 210 | /* Don't make the request block till we finish, if possible. */ |
| 211 | if ( Jetpack_Constants::is_true( 'REST_REQUEST' ) || Jetpack_Constants::is_true('XMLRPC_REQUEST' ) ) { |
| 212 | $this->fastcgi_finish_request(); |
| 213 | } |
| 214 | |
| 215 | $checkout_start_time = microtime( true ); |
| 216 | |
| 217 | $buffer = $queue->checkout_with_memory_limit( $this->dequeue_max_bytes, $this->upload_max_rows ); |
| 218 | |
| 219 | if ( ! $buffer ) { |
| 220 | // buffer has no items |
| 221 | return new WP_Error( 'empty_buffer' ); |
| 222 | } |
| 223 | |
| 224 | if ( is_wp_error( $buffer ) ) { |
| 225 | return $buffer; |
| 226 | } |
| 227 | |
| 228 | $checkout_duration = microtime( true ) - $checkout_start_time; |
| 229 | |
| 230 | list( $items_to_send, $skipped_items_ids, $items, $preprocess_duration ) = $this->get_items_to_send( $buffer, true ); |
| 231 | if ( ! empty( $items_to_send ) ) { |
| 232 | /** |
| 233 | * Fires when data is ready to send to the server. |
| 234 | * Return false or WP_Error to abort the sync (e.g. if there's an error) |
| 235 | * The items will be automatically re-sent later |
| 236 | * |
| 237 | * @since 4.2.0 |
| 238 | * |
| 239 | * @param array $data The action buffer |
| 240 | * @param string $codec The codec name used to encode the data |
| 241 | * @param double $time The current time |
| 242 | * @param string $queue The queue used to send ('sync' or 'full_sync') |
| 243 | */ |
| 244 | Jetpack_Sync_Settings::set_is_sending( true ); |
| 245 | $processed_item_ids = apply_filters( 'jetpack_sync_send_data', $items_to_send, $this->codec->name(), microtime( true ), $queue->id, $checkout_duration, $preprocess_duration ); |
| 246 | Jetpack_Sync_Settings::set_is_sending( false ); |
| 247 | } else { |
| 248 | $processed_item_ids = $skipped_items_ids; |
| 249 | $skipped_items_ids = array(); |
| 250 | } |
| 251 | |
| 252 | if ( ! $processed_item_ids || is_wp_error( $processed_item_ids ) ) { |
| 253 | $checked_in_item_ids = $queue->checkin( $buffer ); |
| 254 | if ( is_wp_error( $checked_in_item_ids ) ) { |
| 255 | error_log( 'Error checking in buffer: ' . $checked_in_item_ids->get_error_message() ); |
| 256 | $queue->force_checkin(); |
| 257 | } |
| 258 | if ( is_wp_error( $processed_item_ids ) ) { |
| 259 | return new WP_Error( 'wpcom_error', $processed_item_ids->get_error_code() ); |
| 260 | } |
| 261 | // returning a WP_Error('wpcom_error') is a sign to the caller that we should wait a while |
| 262 | // before syncing again |
| 263 | return new WP_Error( 'wpcom_error', 'jetpack_sync_send_data_false' ); |
| 264 | } else { |
| 265 | // detect if the last item ID was an error |
| 266 | $had_wp_error = is_wp_error( end( $processed_item_ids ) ); |
| 267 | if ( $had_wp_error ) { |
| 268 | $wp_error = array_pop( $processed_item_ids ); |
| 269 | } |
| 270 | // also checkin any items that were skipped |
| 271 | if ( count( $skipped_items_ids ) > 0 ) { |
| 272 | $processed_item_ids = array_merge( $processed_item_ids, $skipped_items_ids ); |
| 273 | } |
| 274 | $processed_items = array_intersect_key( $items, array_flip( $processed_item_ids ) ); |
| 275 | /** |
| 276 | * Allows us to keep track of all the actions that have been sent. |
| 277 | * Allows us to calculate the progress of specific actions. |
| 278 | * |
| 279 | * @since 4.2.0 |
| 280 | * |
| 281 | * @param array $processed_actions The actions that we send successfully. |
| 282 | */ |
| 283 | do_action( 'jetpack_sync_processed_actions', $processed_items ); |
| 284 | $queue->close( $buffer, $processed_item_ids ); |
| 285 | // returning a WP_Error is a sign to the caller that we should wait a while |
| 286 | // before syncing again |
| 287 | if ( $had_wp_error ) { |
| 288 | return new WP_Error( 'wpcom_error', $wp_error->get_error_code() ); |
| 289 | } |
| 290 | } |
| 291 | return true; |
| 292 | } |
| 293 | |
| 294 | function get_sync_queue() { |
| 295 | return $this->sync_queue; |
| 296 | } |
| 297 | |
| 298 | function get_full_sync_queue() { |
| 299 | return $this->full_sync_queue; |
| 300 | } |
| 301 | |
| 302 | function get_codec() { |
| 303 | return $this->codec; |
| 304 | } |
| 305 | function set_codec() { |
| 306 | if ( function_exists( 'gzinflate' ) ) { |
| 307 | $this->codec = new Jetpack_Sync_JSON_Deflate_Array_Codec(); |
| 308 | } else { |
| 309 | $this->codec = new Jetpack_Sync_Simple_Codec(); |
| 310 | } |
| 311 | } |
| 312 | |
| 313 | function send_checksum() { |
| 314 | require_once 'class.jetpack-sync-wp-replicastore.php'; |
| 315 | $store = new Jetpack_Sync_WP_Replicastore(); |
| 316 | do_action( 'jetpack_sync_checksum', $store->checksum_all() ); |
| 317 | } |
| 318 | |
| 319 | function reset_sync_queue() { |
| 320 | $this->sync_queue->reset(); |
| 321 | } |
| 322 | |
| 323 | function reset_full_sync_queue() { |
| 324 | $this->full_sync_queue->reset(); |
| 325 | } |
| 326 | |
| 327 | function set_dequeue_max_bytes( $size ) { |
| 328 | $this->dequeue_max_bytes = $size; |
| 329 | } |
| 330 | |
| 331 | // in bytes |
| 332 | function set_upload_max_bytes( $max_bytes ) { |
| 333 | $this->upload_max_bytes = $max_bytes; |
| 334 | } |
| 335 | |
| 336 | // in rows |
| 337 | function set_upload_max_rows( $max_rows ) { |
| 338 | $this->upload_max_rows = $max_rows; |
| 339 | } |
| 340 | |
| 341 | // in seconds |
| 342 | function set_sync_wait_time( $seconds ) { |
| 343 | $this->sync_wait_time = $seconds; |
| 344 | } |
| 345 | |
| 346 | function get_sync_wait_time() { |
| 347 | return $this->sync_wait_time; |
| 348 | } |
| 349 | |
| 350 | function set_enqueue_wait_time( $seconds ) { |
| 351 | $this->enqueue_wait_time = $seconds; |
| 352 | } |
| 353 | |
| 354 | function get_enqueue_wait_time() { |
| 355 | return $this->enqueue_wait_time; |
| 356 | } |
| 357 | |
| 358 | // in seconds |
| 359 | function set_sync_wait_threshold( $seconds ) { |
| 360 | $this->sync_wait_threshold = $seconds; |
| 361 | } |
| 362 | |
| 363 | function get_sync_wait_threshold() { |
| 364 | return $this->sync_wait_threshold; |
| 365 | } |
| 366 | |
| 367 | // in seconds |
| 368 | function set_max_dequeue_time( $seconds ) { |
| 369 | $this->max_dequeue_time = $seconds; |
| 370 | } |
| 371 | |
| 372 | |
| 373 | |
| 374 | function set_defaults() { |
| 375 | $this->sync_queue = new Jetpack_Sync_Queue( 'sync' ); |
| 376 | $this->full_sync_queue = new Jetpack_Sync_Queue( 'full_sync' ); |
| 377 | $this->set_codec(); |
| 378 | |
| 379 | // saved settings |
| 380 | Jetpack_Sync_Settings::set_importing( null ); |
| 381 | $settings = Jetpack_Sync_Settings::get_settings(); |
| 382 | $this->set_dequeue_max_bytes( $settings['dequeue_max_bytes'] ); |
| 383 | $this->set_upload_max_bytes( $settings['upload_max_bytes'] ); |
| 384 | $this->set_upload_max_rows( $settings['upload_max_rows'] ); |
| 385 | $this->set_sync_wait_time( $settings['sync_wait_time'] ); |
| 386 | $this->set_enqueue_wait_time( $settings['enqueue_wait_time'] ); |
| 387 | $this->set_sync_wait_threshold( $settings['sync_wait_threshold'] ); |
| 388 | $this->set_max_dequeue_time( Jetpack_Sync_Defaults::get_max_sync_execution_time() ); |
| 389 | } |
| 390 | |
| 391 | function reset_data() { |
| 392 | $this->reset_sync_queue(); |
| 393 | $this->reset_full_sync_queue(); |
| 394 | |
| 395 | foreach ( Jetpack_Sync_Modules::get_modules() as $module ) { |
| 396 | $module->reset_data(); |
| 397 | } |
| 398 | |
| 399 | foreach ( array( 'sync', 'full_sync', 'full-sync-enqueue' ) as $queue_name ) { |
| 400 | delete_option( self::NEXT_SYNC_TIME_OPTION_NAME . '_' . $queue_name ); |
| 401 | } |
| 402 | |
| 403 | Jetpack_Sync_Settings::reset_data(); |
| 404 | } |
| 405 | |
| 406 | function uninstall() { |
| 407 | // Lets delete all the other fun stuff like transient and option and the sync queue |
| 408 | $this->reset_data(); |
| 409 | |
| 410 | // delete the full sync status |
| 411 | delete_option( 'jetpack_full_sync_status' ); |
| 412 | |
| 413 | // clear the sync cron. |
| 414 | wp_clear_scheduled_hook( 'jetpack_sync_cron' ); |
| 415 | wp_clear_scheduled_hook( 'jetpack_sync_full_cron' ); |
| 416 | } |
| 417 | } |
| 418 |