save()` on a CPT store fires it three times. Every firing used to * append a row (an fsync, ~3 ms) after a three-query `wc_get_order()` — 12 * rows and ~66 ms for ONE online order, measured 2026-09-03 on dev-next. * A journal row is a change POINTER (ADR 0033), so one row per order per * request carries the same information. * * Single slot, not a map: a save for a DIFFERENT order flushes the pending * one first (so a bulk loop never holds rows until process end), which * means at most one order is ever pending. Static, not per instance: the * "update row lands before any other-origin row" guarantee must hold for * whichever `Sync_Journal` instance writes the other row. The slot keeps the * blog id so a multisite `switch_to_blog()` between save and flush still * writes to the originating site's table, and the order object the hook * handed us so the flush never refetches. * * Rows land on {@see flush_pending_order_updates()}: at `shutdown` (last, * after WooCommerce's own shutdown saves), before any other-origin row, or * when a different order is saved. Once the shutdown flush has run, later * updates write immediately. * * @var array{blog: int, id: int, order: \WC_Abstract_Order|null}|null */ private static ?array $pending_order_update = null; /** Set by the shutdown flush; afterwards updates are written immediately. */ private static bool $shutdown_flushed = false; /** * Option-name prefix for the per-object-type lossy-prune watermarks. * * The watermark is scoped per object type for the same reason heads are * stream-scoped: the streams share one AUTO_INCREMENT space, so a single * global watermark advanced by an order prune would sit far above a quiet * catalogue stream's head. Every catalogue cursor would then read as * "below the horizon" and rebaseline on EVERY poll, forever — and * symmetrically for the order lane. A stream's horizon may only move for * rows that stream can serve. */ const PRUNE_WATERMARK_OPTION_PREFIX = 'wcpos_change_log_prune_watermark_'; public function table_name(): string { global $wpdb; return $wpdb->prefix . Health::SYNC_JOURNAL_TABLE; } /** * The `revision` column is a per-lane union: a `date_modified` stamp for * catalogue/customer rows, `''` for live order rows (order revisions are * computed at pull time — ADR 0033), `'deleted'` for order tombstones, and * legacy pre-#1746 order rows may still carry stored `sha256:` hashes, * which the pull planner's stored-wins branch serves until they age out. */ public function schema_sql( string $table_name, string $charset_collate = '' ): string { return "CREATE TABLE {$table_name} (\n" . " sequence BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,\n" . " object_type VARCHAR(20) NOT NULL,\n" . " object_id BIGINT UNSIGNED NOT NULL,\n" . " deleted TINYINT(1) NOT NULL DEFAULT 0,\n" . " revision VARCHAR(80) NOT NULL DEFAULT '',\n" . " modified_gmt DATETIME NOT NULL,\n" . " origin VARCHAR(40) NOT NULL DEFAULT 'hook',\n" . " created_gmt DATETIME NOT NULL,\n" . " PRIMARY KEY (sequence),\n" . " KEY type_sequence (object_type, sequence),\n" . " KEY type_object (object_type, object_id)\n" . ") {$charset_collate};"; } public function install(): void { global $wpdb; $table_name = $this->table_name(); $table_existed = Health::table_exists( $table_name ); if ( ! function_exists( 'dbDelta' ) ) { require_once ABSPATH . 'wp-admin/includes/upgrade.php'; } dbDelta( $this->schema_sql( $table_name, $wpdb->get_charset_collate() ) ); if ( ! $table_existed && Health::table_exists( $table_name ) ) { self::reset_prune_watermarks(); } if ( ! Health::table_exists( $table_name ) ) { return; } // The epoch marks a SEQUENCE GENERATION. Regenerate it only when the // table was actually (re)created — dbDelta on a surviving table is a // no-op and every row survives, so activation/upgrade re-runs must not // force every client (both lanes, since the epoch is journal-global) // into a needless resync-from-zero. if ( ! $table_existed ) { $this->regenerate_epoch(); } $backfill = $this->backfill_status(); $backfill_has_progress = 'idle' !== $backfill['status'] || $backfill['nextPage'] > 1 || $backfill['processed'] > 0; if ( 0 === $this->head_sequence() && $backfill_has_progress ) { delete_option( self::BACKFILL_OPTION ); } // A store with no orders has nothing to backfill: mark it complete on a // fresh table so sequence-zero pulls are journal-authoritative from the // first write (Order_Query holds the baseline on the modified-date scan // until the backfill is complete). Stores WITH history stay incomplete // until the admin backfill runs. Runs after the stale-cursor cleanup // above so a carried-over cursor cannot masquerade as completion. if ( ! $table_existed && function_exists( 'wc_get_orders' ) ) { $existing = wc_get_orders( array( 'type' => 'shop_order', 'limit' => 1, 'return' => 'ids', ) ); if ( is_array( $existing ) && array() === $existing ) { update_option( self::BACKFILL_OPTION, array( 'status' => 'complete', 'nextPage' => 1, 'pageSize' => null, 'processed' => 0, 'lastOrderId' => 0, 'lastRunGmt' => gmdate( 'c' ), ), false ); } } } /** Return the stable id for the current sequence generation. */ public function ensure_epoch(): string { $epoch = (string) get_option( self::EPOCH_OPTION, '' ); if ( '' !== $epoch ) { return $epoch; } $epoch = $this->mint_epoch(); add_option( self::EPOCH_OPTION, $epoch, '', true ); return (string) get_option( self::EPOCH_OPTION, $epoch ); } /** Force a new id for the current sequence generation. */ public function regenerate_epoch(): string { $epoch = $this->mint_epoch(); update_option( self::EPOCH_OPTION, $epoch, true ); return $epoch; } private function mint_epoch(): string { return function_exists( 'wp_generate_uuid4' ) ? wp_generate_uuid4() : md5( uniqid( 'wcpos-epoch', true ) ); } /** Append a schema-upgrade customer row for every live user. */ public function append_customer_updates_for_all_users(): bool { global $wpdb; $now = gmdate( 'Y-m-d H:i:s' ); return false !== $wpdb->query( $wpdb->prepare( 'INSERT INTO ' . $this->table_name() . ' (object_type, object_id, deleted, revision, modified_gmt, origin, created_gmt)' . " SELECT 'customer', ID, 0, '', %s, 'schema-upgrade', %s FROM " . $wpdb->users, $now, $now ) ); } /** * Append one tombstone per catalogue post id, in a single statement. * * The per-record `record_post_deleted()` path loads a `WC_Product` for the revision stamp, which * is fine for the handful of records one settings write moves and hopeless for the whole hidden * set of a store that keeps thousands of products online-only. This is the bulk form, shaped like * `append_customer_updates_for_all_users()`: one INSERT ... SELECT, no revision (a tombstone * carries no state the client compares), and the post type read from `wp_posts` so a stale id in * the merchant's list cannot announce a change to an unrelated record. * * @param int[] $ids Product / variation post ids. */ public function append_catalogue_tombstones( array $ids ): bool { global $wpdb; $ids = array_values( array_unique( array_filter( array_map( 'intval', $ids ), static function ( int $id ): bool { return $id > 0; } ) ) ); if ( array() === $ids ) { return true; } $now = gmdate( 'Y-m-d H:i:s' ); $placeholders = implode( ',', array_fill( 0, count( $ids ), '%d' ) ); return false !== $wpdb->query( $wpdb->prepare( 'INSERT INTO ' . $this->table_name() . ' (object_type, object_id, deleted, revision, modified_gmt, origin, created_gmt)' . " SELECT CASE p.post_type WHEN 'product_variation' THEN 'variation' ELSE 'product' END," . " p.ID, 1, '', %s, 'visibility-seed', %s" . " FROM {$wpdb->posts} p" . " WHERE p.post_type IN ('product','product_variation') AND p.ID IN ({$placeholders})" // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- %d placeholder list generated from count(); the ids are bound below. . ' ORDER BY p.ID', $now, $now, ...$ids ) ); } public function register_hooks(): void { add_action( 'woocommerce_new_product', array( $this, 'record_product_created' ), 10, 1 ); add_action( 'woocommerce_update_product', array( $this, 'record_product_updated' ), 10, 1 ); add_action( 'woocommerce_new_product_variation', array( $this, 'record_variation_created' ), 10, 1 ); add_action( 'woocommerce_update_product_variation', array( $this, 'record_variation_updated' ), 10, 1 ); add_action( 'woocommerce_new_coupon', array( $this, 'record_coupon_created' ), 10, 1 ); add_action( 'woocommerce_update_coupon', array( $this, 'record_coupon_updated' ), 10, 1 ); add_action( 'wp_trash_post', array( $this, 'record_post_deleted' ), 10, 1 ); add_action( 'before_delete_post', array( $this, 'record_post_deleted' ), 10, 1 ); add_action( 'untrashed_post', array( $this, 'record_post_untrashed' ), 10, 1 ); add_action( 'woocommerce_tax_rate_added', array( $this, 'record_tax_rate_created' ), 10, 1 ); add_action( 'woocommerce_tax_rate_updated', array( $this, 'record_tax_rate_updated' ), 10, 1 ); add_action( 'woocommerce_tax_rate_deleted', array( $this, 'record_tax_rate_deleted' ), 10, 1 ); add_action( 'created_term', array( $this, 'record_term_created' ), 10, 3 ); add_action( 'edited_term', array( $this, 'record_term_edited' ), 10, 3 ); add_action( 'delete_term', array( $this, 'record_term_deleted' ), 10, 3 ); add_action( 'added_term_meta', array( $this, 'record_term_meta_change' ), 10, 2 ); add_action( 'updated_term_meta', array( $this, 'record_term_meta_change' ), 10, 2 ); add_action( 'deleted_term_meta', array( $this, 'record_term_meta_deleted' ), 10, 2 ); add_action( 'user_register', array( $this, 'record_customer_created' ), 10, 1 ); add_action( 'woocommerce_created_customer', array( $this, 'record_customer_created_persisted' ), 10, 1 ); add_action( 'woocommerce_new_customer', array( $this, 'record_customer_created_persisted' ), 10, 1 ); add_action( 'profile_update', array( $this, 'record_customer_profile_update' ), 10, 2 ); add_action( 'set_user_role', array( $this, 'record_customer_role_change' ), 10, 3 ); add_action( 'add_user_role', array( $this, 'record_customer_role_added' ), 10, 2 ); add_action( 'remove_user_role', array( $this, 'record_customer_role_removed' ), 10, 2 ); add_action( 'woocommerce_update_customer', array( $this, 'record_customer_updated' ), 10, 1 ); add_action( 'delete_user', array( $this, 'record_customer_deleted' ), 10, 1 ); add_action( 'woocommerce_new_order', array( $this, 'record_order_created' ), 10, 1 ); // Two args: the data store passes ($order_id, $order). Keeping the object // lets the coalesced flush read modified_gmt without a refetch. add_action( 'woocommerce_update_order', array( $this, 'record_order_updated' ), 10, 2 ); // Request boundary for the coalesced order update row. LAST on shutdown: // WooCommerce saves the customer at 10 and the session at 20, and any // save those trigger must still find the slot open. Zero accepted args: // do_action( 'shutdown' ) passes an empty string otherwise. add_action( 'shutdown', array( $this, 'flush_pending_order_updates_at_shutdown' ), PHP_INT_MAX, 0 ); add_action( 'woocommerce_before_trash_order', array( $this, 'record_order_deleted' ), 10, 1 ); add_action( 'woocommerce_before_delete_order', array( $this, 'record_order_deleted' ), 10, 1 ); add_action( 'woocommerce_untrash_order', array( $this, 'record_cot_order_untrashed' ), 10, 1 ); add_action( 'woocommerce_pos_invalidate', array( $this, 'record_invalidation' ), 10, 2 ); } /** * Record an out-of-band change announced by an extension. * * Plugins fire `woocommerce_pos_invalidate` when they change a record's * SERVED representation in a way no save hook announces — a filter-only * output change (a pricing filter, an added payload field). The journal * appends a pointer row; clients hydrate pointer rows by sequence, so the * re-served payload carries the plugin's change. Formula fingerprints * (#1742) will eventually make representation changes directly detectable; * until then this action is the documented relief valve. * * `$object_type` is the registry's SINGULAR journal name: `product`, * `variation`, `customer`, `order`, `tax_rate`, and the other journalled * catalogue types. A plural (`products`) or unknown type is logged and * ignored. Rows land with origin `invalidate` on every type. * * @since 1.10.3 * * @param string $object_type Canonical (singular) journal object type. * @param int $object_id Changed object ID. */ public function record_invalidation( $object_type = '', $object_id = 0 ): void { // Loose signature on purpose: a public action handler whose posture is // log-and-ignore — a one-arg or wrong-typed do_action() must not fatal // the calling plugin's request. $object_type = is_scalar( $object_type ) ? (string) $object_type : ''; $object_id = is_scalar( $object_id ) ? (int) $object_id : 0; $collection = Collections::by_object_type( $object_type ); if ( $object_id <= 0 || null === $collection || ! isset( $collection['journal'] ) ) { Logger::log( sprintf( 'WCPOS sync: ignored invalidation for object_type "%s" (id %d)', $object_type, $object_id ) ); return; } if ( 'order' === $object_type ) { $this->record_order_change( $object_id, 'invalidate', false ); return; } $loader = (string) ( $collection['identity']['loader'] ?? '' ); if ( 'product' === $loader ) { $object = function_exists( 'wc_get_product' ) ? wc_get_product( $object_id ) : null; $this->record( $object_type, $object_id, false, self::object_revision( $object ), 'invalidate' ); if ( 'variation' === $object_type ) { // Native variation paths always pair the parent row — the parent // document carries the variable price range — so an invalidation // must too, or the relief valve half-works. Recorded inline (not via // record_variation_parent) so the paired row keeps the 'invalidate' // origin the contract above promises for every row this action lands. $parent_id = function_exists( 'wp_get_post_parent_id' ) ? (int) wp_get_post_parent_id( $object_id ) : 0; if ( $parent_id > 0 ) { $parent = function_exists( 'wc_get_product' ) ? wc_get_product( $parent_id ) : null; $this->record( 'product', $parent_id, false, self::object_revision( $parent ), 'invalidate' ); } } return; } if ( 'customer' === $loader ) { try { $customer = class_exists( '\\WC_Customer' ) ? new \WC_Customer( $object_id ) : null; } catch ( \Exception $e ) { Logger::log( sprintf( 'WCPOS sync: ignored invalidation for missing customer %d', $object_id ) ); return; } $this->record( 'customer', $object_id, false, self::object_revision( $customer ), 'invalidate', true, 'invalidate' ); return; } $this->record( $object_type, $object_id, false, '', 'invalidate' ); } public function record_product_created( int $product_id ): void { $this->record_catalogue_object( 'product', $product_id, false ); } public function record_product_updated( int $product_id ): void { $this->record_catalogue_object( 'product', $product_id, false ); } public function record_variation_created( int $variation_id ): void { $this->record_catalogue_object( 'variation', $variation_id, false ); $this->record_variation_parent( $variation_id ); } public function record_variation_updated( int $variation_id ): void { $this->record_catalogue_object( 'variation', $variation_id, false ); $this->record_variation_parent( $variation_id ); } /** Record the representation change to a variation's parent product. */ private function record_variation_parent( int $variation_id ): void { $parent_id = function_exists( 'wp_get_post_parent_id' ) ? (int) wp_get_post_parent_id( $variation_id ) : 0; if ( $parent_id > 0 ) { $this->record_catalogue_object( 'product', $parent_id, false ); } } public function record_post_deleted( int $post_id ): void { $post_type = get_post_type( $post_id ); if ( 'product' === $post_type ) { $this->record_catalogue_object( 'product', $post_id, true ); return; } if ( 'product_variation' === $post_type ) { $this->record_catalogue_object( 'variation', $post_id, true ); $this->record_variation_parent( $post_id ); return; } if ( 'shop_coupon' === $post_type ) { $this->record( 'coupon', $post_id, true ); return; } if ( 'shop_order' === $post_type ) { $this->record_order_deleted( $post_id ); } } public function record_post_untrashed( int $post_id ): void { $post_type = get_post_type( $post_id ); if ( 'product' === $post_type ) { $this->record_product_updated( $post_id ); return; } if ( 'product_variation' === $post_type ) { $this->record_variation_updated( $post_id ); return; } if ( 'shop_coupon' === $post_type ) { $this->record_coupon_updated( $post_id ); return; } if ( 'shop_order' === $post_type ) { $this->record_order_untrashed( $post_id ); } } public function record_coupon_created( int $coupon_id ): void { $this->record( 'coupon', $coupon_id, false ); } public function record_coupon_updated( int $coupon_id ): void { $this->record( 'coupon', $coupon_id, false ); } /** Map tracked product taxonomies to journal object types. */ private static function term_taxonomy_object_types(): array { $map = array(); foreach ( Collections::with( 'identity' ) as $row ) { if ( isset( $row['identity']['taxonomy'] ) ) { $map[ $row['identity']['taxonomy'] ] = $row['object_type']; } } return $map; } public function record_term_created( int $term_id, int $tt_id, string $taxonomy ): void { $this->record_term_change( $term_id, $taxonomy, 'create' ); } public function record_term_edited( int $term_id, int $tt_id, string $taxonomy ): void { $this->record_term_change( $term_id, $taxonomy, 'update' ); } public function record_term_deleted( int $term_id, int $tt_id, string $taxonomy ): void { $this->record_term_change( $term_id, $taxonomy, 'delete' ); } private function record_term_change( int $term_id, string $taxonomy, string $change_type ): void { $object_type = self::term_taxonomy_object_types()[ $taxonomy ] ?? null; if ( null !== $object_type ) { $this->record( $object_type, $term_id, 'delete' === $change_type ); } } /** Record added or updated metadata for a tracked term. */ public function record_term_meta_change( int $meta_id, int $term_id ): void { $this->record_term_representation_change( $term_id ); } /** Record deleted metadata for a tracked term. */ public function record_term_meta_deleted( array $meta_ids, int $term_id ): void { $this->record_term_representation_change( $term_id ); } /** Resolve a meta hook's term and record it only when tracked. */ private function record_term_representation_change( int $term_id ): void { $term = get_term( $term_id ); if ( is_object( $term ) && isset( $term->taxonomy ) ) { $this->record_term_change( $term_id, (string) $term->taxonomy, 'update' ); } } public function record_tax_rate_created( int $tax_rate_id ): void { $this->record( 'tax_rate', $tax_rate_id, false ); } public function record_tax_rate_updated( int $tax_rate_id ): void { $this->record( 'tax_rate', $tax_rate_id, false ); } public function record_tax_rate_deleted( int $tax_rate_id ): void { $this->record( 'tax_rate', $tax_rate_id, true ); } public function record_customer_created( int $customer_id ): void { $this->record_customer( $customer_id, false, true, 'create' ); } /** WooCommerce create hooks for any user — definitive post-persist create with persisted dedup. */ public function record_customer_created_persisted( int $customer_id ): void { $this->record_customer( $customer_id, false, false, 'create' ); } /** Record the definitive post-persist customer update. */ public function record_customer_updated( int $customer_id ): void { $this->record_customer( $customer_id, false, false ); } /** Record a WordPress profile update as a customer update. */ public function record_customer_profile_update( int $user_id, $old_user_data = null ): void { $this->record_customer( $user_id, false ); } /** Record a customer role replacement. */ public function record_customer_role_change( int $user_id, $role = '', $old_roles = array() ): void { $this->record_customer( $user_id, false ); } /** add_user_role handler — any added role changes the served customer record. */ public function record_customer_role_added( int $user_id, $role = '' ): void { $this->record_customer( $user_id, false ); } /** Record a removed customer role as a present update. */ public function record_customer_role_removed( int $user_id, $role = '' ): void { $this->record_customer( $user_id, false ); } /** delete_user handler — deletion removes any user from the POS customer space. */ public function record_customer_deleted( int $customer_id ): void { $this->record_customer( $customer_id, true, true, 'delete' ); } public function record_order_created( int $order_id ): void { $this->record_order_change( $order_id, 'hook:create', false ); } /** * Mark an order's `hook:update` row as owed; the row lands on flush. * * See {@see $pending_order_updates} for why this is deferred. Direct callers * that need an immediate row use {@see record_order_change()}. * * @param int $order_id Order id from the hook. * @param \WC_Abstract_Order|mixed $order Order object from the hook (second * argument of `woocommerce_update_order`), * or anything else to fall back to a * refetch at flush time. */ public function record_order_updated( int $order_id, $order = null ): void { $order = $order instanceof \WC_Abstract_Order ? $order : null; if ( self::$shutdown_flushed ) { // The request boundary has passed (a save triggered by another // shutdown handler): nothing will flush again, so write now. $this->record_order_change( $order_id, 'hook:update', false, $order ); return; } $blog = get_current_blog_id(); $slot = self::$pending_order_update; if ( null !== $slot && ( $slot['id'] !== $order_id || $slot['blog'] !== $blog ) ) { // A different order began: land what is owed so a bulk loop (WP-CLI // import, Action Scheduler runner) never holds rows until process end. $this->flush_pending_order_updates(); $slot = null; } self::$pending_order_update = array( 'blog' => $blog, 'id' => $order_id, 'order' => $order ?? ( $slot['order'] ?? null ), ); } /** * Write the owed `hook:update` row, if any. * * Called from {@see record_order_change()} before any other-origin row and * from the shutdown flush. Safe to call repeatedly: a flushed order is no * longer pending. */ public function flush_pending_order_updates(): void { $slot = self::$pending_order_update; if ( null === $slot ) { return; } self::$pending_order_update = null; self::in_blog( $slot['blog'], function () use ( $slot ): void { $this->record_order_change( $slot['id'], 'hook:update', false, $slot['order'] ); } ); } /** * The `shutdown` callback: flush, then write every later update immediately. */ public function flush_pending_order_updates_at_shutdown(): void { self::$shutdown_flushed = true; $this->flush_pending_order_updates(); } /** * Discard per-request coalescing state. Tests only: the PHPUnit process * never reaches `shutdown`, so the static slot and flag would leak between * test cases otherwise. * * @internal */ public static function reset_request_state(): void { self::$pending_order_update = null; self::$shutdown_flushed = false; } /** * Run a write under the blog it was recorded on. * * The journal table is blog-scoped, so a deferred write must not follow a * `switch_to_blog()` that happened between the save and the flush. * * @param int $blog_id Blog the write belongs to. * @param callable $write The write. */ private static function in_blog( int $blog_id, callable $write ): void { $switch = is_multisite() && get_current_blog_id() !== $blog_id; if ( $switch ) { switch_to_blog( $blog_id ); } try { $write(); } finally { if ( $switch ) { restore_current_blog(); } } } public function record_order_deleted( int $order_id ): void { $this->record_order_change( $order_id, 'hook:delete', true ); } public function record_order_untrashed( int $order_id ): void { $this->record_order_change( $order_id, 'hook:untrash', false ); } /** * Record an HPOS order's restore once the status change has settled. * * `woocommerce_untrash_order` fires BEFORE the data store restores the * status, so the row cannot be written there. The restore then performs * MORE THAN ONE object save, so the journal row's modified_gmt must be read * from the SETTLED order for checkpoint ordering. The revision is computed * at pull time rather than stored here. * * Measured sequence for an HPOS untrash (status read from wc_orders): * * woocommerce_untrash_order stored=trash * after_order_object_save object=pending stored=wc-pending * after_order_object_save object=pending stored=wc-pending * woocommerce_order_status_changed stored=wc-pending from=trash * * `woocommerce_order_status_changed` fires once, last, with the stored * status settled — so observe that instead. CPT orders never reach here: * their restore fires only `untrashed_post` (see record_post_untrashed). * * @param int $order_id Order being restored. */ public function record_cot_order_untrashed( int $order_id ): void { $handler = function ( $id, $from ) use ( $order_id, &$handler ): void { if ( (int) $id !== $order_id || 'trash' !== $from ) { return; } remove_action( 'woocommerce_order_status_changed', $handler ); $this->record_order_untrashed( $order_id ); }; add_action( 'woocommerce_order_status_changed', $handler, 10, 2 ); } /** * Append one order row immediately. * * @param int $order_id Order id. * @param string $origin Row origin (`hook:create`, `hook:update`, …). * @param bool $deleted Whether the row is a tombstone. * @param \WC_Abstract_Order|mixed $order The order object when the caller already holds it; * anything else triggers a refetch. * * @return bool Whether the insert succeeded. */ public function record_order_change( int $order_id, string $origin, bool $deleted, $order = null ): bool { global $wpdb; if ( 'hook:update' !== $origin ) { $slot = self::$pending_order_update; if ( 'hook:create' === $origin && null !== $slot && $order_id === $slot['id'] && get_current_blog_id() === $slot['blog'] ) { // The Store API saves a checkout-draft several times BEFORE // `woocommerce_new_order` fires. Both rows would point at the same // live record, so the create row makes the owed update row redundant. self::$pending_order_update = null; } else { // Land the owed update row FIRST so the stream never reads as // delete-then-update (a replay would resurrect a trashed order). $this->flush_pending_order_updates(); } } if ( ! $order instanceof \WC_Abstract_Order ) { $order = wc_get_order( $order_id ); } $modified_date = $order ? $order->get_date_modified() : null; $modified = $modified_date ? gmdate( 'Y-m-d H:i:s', $modified_date->getTimestamp() ) : gmdate( 'Y-m-d H:i:s' ); // Order revisions are computed at pull time from the served payload (ADR 0033, // #1746) — an order journal row is a change pointer, not a content stamp. // 'deleted' is kept for wire compatibility (it flows into served checkpoints) // and diagnostics; the planner branches on the `deleted` flag, not this value. $revision = $deleted ? 'deleted' : ''; $now = gmdate( 'Y-m-d H:i:s' ); return false !== $wpdb->insert( $this->table_name(), array( 'object_type' => 'order', 'object_id' => $order_id, 'deleted' => $deleted ? 1 : 0, 'revision' => $revision, 'modified_gmt' => $modified, 'origin' => $origin, 'created_gmt' => $now, ), array( '%s', '%d', '%d', '%s', '%s', '%s', '%s' ) ); } private function record_catalogue_object( string $object_type, int $object_id, bool $deleted ): void { $object = function_exists( 'wc_get_product' ) ? wc_get_product( $object_id ) : null; $this->record( $object_type, $object_id, $deleted, self::object_revision( $object ) ); } private function record_customer( int $customer_id, bool $deleted, bool $dedup = true, string $dedup_namespace = 'update' ): void { $customer = class_exists( '\\WC_Customer' ) ? new \WC_Customer( $customer_id ) : null; $this->record( 'customer', $customer_id, $deleted, self::object_revision( $customer ), 'hook', $dedup, $dedup_namespace ); } private static function object_revision( $object ): string { $date = is_object( $object ) && method_exists( $object, 'get_date_modified' ) ? $object->get_date_modified() : null; return $date && method_exists( $date, 'getTimestamp' ) ? gmdate( 'Y-m-d H:i:s', $date->getTimestamp() ) : ''; } public function record( string $object_type, int $object_id, bool $deleted, string $revision = '', string $origin = 'hook', bool $dedup = true, string $dedup_namespace = '' ): void { global $wpdb; $dedup_key = null; if ( 'customer' === $object_type ) { $dedup_key = ( $dedup ? '' : 'persisted:' ) . $object_id . ':' . ( '' !== $dedup_namespace ? $dedup_namespace : ( $deleted ? 'delete' : 'update' ) ); if ( isset( $this->recorded_this_request[ $dedup_key ] ) ) { if ( $dedup ) { return; // a pre-persist duplicate carries no new state } $wpdb->delete( $this->table_name(), array( 'sequence' => (int) $this->recorded_this_request[ $dedup_key ] ), array( '%d' ) ); } } $started = microtime( true ); $now = gmdate( 'Y-m-d H:i:s' ); $wpdb->insert( $this->table_name(), array( 'object_type' => $object_type, 'object_id' => $object_id, 'deleted' => $deleted ? 1 : 0, 'revision' => $revision, 'modified_gmt' => $now, 'origin' => $origin, 'created_gmt' => $now, ), array( '%s', '%d', '%d', '%s', '%s', '%s', '%s' ) ); if ( null !== $dedup_key ) { $this->recorded_this_request[ $dedup_key ] = $dedup ? true : (int) $wpdb->insert_id; } self::$request_write_ms += ( microtime( true ) - $started ) * 1000; } /** Return the persisted order-backfill cursor. */ public function backfill_status(): array { $status = get_option( self::BACKFILL_OPTION, array() ); $status = is_array( $status ) ? $status : array(); return array( 'status' => isset( $status['status'] ) ? (string) $status['status'] : 'idle', 'nextPage' => isset( $status['nextPage'] ) ? max( 1, (int) $status['nextPage'] ) : 1, 'pageSize' => isset( $status['pageSize'] ) ? max( 1, min( 250, (int) $status['pageSize'] ) ) : null, 'processed' => isset( $status['processed'] ) ? max( 0, (int) $status['processed'] ) : 0, 'lastOrderId' => isset( $status['lastOrderId'] ) ? max( 0, (int) $status['lastOrderId'] ) : 0, 'lastRunGmt' => isset( $status['lastRunGmt'] ) ? (string) $status['lastRunGmt'] : null, ); } /** Clear the full persisted order-backfill cursor. */ public function reset_backfill_state(): void { delete_option( self::BACKFILL_OPTION ); } /** Append one bounded page of existing orders to the journal. */ public function run_backfill_chunk( int $limit ): array { $requested_limit = max( 1, min( 250, $limit ) ); $status = $this->backfill_status(); if ( 'complete' === $status['status'] ) { return array_merge( $status, array( 'processedThisRun' => 0 ) ); } $page_size = null === $status['pageSize'] ? $requested_limit : (int) $status['pageSize']; $last_order_id = $status['lastOrderId']; $query_args = array( 'type' => 'shop_order', 'limit' => $page_size, 'orderby' => 'ID', 'order' => 'ASC', 'return' => 'ids', ); $posts_where = null; if ( $last_order_id > 0 ) { if ( class_exists( OrderUtil::class ) && OrderUtil::custom_orders_table_usage_is_enabled() ) { $query_args['field_query'] = array( array( 'field' => 'id', 'value' => $last_order_id, 'compare' => '>', ), ); } else { $posts_where = static function ( string $where ) use ( $last_order_id ): string { global $wpdb; return $where . $wpdb->prepare( " AND {$wpdb->posts}.ID > %d", $last_order_id ); }; add_filter( 'posts_where', $posts_where ); } } try { /** @var array|mixed $queried_ids The 'return' => 'ids' arg yields ids; the stub over-narrows to WC_Order[]. */ $queried_ids = wc_get_orders( $query_args ); } finally { if ( null !== $posts_where ) { remove_filter( 'posts_where', $posts_where ); } } $ids = is_array( $queried_ids ) ? array_map( 'absint', $queried_ids ) : array(); $processed_this_run = 0; $failed_this_run = 0; foreach ( $ids as $id ) { if ( $this->record_order_change( $id, 'backfill', false ) ) { $processed_this_run++; $last_order_id = $id; } else { $failed_this_run++; break; } } $all_writes_succeeded = 0 === $failed_this_run; $complete = $all_writes_succeeded && count( $ids ) < $page_size; $advance_page = $all_writes_succeeded && count( $ids ) === $page_size; $next_status = array( 'status' => $complete ? 'complete' : 'running', 'nextPage' => $advance_page ? $status['nextPage'] + 1 : $status['nextPage'], 'pageSize' => $page_size, 'processed' => $status['processed'] + $processed_this_run, 'lastOrderId' => $last_order_id, 'lastRunGmt' => gmdate( 'c' ), ); update_option( self::BACKFILL_OPTION, $next_status, false ); return array_merge( $next_status, array( 'processedThisRun' => $processed_this_run ) ); } /** * The catalogue pointer-stream's object types, projected from the registry — * every journal-covered collection except orders (which consume the journal * via the payload-windowed pull lane). Single source for the sequence-log * `all` stream and the purge's per-stream head protection. * * @return string[] */ public static function catalogue_object_types(): array { $types = array(); foreach ( Collections::with( 'journal' ) as $row ) { $object_type = (string) ( $row['journal']['object_type'] ?? '' ); if ( '' !== $object_type && 'order' !== $object_type ) { $types[] = $object_type; } } return $types; } /** * Head of the sequence space — the change stream's current end. * * With `$object_types`, the head is STREAM-SCOPED: the highest sequence any * row of those types holds. Every reader must serve the head of the stream * it serves — orders and catalogue share one AUTO_INCREMENT space, so the * global head moves on foreign writes. A stream-scoped head is what lets a * cursor actually reach `head` (the 304 idle condition) while the other * lane keeps writing, and is the pre-unification semantic of both lanes. * * @param string[] $object_types Empty = global head (retention clamps only). */ public function head_sequence( array $object_types = array() ): int { global $wpdb; if ( ! $this->table_available() ) { return 0; } $types = array_values( array_filter( array_map( 'strval', $object_types ), static fn( string $t ): bool => '' !== $t ) ); if ( array() === $types ) { return (int) $wpdb->get_var( 'SELECT MAX(sequence) FROM ' . $this->table_name() ); } $placeholders = implode( ',', array_fill( 0, count( $types ), '%s' ) ); return (int) $wpdb->get_var( $wpdb->prepare( 'SELECT MAX(sequence) FROM ' . $this->table_name() . " WHERE object_type IN ({$placeholders})", // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- %s placeholder list generated from count(). ...$types ) ); } public function table_available(): bool { global $wpdb; $table = $this->table_name(); return $table === $wpdb->get_var( $wpdb->prepare( 'SHOW TABLES LIKE %s', $table ) ); } /** Return one sequence page for the order pull lane. */ public function rows_after_sequence( int $sequence, int $limit, string $object_type = 'order' ): array { global $wpdb; if ( ! $this->table_available() ) { return array(); } $rows = $wpdb->get_results( $wpdb->prepare( 'SELECT sequence, object_id AS order_id, modified_gmt, revision, deleted, origin, created_gmt FROM ' . $this->table_name() . ' WHERE object_type = %s AND sequence > %d ORDER BY sequence ASC LIMIT %d', $object_type, max( 0, $sequence ), max( 1, min( 251, $limit ) ) ), ARRAY_A ); return is_array( $rows ) ? array_map( array( self::class, 'normalize_order_row' ), $rows ) : array(); } private static function normalize_order_row( array $row ): array { return array( 'sequence' => isset( $row['sequence'] ) ? (int) $row['sequence'] : 0, 'order_id' => isset( $row['order_id'] ) ? (int) $row['order_id'] : 0, 'modified_gmt' => isset( $row['modified_gmt'] ) ? (string) $row['modified_gmt'] : gmdate( 'Y-m-d H:i:s' ), 'revision' => isset( $row['revision'] ) ? (string) $row['revision'] : '', 'deleted' => ! empty( $row['deleted'] ) ? 1 : 0, 'origin' => isset( $row['origin'] ) ? (string) $row['origin'] : 'hook:update', 'created_gmt' => isset( $row['created_gmt'] ) ? (string) $row['created_gmt'] : gmdate( 'Y-m-d H:i:s' ), ); } /** Oldest sequence that remains available to incremental-sync clients. */ public function oldest_sequence(): int { global $wpdb; return (int) $wpdb->get_var( 'SELECT MIN(sequence) FROM ' . $this->table_name() ); } /** Delete one batch of superseded rows through the supplied sequence. */ public function compact( int $cutoff_sequence, string $cutoff_gmt, int $batch ): int { global $wpdb; if ( $cutoff_sequence <= 0 || $batch <= 0 ) { return 0; } $table = $this->table_name(); $deleted = $wpdb->query( $wpdb->prepare( 'DELETE FROM ' . $table . ' WHERE sequence <= %d AND sequence IN (' . ' SELECT sequence FROM (' . ' SELECT stale.sequence FROM ' . $table . ' stale' . ' WHERE stale.sequence <= %d' . ' AND stale.created_gmt < %s' . ' AND EXISTS (' . ' SELECT 1 FROM ' . $table . ' newer' . ' WHERE newer.object_type = stale.object_type' . ' AND newer.object_id = stale.object_id' . ' AND newer.sequence > stale.sequence' . ' ) ORDER BY stale.sequence ASC LIMIT %d' . ' ) compactable' . ' )', $cutoff_sequence, $cutoff_sequence, $cutoff_gmt, $batch ) ); return false === $deleted ? 0 : (int) $deleted; } /** * Delete one batch of expired tombstones — the log's only LOSSY deletion. * * @param int $cutoff_sequence Highest sequence this batch may remove. * @param string $cutoff_gmt Rows created before this are expired. * @param int $batch Maximum rows to remove. * @param array $object_types Restrict to one stream's types; empty = every type. */ public function prune_tombstones( int $cutoff_sequence, string $cutoff_gmt, int $batch, array $object_types = array() ): array { global $wpdb; $none = array( 'deleted' => 0, 'watermark' => 0, ); if ( $cutoff_sequence <= 0 || $batch <= 0 ) { return $none; } // Each stream is pruned under its OWN cutoff (Sync_Journal_Purge), so the // batch must not reach across streams: an order tombstone is not eligible // merely because the catalogue's cutoff cleared it. $types = array_values( array_unique( array_filter( array_map( 'strval', $object_types ), static fn( string $type ): bool => '' !== $type ) ) ); $type_where = ''; $type_args = array(); if ( array() !== $types ) { $type_where = ' AND object_type IN (' . implode( ',', array_fill( 0, \count( $types ), '%s' ) ) . ')'; $type_args = $types; } $rows = $wpdb->get_results( $wpdb->prepare( 'SELECT sequence, object_type FROM ' . $this->table_name() . ' WHERE deleted = 1 AND sequence <= %d AND created_gmt < %s' . $type_where // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- %s placeholder list generated from count(). . ' ORDER BY sequence ASC LIMIT %d', $cutoff_sequence, $cutoff_gmt, ...array_merge( $type_args, array( $batch ) ) ), ARRAY_A ); if ( empty( $rows ) ) { return $none; } // Each pruned row raises ONLY its own object type's watermark: a client // reading a stream that never held these rows has missed nothing. $sequences = array(); $per_type = array(); foreach ( $rows as $row ) { $sequence = (int) $row['sequence']; $object_type = (string) $row['object_type']; $sequences[] = $sequence; $per_type[ $object_type ] = max( $per_type[ $object_type ] ?? 0, $sequence ); } $watermark = max( $sequences ); // Publish every watermark BEFORE deleting anything. A row deleted while // its type's horizon still reads below it is silently lost history — no // client would ever learn to reconcile it — so a failed write aborts the // whole batch and the next run retries it. foreach ( $per_type as $object_type => $sequence ) { $this->advance_prune_watermark( $object_type, $sequence ); if ( $this->prune_watermark( array( $object_type ) ) < $sequence ) { return $none; } } $deleted = $wpdb->query( 'DELETE FROM ' . $this->table_name() . ' WHERE sequence IN (' . implode( ',', $sequences ) . ')' ); return array( 'deleted' => false === $deleted ? 0 : (int) $deleted, 'watermark' => $watermark, ); } /** Resolve a wall-clock cutoff to one stable sequence boundary. */ public function sequence_at_or_before( string $cutoff_gmt ): int { global $wpdb; return (int) $wpdb->get_var( $wpdb->prepare( 'SELECT MAX(sequence) FROM ' . $this->table_name() . ' WHERE created_gmt < %s', $cutoff_gmt ) ); } /** Option holding one object type's lossy-prune watermark. */ public static function prune_watermark_option( string $object_type ): string { return self::PRUNE_WATERMARK_OPTION_PREFIX . $object_type; } /** * Highest sequence ever removed by lossy tombstone pruning FROM A STREAM. * * Mirrors head_sequence(): the caller names the object types its stream * serves and gets that stream's boundary. An empty list means every * registered type — the whole journal's boundary. * * @param array $object_types Object types the reading stream serves. */ public function prune_watermark( array $object_types = array() ): int { $watermark = 0; foreach ( self::watermark_object_types( $object_types ) as $object_type ) { $watermark = max( $watermark, (int) get_option( self::prune_watermark_option( $object_type ), 0 ) ); } return $watermark; } /** * Advance one object type's persisted watermark (never moves backwards). * * @param string $object_type Journal object type the pruned rows belonged to. * @param int $sequence Highest sequence pruned for that type. */ public function advance_prune_watermark( string $object_type, int $sequence ): void { global $wpdb; if ( $sequence <= 0 || '' === $object_type ) { return; } $option = self::prune_watermark_option( $object_type ); $wpdb->query( $wpdb->prepare( 'INSERT INTO ' . $wpdb->options . " (option_name, option_value, autoload) VALUES (%s, %d, 'yes')" . ' ON DUPLICATE KEY UPDATE option_value = GREATEST(CAST(option_value AS UNSIGNED), %d)', $option, $sequence, $sequence ) ); wp_cache_delete( $option, 'options' ); wp_cache_delete( 'alloptions', 'options' ); wp_cache_delete( 'notoptions', 'options' ); } /** Drop every per-type watermark — the journal's history is starting over. */ public static function reset_prune_watermarks(): void { global $wpdb; // Trailing separator trimmed from the LIKE so this also clears the single // pre-stream-scoping watermark a pre-release install may still carry. $names = $wpdb->get_col( $wpdb->prepare( "SELECT option_name FROM {$wpdb->options} WHERE option_name LIKE %s", $wpdb->esc_like( rtrim( self::PRUNE_WATERMARK_OPTION_PREFIX, '_' ) ) . '%' ) ); foreach ( (array) $names as $name ) { delete_option( (string) $name ); } } /** * Object types a watermark read covers: the named stream, or all of them. * * @param array $object_types Object types named by the caller. */ private static function watermark_object_types( array $object_types ): array { $types = array_values( array_filter( array_map( 'strval', $object_types ), static fn( string $type ): bool => '' !== $type ) ); if ( array() !== $types ) { return array_unique( $types ); } return array_unique( array_merge( array( 'order' ), self::catalogue_object_types() ) ); } /** One page of the change stream past a cursor, plus the head it was read against. */ public function page( array $object_types, int $since, int $limit ): array { global $wpdb; $types = array(); foreach ( $object_types as $object_type ) { $object_type = (string) $object_type; if ( '' !== $object_type ) { $types[] = $object_type; } } $sql = 'SELECT sequence, object_id, object_type, deleted, revision, modified_gmt FROM ' . $this->table_name() . ' WHERE '; $args = array(); if ( array() !== $types ) { $sql .= 'object_type IN (' . implode( ',', array_fill( 0, count( $types ), '%s' ) ) . ') AND '; $args = $types; } $sql .= 'sequence > %d ORDER BY sequence ASC LIMIT %d'; $args[] = max( 0, $since ); $args[] = max( 1, $limit ); $rows = $wpdb->get_results( $wpdb->prepare( $sql, ...$args ), ARRAY_A ); return array( 'rows' => array_map( static function ( array $row ): array { return self::normalize_row( $row ); }, is_array( $rows ) ? $rows : array() ), 'head' => $this->head_sequence( $types ), ); } /** Coerce a raw journal row to the served scalar types. */ private static function normalize_row( array $row ): array { return array( 'sequence' => isset( $row['sequence'] ) ? (int) $row['sequence'] : 0, 'object_id' => isset( $row['object_id'] ) ? (int) $row['object_id'] : 0, 'object_type' => isset( $row['object_type'] ) ? (string) $row['object_type'] : '', 'deleted' => ! empty( $row['deleted'] ) ? 1 : 0, 'revision' => isset( $row['revision'] ) ? (string) $row['revision'] : '', 'modified_gmt' => isset( $row['modified_gmt'] ) ? (string) $row['modified_gmt'] : gmdate( 'Y-m-d H:i:s' ), ); } }