PluginProbe
Jetpack – WP Security, Backup, Speed, & Growth / 16.3-a.5
Jetpack – WP Security, Backup, Speed, & Growth v16.3-a.5
16.3-a.5 16.3-a.7 16.3-a.3 16.3-a.1 16.2 16.2-beta 12.0.3 12.1.3 12.2.3 12.3.2 12.4.2 12.5.2 12.6.4 12.7.3 12.8.3 12.9.5 13.0.2 13.1.5 13.2.4 13.3.3 13.4.5 13.5.2 13.6.2 13.7.2 13.8.3 All 506 releases
← All changes | jetpack_vendor/automattic/jetpack-sync/src/class-queue.php +141 -174 12.0.3 → 16.3-a.5 View file →
@@ -6,8 +6,10 @@
6 6 */
7 7
8 8 namespace Automattic\Jetpack\Sync;
9 9
10 +use Automattic\Jetpack\Sync\Queue\Queue_Storage_Options;
11 +use Automattic\Jetpack\Sync\Queue\Queue_Storage_Table;
10 12 use WP_Error;
11 13
12 14 /**
13 15 * A persistent queue that can be flushed in increments of N items,
@@ -37,8 +39,17 @@
37 39 */
38 40 public $random_int;
39 41
40 42 /**
43 + * Queue Storage instance where we'll store the queue items.
44 + *
45 + * For now, it's only the Options table. To be updated to include the Custom table in future updates.
46 + *
47 + * @var Queue_Storage_Options|Queue_Storage_Table|null
48 + */
49 + public $queue_storage = null;
50 +
51 + /**
41 52 * Queue constructor.
42 53 *
43 54 * @param string $id Name of the queue.
44 55 */
@@ -45,8 +56,18 @@
45 56 public function __construct( $id ) {
46 57 $this->id = str_replace( '-', '_', $id ); // Necessary to ensure we don't have ID collisions in the SQL.
47 58 $this->row_iterator = 0;
48 59 $this->random_int = wp_rand( 1, 1000000 );
60 +
61 + /**
62 + * If the Custom queue table is enabled - let's use it as a backend. Otherwise, fall back to the Options table.
63 + */
64 + if ( Settings::is_custom_queue_table_enabled() ) {
65 + $this->queue_storage = new Queue_Storage_Table( $this->id );
66 + } else {
67 + // Initialize the storage with the Options table backend. To be changed in subsequent updates to include the logic to switch to Custom Table.
68 + $this->queue_storage = new Queue_Storage_Options( $this->id );
69 + }
49 70 }
50 71
51 72 /**
52 73 * Add a single item to the queue.
@@ -51,16 +72,17 @@
51 72 /**
52 73 * Add a single item to the queue.
53 74 *
54 75 * @param object $item Event object to add to queue.
76 + *
77 + * @return bool|WP_Error
55 78 */
56 79 public function add( $item ) {
57 - global $wpdb;
58 80 $added = false;
59 81
60 82 // If empty, don't add.
61 83 if ( empty( $item ) ) {
62 - return;
84 + return false;
63 85 }
64 86
65 87 // Attempt to serialize data, if an exception (closures) return early.
66 88 try {
@@ -65,24 +87,18 @@
65 87 // Attempt to serialize data, if an exception (closures) return early.
66 88 try {
67 89 $item = serialize( $item ); // phpcs:ignore WordPress.PHP.DiscouragedPHPFunctions.serialize_serialize
68 90 } catch ( \Exception $ex ) {
69 - return;
91 + return new WP_Error( 'queue_unable_to_serialize', 'Unable to serialize item' );
70 92 }
71 93
72 94 // This basically tries to add the option until enough time has elapsed that
73 95 // it has a unique (microtime-based) option key.
74 96 while ( ! $added ) {
75 - $rows_added = $wpdb->query(
76 - $wpdb->prepare(
77 - "INSERT INTO $wpdb->options (option_name, option_value, autoload) VALUES (%s, %s,%s)",
78 - $this->get_next_data_row_option_name(),
79 - $item,
80 - 'no'
81 - )
82 - );
83 - $added = ( 0 !== $rows_added );
97 + $added = $this->queue_storage->insert_item( $this->get_next_data_row_option_name(), $item );
84 98 }
99 +
100 + return $added;
85 101 }
86 102
87 103 /**
88 104 * Insert all the items in a single SQL query. May be subject to query size limits!
@@ -91,40 +107,22 @@
91 107 *
92 108 * @return bool|\WP_Error
93 109 */
94 110 public function add_all( $items ) {
95 - global $wpdb;
111 + // TODO check and figure out if it's used at all and if we can optimize it.
96 112 $base_option_name = $this->get_next_data_row_option_name();
97 113
98 - $query = "INSERT INTO $wpdb->options (option_name, option_value, autoload) VALUES ";
114 + $rows_added = $this->queue_storage->add_all( $items, $base_option_name );
99 115
100 - $rows = array();
101 - $count_items = count( $items );
102 - for ( $i = 0; $i < $count_items; ++$i ) {
103 - // skip empty items.
104 - if ( empty( $items[ $i ] ) ) {
105 - continue;
106 - }
107 - try {
108 - $option_name = esc_sql( $base_option_name . '-' . $i );
109 - $option_value = esc_sql( serialize( $items[ $i ] ) ); // phpcs:ignore WordPress.PHP.DiscouragedPHPFunctions.serialize_serialize
110 - $rows[] = "('$option_name', '$option_value', 'no')";
111 - } catch ( \Exception $e ) {
112 - // Item cannot be serialized so skip.
113 - continue;
114 - }
115 - }
116 -
117 - $rows_added = $wpdb->query( $query . join( ',', $rows ) ); // phpcs:ignore WordPress.DB.PreparedSQL.NotPrepared
118 -
119 116 if ( count( $items ) !== $rows_added ) {
120 117 return new WP_Error( 'row_count_mismatch', "The number of rows inserted didn't match the size of the input array" );
121 118 }
119 +
122 120 return true;
123 121 }
124 122
125 123 /**
126 - * Get the front-most item on the queue without checking it out.
124 + * Get the front-most items on the queue without checking them out.
127 125 *
128 126 * @param int $count Number of items to return when looking at the items.
129 127 *
130 128 * @return array
@@ -138,8 +136,24 @@
138 136 return array();
139 137 }
140 138
141 139 /**
140 + * Get the last-added items on the queue without checking them out.
141 + *
142 + * @param int $count Number of items to return when looking at the items.
143 + *
144 + * @return array
145 + */
146 + public function peek_newest( $count = 1 ) {
147 + $items = $this->fetch_items( $count, 'DESC' );
148 + if ( $items ) {
149 + return Utils::get_item_values( $items );
150 + }
151 +
152 + return array();
153 + }
154 +
155 + /**
142 156 * Gets items with particular IDs.
143 157 *
144 158 * @param array $item_ids Array of item IDs to retrieve.
145 159 *
@@ -158,37 +172,14 @@
158 172 * Gets the queue lag.
159 173 * Lag is the difference in time between the age of the oldest item
160 174 * (aka first or frontmost item) and the current time.
161 175 *
162 - * @param microtime $now The current time in microtime.
176 + * @param float $now The current time in microtime.
163 177 *
164 - * @return float|int|mixed|null
178 + * @return float
165 179 */
166 180 public function lag( $now = null ) {
167 - global $wpdb;
168 -
169 - $first_item_name = $wpdb->get_var(
170 - $wpdb->prepare(
171 - "SELECT option_name FROM $wpdb->options WHERE option_name LIKE %s ORDER BY option_name ASC LIMIT 1",
172 - "jpsq_{$this->id}-%"
173 - )
174 - );
175 -
176 - if ( ! $first_item_name ) {
177 - return 0;
178 - }
179 -
180 - if ( null === $now ) {
181 - $now = microtime( true );
182 - }
183 -
184 - // Break apart the item name to get the timestamp.
185 - $matches = null;
186 - if ( preg_match( '/^jpsq_' . $this->id . '-(\d+\.\d+)-/', $first_item_name, $matches ) ) {
187 - return $now - (float) $matches[1];
188 - } else {
189 - return 0;
190 - }
181 + return (float) $this->queue_storage->get_lag( $now );
191 182 }
192 183
193 184 /**
194 185 * Resets the queue.
@@ -193,16 +184,11 @@
193 184 /**
194 185 * Resets the queue.
195 186 */
196 187 public function reset() {
197 - global $wpdb;
198 188 $this->delete_checkout_id();
199 - $wpdb->query(
200 - $wpdb->prepare(
201 - "DELETE FROM $wpdb->options WHERE option_name LIKE %s",
202 - "jpsq_{$this->id}-%"
203 - )
204 - );
189 +
190 + $this->queue_storage->clear_queue();
205 191 }
206 192
207 193 /**
208 194 * Return the size of the queue.
@@ -209,16 +195,9 @@
209 195 *
210 196 * @return int
211 197 */
212 198 public function size() {
213 - global $wpdb;
214 -
215 - return (int) $wpdb->get_var(
216 - $wpdb->prepare(
217 - "SELECT count(*) FROM $wpdb->options WHERE option_name LIKE %s",
218 - "jpsq_{$this->id}-%"
219 - )
220 - );
199 + return $this->queue_storage->get_item_count();
221 200 }
222 201
223 202 /**
224 203 * Lets you know if there is any items in the queue.
@@ -227,17 +206,9 @@
227 206 *
228 207 * @return bool
229 208 */
230 209 public function has_any_items() {
231 - global $wpdb;
232 - $value = $wpdb->get_var(
233 - $wpdb->prepare(
234 - "SELECT exists( SELECT option_name FROM $wpdb->options WHERE option_name LIKE %s )",
235 - "jpsq_{$this->id}-%"
236 - )
237 - );
238 -
239 - return ( '1' === $value );
210 + return $this->size() > 0;
240 211 }
241 212
242 213 /**
243 214 * Used to checkout the queue.
@@ -243,9 +214,9 @@
243 214 * Used to checkout the queue.
244 215 *
245 216 * @param int $buffer_size Size of the buffer to checkout.
246 217 *
247 - * @return Automattic\Jetpack\Sync\Queue_Buffer|bool|int|\WP_Error
218 + * @return \Automattic\Jetpack\Sync\Queue_Buffer|bool|int|\WP_Error
248 219 */
249 220 public function checkout( $buffer_size ) {
250 221 if ( $this->get_checkout_id() ) {
251 222 return new WP_Error( 'unclosed_buffer', 'There is an unclosed buffer' );
@@ -250,9 +221,10 @@
250 221 if ( $this->get_checkout_id() ) {
251 222 return new WP_Error( 'unclosed_buffer', 'There is an unclosed buffer' );
252 223 }
253 224
254 - $buffer_id = uniqid();
225 + // TODO check if adding a prefix is going to be a problem
226 + $buffer_id = uniqid( '', true );
255 227
256 228 $result = $this->set_checkout_id( $buffer_id );
257 229
258 230 if ( ! $result || is_wp_error( $result ) ) {
@@ -260,15 +232,17 @@
260 232 }
261 233
262 234 $items = $this->fetch_items( $buffer_size );
263 235
236 + if ( ! is_countable( $items ) ) {
237 + return false;
238 + }
239 +
264 240 if ( count( $items ) === 0 ) {
265 241 return false;
266 242 }
267 243
268 - $buffer = new Queue_Buffer( $buffer_id, array_slice( $items, 0, $buffer_size ) );
269 -
270 - return $buffer;
244 + return new Queue_Buffer( $buffer_id, array_slice( $items, 0, $buffer_size ) );
271 245 }
272 246
273 247 /**
274 248 * Given a list of items return the items ids.
@@ -286,11 +260,11 @@
286 260 );
287 261 }
288 262
289 263 /**
290 - * Pop elements from the queue.
264 + * Remove the oldest items from the queue.
291 265 *
292 - * @param int $limit Number of items to pop from the queue.
266 + * @param int $limit Number of items to remove from the queue.
293 267 *
294 268 * @return array|object|null
295 269 */
296 270 public function pop( $limit ) {
@@ -295,8 +269,12 @@
295 269 */
296 270 public function pop( $limit ) {
297 271 $items = $this->fetch_items( $limit );
298 272
273 + if ( ! $items ) {
274 + return array();
275 + }
276 +
299 277 $ids = $this->get_ids( $items );
300 278
301 279 $this->delete( $ids );
302 280
@@ -303,8 +281,29 @@
303 281 return $items;
304 282 }
305 283
306 284 /**
285 + * Remove the newest items from the queue.
286 + *
287 + * @param int $limit Number of items to remove from the queue.
288 + *
289 + * @return array|object|null
290 + */
291 + public function pop_newest( $limit ) {
292 + $items = $this->fetch_items( $limit, 'DESC' );
293 +
294 + if ( ! $items ) {
295 + return array();
296 + }
297 +
298 + $ids = $this->get_ids( $items );
299 +
300 + $this->delete( $ids );
301 +
302 + return $items;
303 + }
304 +
305 + /**
307 306 * Get the items from the queue with a memory limit.
308 307 *
309 308 * This checks out rows until it either empties the queue or hits a certain memory limit
310 309 * it loads the sizes from the DB first so that it doesn't accidentally
@@ -314,9 +313,9 @@
314 313 *
315 314 * @param int $max_memory (bytes) Maximum memory threshold.
316 315 * @param int $max_buffer_size Maximum buffer size (number of items).
317 316 *
318 - * @return Automattic\Jetpack\Sync\Queue_Buffer|bool|int|\WP_Error
317 + * @return \Automattic\Jetpack\Sync\Queue_Buffer|bool|int|\WP_Error
319 318 */
320 319 public function checkout_with_memory_limit( $max_memory, $max_buffer_size = 500 ) {
321 320 if ( $this->get_checkout_id() ) {
322 321 return new WP_Error( 'unclosed_buffer', 'There is an unclosed buffer' );
@@ -321,9 +320,9 @@
321 320 if ( $this->get_checkout_id() ) {
322 321 return new WP_Error( 'unclosed_buffer', 'There is an unclosed buffer' );
323 322 }
324 323
325 - $buffer_id = uniqid();
324 + $buffer_id = uniqid( '', true );
326 325
327 326 $result = $this->set_checkout_id( $buffer_id );
328 327
329 328 if ( ! $result || is_wp_error( $result ) ) {
@@ -329,29 +328,24 @@
329 328 if ( ! $result || is_wp_error( $result ) ) {
330 329 return $result;
331 330 }
332 331
333 - // Get the map of buffer_id -> memory_size.
334 - global $wpdb;
332 + // How much memory is currently being used by the items.
333 + $total_memory = 0;
335 334
336 - $items_with_size = $wpdb->get_results(
337 - $wpdb->prepare(
338 - "SELECT option_name AS id, LENGTH(option_value) AS value_size FROM $wpdb->options WHERE option_name LIKE %s ORDER BY option_name ASC LIMIT %d",
339 - "jpsq_{$this->id}-%",
340 - $max_buffer_size
341 - ),
342 - OBJECT
343 - );
335 + // Store the items to return
336 + $items = array();
344 337
345 - if ( count( $items_with_size ) === 0 ) {
338 + $current_items_ids = $this->queue_storage->get_items_ids_with_size( $max_buffer_size - count( $items ) );
339 +
340 + // If no valid items are returned or no items are returned, continue.
341 + if ( ! is_countable( $current_items_ids ) || count( $current_items_ids ) === 0 ) {
346 342 return false;
347 343 }
348 344
349 - $total_memory = 0;
350 - $max_item_id = $items_with_size[0]->id;
351 - $min_item_id = $max_item_id;
345 + $item_ids_to_fetch = array();
352 346
353 - foreach ( $items_with_size as $id => $item_with_size ) {
347 + foreach ( $current_items_ids as $id => $item_with_size ) {
354 348 $total_memory += $item_with_size->value_size;
355 349
356 350 // If this is the first item and it exceeds memory, allow loop to continue
357 351 // we will exit on the next iteration instead.
@@ -358,22 +352,29 @@
358 352 if ( $total_memory > $max_memory && $id > 0 ) {
359 353 break;
360 354 }
361 355
362 - $max_item_id = $item_with_size->id;
356 + $item_ids_to_fetch[] = $item_with_size->id;
363 357 }
364 358
365 - $query = $wpdb->prepare(
366 - "SELECT option_name AS id, option_value AS value FROM $wpdb->options WHERE option_name >= %s and option_name <= %s ORDER BY option_name ASC",
367 - $min_item_id,
368 - $max_item_id
369 - );
359 + $current_items = $this->queue_storage->fetch_items_by_ids( $item_ids_to_fetch );
370 360
371 - $items = $wpdb->get_results( $query, OBJECT ); // phpcs:ignore WordPress.DB.PreparedSQL.NotPrepared
372 - foreach ( $items as $item ) {
373 - // @codingStandardsIgnoreStart
374 - $item->value = @unserialize( $item->value );
375 - // @codingStandardsIgnoreEnd
361 + $items_count = is_countable( $current_items ) ? count( $current_items ) : 0;
362 +
363 + if ( $items_count > 0 ) {
364 + /**
365 + * Save some memory by moving things one by one to the array of items being returned, instead of
366 + * unserializing all and then merging them with other items.
367 + *
368 + * PHPCS ignore is because this is the expected behavior - we're assigning a variable in the condition part of the loop.
369 + */
370 + // phpcs:ignore Generic.CodeAnalysis.AssignmentInCondition.FoundInWhileCondition
371 + while ( ( $current_item = array_shift( $current_items ) ) !== null ) {
372 + // @codingStandardsIgnoreStart
373 + $current_item->value = @unserialize( $current_item->value );
374 + // @codingStandardsIgnoreEnd
375 + $items[] = $current_item;
376 + }
376 377 }
377 378
378 379 if ( count( $items ) === 0 ) {
379 380 $this->delete_checkout_id();
@@ -380,17 +381,15 @@
380 381
381 382 return false;
382 383 }
383 384
384 - $buffer = new Queue_Buffer( $buffer_id, $items );
385 -
386 - return $buffer;
385 + return new Queue_Buffer( $buffer_id, $items );
387 386 }
388 387
389 388 /**
390 389 * Check in the queue.
391 390 *
392 - * @param Automattic\Jetpack\Sync\Queue_Buffer $buffer Queue_Buffer object.
391 + * @param \Automattic\Jetpack\Sync\Queue_Buffer $buffer Queue_Buffer object.
393 392 *
394 393 * @return bool|\WP_Error
395 394 */
396 395 public function checkin( $buffer ) {
@@ -407,10 +406,10 @@
407 406
408 407 /**
409 408 * Close the buffer.
410 409 *
411 - * @param Automattic\Jetpack\Sync\Queue_Buffer $buffer Queue_Buffer object.
412 - * @param null|array $ids_to_remove Ids to remove from the queue.
410 + * @param \Automattic\Jetpack\Sync\Queue_Buffer $buffer Queue_Buffer object.
411 + * @param null|array $ids_to_remove Ids to remove from the queue.
413 412 *
414 413 * @return bool|\WP_Error
415 414 */
416 415 public function close( $buffer, $ids_to_remove = null ) {
@@ -444,16 +443,15 @@
444 443 *
445 444 * @return bool|int
446 445 */
447 446 private function delete( $ids ) {
448 - if ( 0 === count( $ids ) ) {
447 + if ( array() === $ids ) {
449 448 return 0;
450 449 }
451 - global $wpdb;
452 - $sql = "DELETE FROM $wpdb->options WHERE option_name IN (" . implode( ', ', array_fill( 0, count( $ids ), '%s' ) ) . ')';
453 - $query = call_user_func_array( array( $wpdb, 'prepare' ), array_merge( array( $sql ), $ids ) );
454 450
455 - return $wpdb->query( $query ); // phpcs:ignore WordPress.DB.PreparedSQL.NotPrepared
451 + $this->queue_storage->delete_items_by_ids( $ids );
452 +
453 + return true;
456 454 }
457 455
458 456 /**
459 457 * Flushes all items from the queue.
@@ -562,11 +560,14 @@
562 560 )
563 561 );
564 562
565 563 if ( $checkout_value ) {
566 - list( $checkout_id, $timestamp ) = explode( ':', $checkout_value );
567 - if ( (int) $timestamp > time() ) {
568 - return $checkout_id;
564 + $parts = explode( ':', $checkout_value, 2 );
565 + if ( count( $parts ) === 2 ) {
566 + list( $checkout_id, $timestamp ) = $parts;
567 + if ( (int) $timestamp > time() ) {
568 + return $checkout_id;
569 + }
569 570 }
570 571 }
571 572
572 573 return false;
@@ -651,32 +652,17 @@
651 652 /**
652 653 * Return the items in the queue.
653 654 *
654 655 * @param null|int $limit Limit to the number of items we fetch at once.
656 + * @param string $order Sort direction for the items. Accepts 'ASC' or 'DESC'.
657 + * Any other value will be treated as 'ASC'.
655 658 *
656 659 * @return array|object|null
657 660 */
658 - private function fetch_items( $limit = null ) {
659 - global $wpdb;
661 + private function fetch_items( $limit = null, $order = 'ASC' ) {
662 + $order = 'DESC' === $order ? 'DESC' : 'ASC';
660 663
661 - if ( $limit ) {
662 - $items = $wpdb->get_results(
663 - $wpdb->prepare(
664 - "SELECT option_name AS id, option_value AS value FROM $wpdb->options WHERE option_name LIKE %s ORDER BY option_name ASC LIMIT %d",
665 - "jpsq_{$this->id}-%",
666 - $limit
667 - ),
668 - OBJECT
669 - );
670 - } else {
671 - $items = $wpdb->get_results(
672 - $wpdb->prepare(
673 - "SELECT option_name AS id, option_value AS value FROM $wpdb->options WHERE option_name LIKE %s ORDER BY option_name ASC",
674 - "jpsq_{$this->id}-%"
675 - ),
676 - OBJECT
677 - );
678 - }
664 + $items = $this->queue_storage->fetch_items( $limit, $order );
679 665
680 666 return $this->unserialize_values( $items );
681 667 }
682 668
@@ -687,28 +673,9 @@
687 673 *
688 674 * @return array|object|null
689 675 */
690 676 private function fetch_items_by_id( $items_ids ) {
691 - global $wpdb;
692 -
693 - // return early if $items_ids is empty or not an array.
694 - if ( empty( $items_ids ) || ! is_array( $items_ids ) ) {
695 - return null;
696 - }
697 -
698 - $ids_placeholders = implode( ', ', array_fill( 0, count( $items_ids ), '%s' ) );
699 - $query_with_placeholders = "SELECT option_name AS id, option_value AS value
700 - FROM $wpdb->options
701 - WHERE option_name IN ( $ids_placeholders )";
702 - $items = $wpdb->get_results(
703 - $wpdb->prepare(
704 - $query_with_placeholders, // phpcs:ignore WordPress.DB.PreparedSQL.NotPrepared
705 - $items_ids
706 - ),
707 - OBJECT
708 - );
709 -
710 - return $this->unserialize_values( $items );
677 + return $this->unserialize_values( $this->queue_storage->fetch_items_by_ids( $items_ids ) );
711 678 }
712 679
713 680 /**
714 681 * Unserialize item values.
@@ -732,9 +699,9 @@
732 699
733 700 /**
734 701 * Return true if the buffer is still valid or an Error other wise.
735 702 *
736 - * @param Automattic\Jetpack\Sync\Queue_Buffer $buffer The Queue_Buffer.
703 + * @param \Automattic\Jetpack\Sync\Queue_Buffer $buffer The Queue_Buffer.
737 704 *
738 705 * @return bool|WP_Error
739 706 */
740 707 private function validate_checkout( $buffer ) {