| @@ -55,9 +55,20 @@ | ||
| 55 | 55 | const OM_OFFLOADED_FLAG = 'om_image_offloaded'; |
| 56 | 56 | const POST_OFFLOADED_FLAG = 'optimole_offload_post'; |
| 57 | 57 | const POST_ROLLBACK_FLAG = 'optimole_rollback_post'; |
| 58 | 58 | const RETRYABLE_META_COUNTER = '_optimole_retryable_errors'; |
| 59 | + | |
| 59 | 60 | /** |
| 61 | + * Transient name for the transfer lock. | |
| 62 | + */ | |
| 63 | + const TRANSFER_LOCK_TRANSIENT = 'optml_transfer_lock'; | |
| 64 | + | |
| 65 | + /** | |
| 66 | + * Time to live for the transfer lock, in seconds. | |
| 67 | + */ | |
| 68 | + const TRANSFER_LOCK_TTL = 600; | |
| 69 | + | |
| 70 | + /** | |
| 60 | 71 | * Flag used inside wp_get_attachment url filter. |
| 61 | 72 | * |
| 62 | 73 | * @var bool Whether or not to return the original url of the image. |
| 63 | 74 | */ |
| @@ -172,9 +183,9 @@ | ||
| 172 | 183 | } else { |
| 173 | 184 | add_filter( 'wp_insert_attachment_data', [ self::$instance, 'insert' ], 10, 4 ); |
| 174 | 185 | } |
| 175 | 186 | |
| 176 | - add_action( 'optml_start_processing_images', [ self::$instance, 'start_processing_images' ], 10, 5 ); | |
| 187 | + add_action( 'optml_start_processing_images', [ self::$instance, 'start_processing_images' ], 10, 6 ); | |
| 177 | 188 | add_action( |
| 178 | 189 | 'optml_move_images_by_id', |
| 179 | 190 | [ |
| 180 | 191 | self::$instance, |
| @@ -199,10 +210,11 @@ | ||
| 199 | 210 | * |
| 200 | 211 | * @return void |
| 201 | 212 | */ |
| 202 | 213 | public function maybe_reschedule() { |
| 214 | + $lock = get_transient( self::TRANSFER_LOCK_TRANSIENT ); | |
| 203 | 215 | // If this is in pending, we do nothing. |
| 204 | - if ( self::is_scheduled( 'optml_start_processing_images' ) ) { | |
| 216 | + if ( false !== $lock ) { | |
| 205 | 217 | return; |
| 206 | 218 | } |
| 207 | 219 | // If there is no transfer in progress, we do nothing. |
| 208 | 220 | if ( self::$instance->settings->get( 'transfer_status' ) === 'disabled' ) { |
| @@ -1898,20 +1910,27 @@ | ||
| 1898 | 1910 | 'status' => $in_progress, |
| 1899 | 1911 | 'action' => $type, |
| 1900 | 1912 | ]; |
| 1901 | 1913 | } |
| 1902 | - $total = ceil( $count / $batch ); | |
| 1903 | - self::schedule_action( | |
| 1904 | - time(), | |
| 1905 | - 'optml_start_processing_images', | |
| 1906 | - [ | |
| 1907 | - $action, | |
| 1908 | - $batch, | |
| 1909 | - 1, | |
| 1910 | - $total, | |
| 1911 | - $step, | |
| 1912 | - ] | |
| 1913 | - ); | |
| 1914 | + | |
| 1915 | + // We acquire a lock to prevent multiple workers from running the same action concurrently. | |
| 1916 | + $lock_token = self::acquire_transfer_lock( $action ); | |
| 1917 | + | |
| 1918 | + if ( false !== $lock_token ) { | |
| 1919 | + $total = ceil( $count / $batch ); | |
| 1920 | + self::schedule_action( | |
| 1921 | + time(), | |
| 1922 | + 'optml_start_processing_images', | |
| 1923 | + [ | |
| 1924 | + $action, | |
| 1925 | + $batch, | |
| 1926 | + 1, | |
| 1927 | + $total, | |
| 1928 | + $step, | |
| 1929 | + $lock_token, | |
| 1930 | + ] | |
| 1931 | + ); | |
| 1932 | + } | |
| 1914 | 1933 | } |
| 1915 | 1934 | |
| 1916 | 1935 | $response = [ |
| 1917 | 1936 | 'count' => $count, |
| @@ -1970,8 +1989,79 @@ | ||
| 1970 | 1989 | } |
| 1971 | 1990 | } |
| 1972 | 1991 | |
| 1973 | 1992 | /** |
| 1993 | + * Attempt to acquire the transfer lock for a given action. | |
| 1994 | + * | |
| 1995 | + * @param string $action The transfer action ('offload_images'|'rollback_images'). | |
| 1996 | + * | |
| 1997 | + * @return string|false The lock token on success, false if another worker already holds the lock. | |
| 1998 | + */ | |
| 1999 | + public static function acquire_transfer_lock( $action ) { | |
| 2000 | + $lock = get_transient( self::TRANSFER_LOCK_TRANSIENT ); | |
| 2001 | + if ( false !== $lock ) { | |
| 2002 | + return false; | |
| 2003 | + } | |
| 2004 | + | |
| 2005 | + $token = wp_generate_uuid4(); | |
| 2006 | + | |
| 2007 | + set_transient( | |
| 2008 | + self::TRANSFER_LOCK_TRANSIENT, | |
| 2009 | + [ | |
| 2010 | + 'token' => $token, | |
| 2011 | + 'action' => $action, | |
| 2012 | + ], | |
| 2013 | + self::TRANSFER_LOCK_TTL | |
| 2014 | + ); | |
| 2015 | + | |
| 2016 | + return $token; | |
| 2017 | + } | |
| 2018 | + | |
| 2019 | + /** | |
| 2020 | + * Renew the transfer lock if we still own it, extending its expiration. | |
| 2021 | + * | |
| 2022 | + * @param string $token The lock token this worker was given when it started the chain. | |
| 2023 | + * @param string $action The transfer action currently being processed. | |
| 2024 | + * | |
| 2025 | + * @return bool True if we still own the lock and renewed it, false if ownership was lost. | |
| 2026 | + */ | |
| 2027 | + public static function renew_transfer_lock( $token, $action ) { | |
| 2028 | + $lock = get_transient( self::TRANSFER_LOCK_TRANSIENT ); | |
| 2029 | + | |
| 2030 | + if ( ! is_array( $lock ) || ! isset( $lock['token'] ) || $lock['token'] !== $token ) { | |
| 2031 | + return false; | |
| 2032 | + } | |
| 2033 | + | |
| 2034 | + set_transient( | |
| 2035 | + self::TRANSFER_LOCK_TRANSIENT, | |
| 2036 | + [ | |
| 2037 | + 'token' => $token, | |
| 2038 | + 'action' => $action, | |
| 2039 | + ], | |
| 2040 | + self::TRANSFER_LOCK_TTL | |
| 2041 | + ); | |
| 2042 | + | |
| 2043 | + return true; | |
| 2044 | + } | |
| 2045 | + | |
| 2046 | + /** | |
| 2047 | + * Release the transfer lock if we still own it, allowing another worker to acquire it. | |
| 2048 | + * | |
| 2049 | + * @param string $token The lock token to release. | |
| 2050 | + * | |
| 2051 | + * @return void | |
| 2052 | + */ | |
| 2053 | + public static function release_transfer_lock( $token ) { | |
| 2054 | + $lock = get_transient( self::TRANSFER_LOCK_TRANSIENT ); | |
| 2055 | + | |
| 2056 | + if ( ! is_array( $lock ) || ! isset( $lock['token'] ) || $lock['token'] !== $token ) { | |
| 2057 | + return; | |
| 2058 | + } | |
| 2059 | + | |
| 2060 | + delete_transient( self::TRANSFER_LOCK_TRANSIENT ); | |
| 2061 | + } | |
| 2062 | + | |
| 2063 | + /** | |
| 1974 | 2064 | * Start Processing Images by IDs |
| 1975 | 2065 | * |
| 1976 | 2066 | * @param string $action The action for which to get the number of images. |
| 1977 | 2067 | * @param int $id The images to process. |
| @@ -2028,19 +2118,26 @@ | ||
| 2028 | 2118 | * @param int $batch The batch of images to process. |
| 2029 | 2119 | * @param int $page The page of images to process. |
| 2030 | 2120 | * @param int $total The total number of pages. |
| 2031 | 2121 | * @param int $step The current step. |
| 2122 | + * @param string $lock_token The transfer lock token owned by this processing chain. | |
| 2032 | 2123 | * |
| 2033 | 2124 | * @return void |
| 2034 | 2125 | */ |
| 2035 | - public function start_processing_images( $action, $batch, $page, $total, $step ) { | |
| 2126 | + public function start_processing_images( $action, $batch, $page, $total, $step, $lock_token = '' ) { | |
| 2036 | 2127 | $option = 'offload_images' === $action ? 'offloading_status' : 'rollback_status'; |
| 2037 | 2128 | $type = 'offload_images' === $action ? 'offload' : 'rollback'; |
| 2038 | 2129 | |
| 2039 | 2130 | if ( self::$instance->settings->get( $option ) === 'disabled' ) { |
| 2131 | + self::release_transfer_lock( $lock_token ); | |
| 2040 | 2132 | return; |
| 2041 | 2133 | } |
| 2042 | 2134 | |
| 2135 | + // If we don't own the lock anymore, stop processing. | |
| 2136 | + if ( ! self::renew_transfer_lock( $lock_token, $action ) ) { | |
| 2137 | + return; | |
| 2138 | + } | |
| 2139 | + | |
| 2043 | 2140 | if ( $step > $total || 0 === $total ) { |
| 2044 | 2141 | $meta = self::get_process_meta(); |
| 2045 | 2142 | self::$instance->logger->add_log( $type, 'Process finished with ' . $meta['count'] . ' items in ' . $meta['time_passed'] . ' minutes.' ); |
| 2046 | 2143 | |
| @@ -2047,8 +2144,11 @@ | ||
| 2047 | 2144 | self::$instance->settings->update( $option, 'disabled' ); |
| 2048 | 2145 | |
| 2049 | 2146 | self::$instance->settings->update( 'show_offload_finish_notice', $type ); |
| 2050 | 2147 | |
| 2148 | + // Transfer completed successfully: release the lock. | |
| 2149 | + self::release_transfer_lock( $lock_token ); | |
| 2150 | + | |
| 2051 | 2151 | return; |
| 2052 | 2152 | } |
| 2053 | 2153 | |
| 2054 | 2154 | set_time_limit( 0 ); |
| @@ -2079,12 +2179,14 @@ | ||
| 2079 | 2179 | $batch, |
| 2080 | 2180 | $page, |
| 2081 | 2181 | $total, |
| 2082 | 2182 | $step, |
| 2183 | + $lock_token, | |
| 2083 | 2184 | ] |
| 2084 | 2185 | ); |
| 2085 | 2186 | } catch ( Exception $e ) { |
| 2086 | 2187 | // Reschedule the cron to run again after a delay. Sometimes memory limit is exausted. |
| 2188 | + // This is a retryable error, so the lock is kept rather than released. | |
| 2087 | 2189 | $delay_in_seconds = 10; |
| 2088 | 2190 | self::$instance->logger->add_log( $type, $e->getMessage() ); |
| 2089 | 2191 | |
| 2090 | 2192 | self::schedule_action( |
| @@ -2095,8 +2197,9 @@ | ||
| 2095 | 2197 | $batch, |
| 2096 | 2198 | $page, |
| 2097 | 2199 | $total, |
| 2098 | 2200 | $step, |
| 2201 | + $lock_token, | |
| 2099 | 2202 | ] |
| 2100 | 2203 | ); |
| 2101 | 2204 | } |
| 2102 | 2205 | } |
| @@ -2672,16 +2775,58 @@ | ||
| 2672 | 2775 | } |
| 2673 | 2776 | |
| 2674 | 2777 | /** |
| 2675 | 2778 | * Cleanup the offload errors meta. |
| 2779 | + * | |
| 2780 | + * @param string $meta_key The meta key to delete. Defaults to the offload error key. | |
| 2781 | + * | |
| 2782 | + * @return int|bool Number of rows affected/selected or false on error. | |
| 2676 | 2783 | */ |
| 2677 | - public static function clear_offload_errors_meta() { | |
| 2784 | + public static function clear_offload_errors_meta( $meta_key = '' ) { | |
| 2678 | 2785 | global $wpdb; |
| 2679 | 2786 | |
| 2680 | - return $wpdb->query( | |
| 2787 | + if ( empty( $meta_key ) ) { | |
| 2788 | + $meta_key = self::META_KEYS['offload_error']; | |
| 2789 | + } | |
| 2790 | + | |
| 2791 | + // Collect the affected attachments before the bulk delete so their object | |
| 2792 | + // caches can be invalidated. A raw DELETE bypasses the meta/query caches, | |
| 2793 | + // which would otherwise leave stale WP_Query results for subsequent queries. | |
| 2794 | + $post_ids = $wpdb->get_col( | |
| 2681 | 2795 | $wpdb->prepare( |
| 2796 | + "SELECT post_id FROM {$wpdb->postmeta} WHERE meta_key = %s", | |
| 2797 | + $meta_key | |
| 2798 | + ) | |
| 2799 | + ); | |
| 2800 | + | |
| 2801 | + $result = $wpdb->query( | |
| 2802 | + $wpdb->prepare( | |
| 2682 | 2803 | "DELETE FROM {$wpdb->postmeta} WHERE meta_key = %s", |
| 2683 | - self::META_KEYS['offload_error'] | |
| 2804 | + $meta_key | |
| 2684 | 2805 | ) |
| 2685 | 2806 | ); |
| 2807 | + | |
| 2808 | + foreach ( $post_ids as $post_id ) { | |
| 2809 | + wp_cache_delete( (int) $post_id, 'post_meta' ); | |
| 2810 | + } | |
| 2811 | + | |
| 2812 | + // Bump the posts last_changed so cached WP_Query results (which are keyed | |
| 2813 | + // on it) are recomputed on the next query. The raw DELETE above does not | |
| 2814 | + // touch the object cache, so without this the retried rollback/offload | |
| 2815 | + // query could return a stale set that still excludes the cleared posts. | |
| 2816 | + wp_cache_set( 'last_changed', microtime(), 'posts' ); | |
| 2817 | + | |
| 2818 | + return $result; | |
| 2819 | + } | |
| 2820 | + | |
| 2821 | + /** | |
| 2822 | + * Cleanup the rollback errors meta. | |
| 2823 | + * | |
| 2824 | + * Used when the user retries the rollback process so previously errored | |
| 2825 | + * attachments are considered again for restore. | |
| 2826 | + * | |
| 2827 | + * @return int|bool Number of rows affected/selected or false on error. | |
| 2828 | + */ | |
| 2829 | + public static function clear_rollback_errors_meta() { | |
| 2830 | + return self::clear_offload_errors_meta( self::META_KEYS['rollback_error'] ); | |
| 2686 | 2831 | } |
| 2687 | 2832 | } |