| @@ -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 | } |