is_locked() ) { return new WP_Error( 'locked_queue_' . $queue->id ); } if ( $queue->size() === 0 ) { return new WP_Error( 'empty_queue_' . $queue->id ); } if ( get_transient( Sender::TEMP_SYNC_DISABLE_TRANSIENT_NAME ) ) { return new WP_Error( 'sender_temporarily_disabled_while_pulling' ); } // Return early if we've gotten a retry-after header response that is not expired. $retry_time = get_option( Actions::RETRY_AFTER_PREFIX . $queue->id ); if ( $retry_time && $retry_time >= microtime( true ) ) { return new WP_Error( 'retry_after_' . $queue->id ); } // Don't sync if we are throttled. $sync_next_time = Sender::get_instance()->get_next_sync_time( $queue->id ); if ( $sync_next_time > microtime( true ) ) { return new WP_Error( 'sync_throttled_' . $queue->id ); } /** * How much time to wait before we start suspecting Dedicated Sync is in trouble. */ $queue_send_time_threshold = 30 * MINUTE_IN_SECONDS; $queue_lag = $queue->lag(); /** * Try to acquire a request lock, so we don't spawn multiple requests at the same time. * This should prevent cases where sites might have limits on the amount of simultaneous requests. */ $request_lock = self::try_lock_spawn_request(); if ( ! $request_lock ) { return new WP_Error( 'dedicated_request_lock', 'Unable to acquire request lock' ); } /** * If the queue lag is bigger than the threshold, we want to check if Dedicated Sync is working correctly. * We will do by sending a test request and disabling Dedicated Sync if it's not working. We will also exit early * in case we send the test request since it is a blocking request. */ if ( $queue_lag > $queue_send_time_threshold ) { if ( false === get_transient( self::DEDICATED_SYNC_CHECK_TRANSIENT ) ) { if ( ! self::can_spawn_dedicated_sync_request() ) { self::on_dedicated_sync_lag_not_sending_threshold_reached(); return new WP_Error( 'dedicated_sync_not_sending', 'Dedicated Sync is not successfully sending events' ); } return true; } } $url = rest_url( 'jetpack/v4/sync/spawn-sync' ); $url = add_query_arg( 'time', time(), $url ); // Enforce Cache busting. $url = add_query_arg( self::DEDICATED_SYNC_REQUEST_LOCK_QUERY_PARAM_NAME, $request_lock, $url ); $args = array( 'cookies' => $_COOKIE, 'blocking' => false, 'timeout' => 0.01, /** This filter is documented in wp-includes/class-wp-http-streams.php */ 'sslverify' => apply_filters( 'https_local_ssl_verify', false ), ); $result = wp_remote_get( $url, $args ); if ( is_wp_error( $result ) ) { return $result; } return true; } /** * Attempt to acquire a request lock. * * To avoid spawning multiple requests at the same time, we need to have a quick lock that will * allow only a single request to continue if we try to spawn multiple at the same time. * * @return string|false */ public static function try_lock_spawn_request() { $option_name = self::DEDICATED_SYNC_REQUEST_LOCK_OPTION_NAME; $expires_name = $option_name . '_expires'; $ttl = self::DEDICATED_SYNC_REQUEST_LOCK_TIMEOUT; $lock_id = wp_generate_uuid4(); $now = microtime( true ); // Fast path: external object cache is atomic. if ( wp_using_ext_object_cache() ) { if ( wp_cache_add( $option_name, $lock_id, 'jetpack', $ttl ) ) { return $lock_id; } return false; // Worker already active } global $wpdb; // 1) Check & clear expired lock (best effort; failure here is harmless) $expiry = (float) \Jetpack_Options::get_raw_option( $expires_name, 0 ); if ( ! $expiry || $expiry < $now ) { // Either missing (edge case) or expired → clean up \Jetpack_Options::delete_raw_option( $option_name ); \Jetpack_Options::delete_raw_option( $expires_name ); } // 2) Atomic acquisition: INSERT IGNORE (succeeds only if the lock doesn't exist) $inserted = $wpdb->query( // phpcs:disable WordPress.DB.DirectDatabaseQuery.DirectQuery,WordPress.DB.DirectDatabaseQuery.NoCaching --- Ensure atomicity. $wpdb->prepare( "INSERT IGNORE INTO $wpdb->options ( option_name, option_value, autoload ) VALUES ( %s, %s, 'no' )", $option_name, maybe_serialize( $lock_id ) ) ); if ( $inserted ) { // 3) We own the lock — store expiry separately \Jetpack_Options::update_raw_option( $expires_name, $now + $ttl, false ); return $lock_id; // Success } // Lock already present → normal state → do not spawn return false; } /** * Attempt to release the request lock. * * @param string $lock_id The request lock that's currently being held. * * @return bool|WP_Error */ public static function try_release_lock_spawn_request( $lock_id = '' ) { // Try to get the lock_id from the current request if it's not supplied. if ( empty( $lock_id ) ) { $lock_id = self::get_request_lock_id_from_request(); } // If it's still not a valid lock_id, throw an error and let the lock process figure it out. if ( empty( $lock_id ) ) { return new WP_Error( 'dedicated_request_lock_invalid', 'Invalid lock_id supplied for unlock' ); } if ( wp_using_ext_object_cache() ) { $cached = wp_cache_get( self::DEDICATED_SYNC_REQUEST_LOCK_OPTION_NAME, 'jetpack', true ); if ( (string) $lock_id === $cached ) { wp_cache_delete( self::DEDICATED_SYNC_REQUEST_LOCK_OPTION_NAME, 'jetpack' ); return true; } return false; } // If this is the flow that has the lock, let's release it so we can spawn other requests afterwards $current_lock_value = \Jetpack_Options::get_raw_option( self::DEDICATED_SYNC_REQUEST_LOCK_OPTION_NAME, null ); if ( (string) $lock_id === $current_lock_value ) { \Jetpack_Options::delete_raw_option( self::DEDICATED_SYNC_REQUEST_LOCK_OPTION_NAME ); return true; } return false; } /** * Try to get the request lock id from the current request. * * @return array|string|string[]|null */ public static function get_request_lock_id_from_request() { // phpcs:ignore WordPress.Security.NonceVerification.Recommended if ( ! isset( $_GET[ self::DEDICATED_SYNC_REQUEST_LOCK_QUERY_PARAM_NAME ] ) ) { return null; } // phpcs:ignore WordPress.Security.NonceVerification.Recommended,WordPress.Security.ValidatedSanitizedInput.InputNotSanitized return wp_unslash( $_GET[ self::DEDICATED_SYNC_REQUEST_LOCK_QUERY_PARAM_NAME ] ); } /** * Test Sync spawning functionality by making a request to the * Sync spawning endpoint and storing the result (status code) in a transient. * * @since 1.34.0 * * @return bool True if we got a successful response, false otherwise. */ public static function can_spawn_dedicated_sync_request() { $dedicated_sync_check_transient = self::DEDICATED_SYNC_CHECK_TRANSIENT; $dedicated_sync_response_body = get_transient( $dedicated_sync_check_transient ); if ( false === $dedicated_sync_response_body ) { $url = rest_url( 'jetpack/v4/sync/spawn-sync' ); $url = add_query_arg( 'time', time(), $url ); // Enforce Cache busting. $args = array( 'cookies' => $_COOKIE, 'timeout' => 30, /** This filter is documented in wp-includes/class-wp-http-streams.php */ 'sslverify' => apply_filters( 'https_local_ssl_verify', false ), ); $response = wp_remote_get( $url, $args ); $dedicated_sync_response_code = wp_remote_retrieve_response_code( $response ); $dedicated_sync_response_body = trim( wp_remote_retrieve_body( $response ) ); /** * Limit the size of the body that we save in the transient to avoid cases where an error * occurs and a whole generated HTML page is returned. We don't need to store the whole thing. * * The regexp check is done to make sure we can detect the string even if the body returns some additional * output, like some caching plugins do when they try to pad the request. */ $regexp = '!' . preg_quote( self::DEDICATED_SYNC_VALIDATION_STRING, '!' ) . '!uis'; if ( preg_match( $regexp, $dedicated_sync_response_body ) ) { $saved_response_body = self::DEDICATED_SYNC_VALIDATION_STRING; } else { $saved_response_body = time(); } set_transient( $dedicated_sync_check_transient, $saved_response_body, HOUR_IN_SECONDS ); // Send a bit more information to WordPress.com to help debugging issues. if ( $saved_response_body !== self::DEDICATED_SYNC_VALIDATION_STRING ) { $data = array( 'timestamp' => microtime( true ), 'response_code' => $dedicated_sync_response_code, 'response_body' => $dedicated_sync_response_body, // Send the flow type that was attempted. 'sync_flow_type' => 'dedicated', ); $sender = Sender::get_instance(); $sender->send_action( 'jetpack_sync_flow_error_enable', $data ); } } return self::DEDICATED_SYNC_VALIDATION_STRING === $dedicated_sync_response_body; } /** * Disable dedicated sync and set a transient to prevent re-enabling it for some time. * * @return void */ public static function on_dedicated_sync_lag_not_sending_threshold_reached() { set_transient( self::DEDICATED_SYNC_TEMPORARY_DISABLE_FLAG, true, 6 * HOUR_IN_SECONDS ); Settings::update_settings( array( 'dedicated_sync_enabled' => 0, ) ); // Inform that we had to temporarily disable Dedicated Sync $data = array( 'timestamp' => microtime( true ), // Send the flow type that was attempted. 'sync_flow_type' => 'dedicated', ); $sender = Sender::get_instance(); $sender->send_action( 'jetpack_sync_flow_error_temp_disable', $data ); } /** * Disable or enable Dedicated Sync sender based on the header value returned from WordPress.com * * @param string $dedicated_sync_header The Dedicated Sync header value - `on` or `off`. * * @return bool Whether Dedicated Sync is going to be enabled or not. */ public static function maybe_change_dedicated_sync_status_from_wpcom_header( $dedicated_sync_header ) { $dedicated_sync_enabled = 'on' === $dedicated_sync_header ? 1 : 0; // Prevent enabling of Dedicated sync via header flag if we're in an autoheal timeout. if ( $dedicated_sync_enabled ) { $check_transient = get_transient( self::DEDICATED_SYNC_TEMPORARY_DISABLE_FLAG ); if ( $check_transient ) { // Something happened and Dedicated Sync should not be automatically re-enabled. return false; } } $current_setting = Settings::is_dedicated_sync_enabled(); // No need to update if current setting matches header value. if ( $current_setting === (bool) $dedicated_sync_enabled ) { return $current_setting; } Settings::update_settings( array( 'dedicated_sync_enabled' => $dedicated_sync_enabled, ) ); return Settings::is_dedicated_sync_enabled(); } }