v2 legacy meta orchestration if ( ! defined( 'ABSPATH' ) ) { exit; } // phpcs:disable Squiz.Commenting.FunctionComment,Squiz.Commenting.VariableComment,Squiz.Commenting.ClassComment,WordPress.DB.PreparedSQL.InterpolatedNotPrepared,WordPress.Security.EscapeOutput.ExceptionNotEscaped,WordPress.PHP.YodaConditions,Squiz.Commenting.InlineComment.InvalidEndChar -- Internal state-machine helpers and Throwable messages — not user output. SQL composed with sanitize_column_name() + $wpdb->prepare() values. /** * One-click User Migration Orchestrator — drives the demote → backfill → * verify → promote → cleanup state machine for the user entity. * * @since 2.8.0 */ class TMDO_Migration_Orchestrator { const OPT_JOB = 'wpdo_migration_job'; const OPT_LOCK = 'wpdo_migration_lock'; const TRANS_PROGRESS = 'wpdo_migration_progress'; const TRANS_NEEDS = 'wpdo_migration_needs_attention'; const CRON_HOOK = 'wpdo_migration_tick'; const BACKUP_DIR_REL = 'wpdo-backups'; const MAX_LOG_LINES = 50; const SYNC_THRESHOLD_SEC = 30; const SYNC_DEADLINE_SEC = 110; const VERIFY_SAMPLE_MIN = 500; const VERIFY_SAMPLE_RATIO = 0.10; const LOCK_TTL_SEC = 1800; const NEEDS_TTL_SEC = 300; const ENTITY_TYPE = 'user'; /** * Phase order. Each phase is idempotent — re-entering returns ok if * the desired state is already achieved. */ private const PHASES = array( 'diagnose', 'backup', 'demote', 'install_schema', 'backfill_bulk', 'backfill_unserialize', 'promote_shadow', 'verify_sample', 'promote_aeav', 'cleanup', 'completed', ); // ───────────────────────────────────────────────────────────────────── // Public API // ───────────────────────────────────────────────────────────────────── /** * Read-only inspection — returns what `start()` would do. * * @return array{managed_keys:array,eav_rows:int,users:int,usermeta:int,ratio:float,mode:string,groups:array,estimated_strategy:string,estimated_sec:float} */ public static function preflight(): array { global $wpdb; $managed_keys = self::get_managed_keys(); $eav_rows = self::count_eav_residue( $managed_keys ); $users = (int) $wpdb->get_var( "SELECT COUNT(*) FROM {$wpdb->users}" ); $usermeta = (int) $wpdb->get_var( "SELECT COUNT(*) FROM {$wpdb->usermeta}" ); $mode = TMDO_Mode_Manager::get( static::ENTITY_TYPE ); $groups = array(); foreach ( TMDO_Entity_Registry::get_groups_for_type( static::ENTITY_TYPE ) as $group ) { $keys = TMDO_Entity_Registry::get_group_keys( static::ENTITY_TYPE, $group ); $has_json = self::group_has_json_field( $group ); $residue = $keys ? self::count_eav_residue( $keys ) : 0; $flat_table = TMDO_Schema_Manager::get_table_name( static::ENTITY_TYPE, $group ); $flat_rows = TMDO_Schema_Manager::table_exists( $flat_table ) ? (int) $wpdb->get_var( "SELECT COUNT(*) FROM `{$flat_table}`" ) : 0; $groups[ $group ] = array( 'keys' => $keys, 'residue' => $residue, 'flat_rows' => $flat_rows, 'has_json' => $has_json, 'strategy' => $has_json ? 'row_by_row' : 'bulk_pivot', ); } // Heuristic: 5ms per bulk row, 80ms per per-row PHP entity. $bulk_rows = 0; $row_rows = 0; foreach ( $groups as $g ) { if ( $g['has_json'] ) { $row_rows += $g['residue']; } else { $bulk_rows += $g['residue']; } } $estimated_sec = ( $bulk_rows * 0.005 ) + ( $row_rows * 0.080 ) + 2.0; $strategy = $estimated_sec <= self::SYNC_THRESHOLD_SEC ? 'sync' : 'async'; return array( 'managed_keys' => $managed_keys, 'eav_rows' => $eav_rows, 'users' => $users, 'usermeta' => $usermeta, 'ratio' => $users > 0 ? round( $usermeta / $users, 2 ) : 0, 'mode' => $mode, 'groups' => $groups, 'estimated_strategy' => $strategy, 'estimated_sec' => round( $estimated_sec, 1 ), ); } /** * Cached attention summary for dashboard widget + tab-nav red dot. * * Hits a single SELECT against wp_usermeta + an in-process ratio * calculation; cached for NEEDS_TTL_SEC (5 min) so dashboard widgets * stay snappy. Bust the cache after migrations or mode changes via * `bust_attention_cache()`. * * @return array{needs:bool,eav_rows:int,groups_with_residue:int,ratio:float,mode:string,job_state:string} */ public static function needs_attention(): array { $cached = get_transient( self::TRANS_NEEDS ); if ( is_array( $cached ) ) { return $cached; } global $wpdb; $managed_keys = self::get_managed_keys(); $eav_rows = self::count_eav_residue( $managed_keys ); $groups_with_residue = 0; foreach ( TMDO_Entity_Registry::get_groups_for_type( static::ENTITY_TYPE ) as $group ) { $keys = TMDO_Entity_Registry::get_group_keys( static::ENTITY_TYPE, $group ); if ( $keys && self::count_eav_residue( $keys ) > 0 ) { ++$groups_with_residue; } } $users = (int) $wpdb->get_var( "SELECT COUNT(*) FROM {$wpdb->users}" ); $usermeta = (int) $wpdb->get_var( "SELECT COUNT(*) FROM {$wpdb->usermeta}" ); $mode = TMDO_Mode_Manager::get( static::ENTITY_TYPE ); $job_state = 'idle'; $job = get_option( self::OPT_JOB ); if ( is_array( $job ) && ! empty( $job['state'] ) ) { $job_state = (string) $job['state']; } $result = array( 'needs' => $eav_rows > 0, 'eav_rows' => $eav_rows, 'groups_with_residue' => $groups_with_residue, 'ratio' => $users > 0 ? round( $usermeta / $users, 2 ) : 0, 'mode' => $mode, 'job_state' => $job_state, ); set_transient( self::TRANS_NEEDS, $result, self::NEEDS_TTL_SEC ); return $result; } /** * Bust the attention cache — call after migrations, mode changes, * group registrations, etc. */ public static function bust_attention_cache(): void { delete_transient( self::TRANS_NEEDS ); } /** * Begin a migration job. If small, runs to completion in this request. * Otherwise schedules cron-driven ticks and returns immediately. * * @param array{verify_strict?:bool,verify_24h?:bool,auto_backup?:bool,force_async?:bool,dry_run?:bool} $options * @return array{ok:bool,job_id?:string,strategy?:string,error?:string,reason?:string} */ public static function start( array $options = array() ): array { if ( ! self::acquire_lock() ) { return array( 'ok' => false, 'error' => __( '另一個 migration job 正在執行;請先取消或等候完成。', '2meet-data-optimizer' ), ); } $preflight = self::preflight(); // Idempotency: nothing to do. if ( $preflight['eav_rows'] === 0 && TMDO_Mode_Manager::MODE_AEAV_ONLY === $preflight['mode'] ) { self::release_lock(); self::bust_attention_cache(); return array( 'ok' => false, 'reason' => 'nothing_to_do', 'error' => sprintf( /* translators: %s: ratio. */ __( '所有 entity group 已 aeav_only / 0 EAV 殘留。當前 ratio %s。', '2meet-data-optimizer' ), (string) $preflight['ratio'] ), ); } $options = wp_parse_args( $options, array( 'verify_strict' => true, 'verify_24h' => false, 'auto_backup' => true, 'force_async' => false, 'dry_run' => false, ) ); $job = array( 'job_id' => 'mig_' . wp_generate_password( 12, false ), 'started_at' => time(), 'updated_at' => time(), 'phase' => 'diagnose', 'phase_index' => 0, 'phase_progress' => 0, 'overall_progress' => 0, 'log' => array(), 'options' => $options, 'metrics' => array( 'mode_start' => $preflight['mode'], 'ratio_start' => $preflight['ratio'], 'ratio_now' => $preflight['ratio'], 'eav_rows_start' => $preflight['eav_rows'], 'eav_rows_now' => $preflight['eav_rows'], 'users' => $preflight['users'], ), 'strategy' => $options['force_async'] ? 'async' : $preflight['estimated_strategy'], 'state' => 'running', 'errors' => array(), 'backup_path' => '', ); self::log( $job, sprintf( 'Job started — strategy=%s, residue=%d rows across %d groups, mode=%s', $job['strategy'], $preflight['eav_rows'], count( array_filter( $preflight['groups'], fn( $g ) => $g['residue'] > 0 ) ), $preflight['mode'] ) ); self::persist_job( $job ); if ( 'sync' === $job['strategy'] ) { self::run_sync_loop( $job ); } else { self::schedule_next_tick(); } return array( 'ok' => true, 'job_id' => $job['job_id'], 'strategy' => $job['strategy'], ); } /** * Advance one phase. Used by cron and by sync inline loop. * * @return array{done:bool,phase:string,error?:string} * @throws \RuntimeException When a phase callback throws or returns non-ok status; caught internally and converted to an `error` array entry. */ public static function tick(): array { $job = get_option( self::OPT_JOB ); if ( ! is_array( $job ) || empty( $job['job_id'] ) ) { return array( 'done' => true, 'phase' => 'idle', 'error' => 'No active job', ); } if ( 'running' !== ( $job['state'] ?? '' ) ) { return array( 'done' => true, 'phase' => $job['phase'] ?? 'unknown', ); } $current = $job['phase']; try { $phase_started = microtime( true ); $phase = self::make_phase( $current ); $result = $phase->execute( $job ); if ( 'in_progress' === ( $result['status'] ?? '' ) ) { self::persist_job( $job ); if ( 'async' === $job['strategy'] ) { self::schedule_next_tick(); } return array( 'done' => false, 'phase' => $current, ); } if ( 'ok' !== ( $result['status'] ?? '' ) ) { throw new \RuntimeException( $result['message'] ?? 'Phase returned non-ok status' ); } $next = self::next_phase( $current ); self::log( $job, sprintf( '✓ %s (%.2fs)', $current, microtime( true ) - $phase_started ) ); $job['phase'] = $next; $job['phase_index'] = array_search( $next, self::PHASES, true ); $job['phase_progress'] = 0; $job['overall_progress'] = (int) round( ( $job['phase_index'] / ( count( self::PHASES ) - 1 ) ) * 100 ); $job['updated_at'] = time(); if ( 'completed' === $next ) { self::make_phase( 'completed' )->execute( $job ); self::persist_job( $job ); self::release_lock(); self::bust_attention_cache(); return array( 'done' => true, 'phase' => 'completed', ); } self::persist_job( $job ); if ( 'async' === $job['strategy'] ) { self::schedule_next_tick(); } return array( 'done' => false, 'phase' => $next, ); } catch ( \Throwable $e ) { self::handle_failure( $job, $e ); return array( 'done' => true, 'phase' => 'failed', 'error' => $e->getMessage(), ); } } /** * Inspection — read-only snapshot for the polling UI. */ public static function get_status(): array { $cached = get_transient( self::TRANS_PROGRESS ); if ( is_array( $cached ) ) { return $cached; } $job = get_option( self::OPT_JOB ); if ( ! is_array( $job ) ) { return array( 'state' => 'idle' ); } return self::project_status( $job ); } public static function cancel(): bool { $job = get_option( self::OPT_JOB ); if ( ! is_array( $job ) || 'running' !== ( $job['state'] ?? '' ) ) { return false; } // Auto-rollback to safe state (dual_write) before clearing. try { $current_mode = TMDO_Mode_Manager::get( static::ENTITY_TYPE ); if ( TMDO_Mode_Manager::MODE_AEAV_ONLY === $current_mode ) { TMDO_Mode_Manager::set( static::ENTITY_TYPE, TMDO_Mode_Manager::MODE_SHADOW_READ ); TMDO_Mode_Manager::set( static::ENTITY_TYPE, TMDO_Mode_Manager::MODE_DUAL_WRITE ); } elseif ( TMDO_Mode_Manager::MODE_SHADOW_READ === $current_mode ) { TMDO_Mode_Manager::set( static::ENTITY_TYPE, TMDO_Mode_Manager::MODE_DUAL_WRITE ); } } catch ( \Throwable $e ) { // Logged but non-fatal. TMDO_Logger::warning( 'migration_cancel_rollback_failed', array( 'error' => $e->getMessage() ) ); } $job['state'] = 'cancelled'; $job['cancelled_at'] = time(); self::log( $job, '⚠ Cancelled by operator — rolled back to dual_write.' ); self::persist_job( $job ); self::release_lock(); self::bust_attention_cache(); wp_clear_scheduled_hook( self::CRON_HOOK ); return true; } public static function resume(): array { $job = get_option( self::OPT_JOB ); if ( ! is_array( $job ) ) { return array( 'ok' => false, 'error' => 'No job to resume', ); } if ( ! in_array( $job['state'] ?? '', array( 'failed', 'paused' ), true ) ) { return array( 'ok' => false, 'error' => 'Job is not in a resumable state', ); } $job['state'] = 'running'; $job['updated_at'] = time(); self::log( $job, '↻ Resumed by operator.' ); self::persist_job( $job ); self::acquire_lock(); if ( 'sync' === $job['strategy'] ) { self::run_sync_loop( $job ); } else { self::schedule_next_tick(); } return array( 'ok' => true ); } // ───────────────────────────────────────────────────────────────────── // Sync inline loop — for sites where total work fits within // SYNC_DEADLINE_SEC. Heartbeats progress via transient on every phase. // ───────────────────────────────────────────────────────────────────── private static function run_sync_loop( array $job ): void { set_time_limit( self::SYNC_DEADLINE_SEC + 10 ); $deadline = microtime( true ) + self::SYNC_DEADLINE_SEC; while ( microtime( true ) < $deadline ) { $result = self::tick(); if ( $result['done'] ) { return; } // Tiny pause to let DB breathe and avoid 100% CPU pegs. usleep( 5000 ); $job = get_option( self::OPT_JOB ); if ( ! is_array( $job ) || 'running' !== ( $job['state'] ?? '' ) ) { return; } } // Deadline reached but not done — convert to async. $job = get_option( self::OPT_JOB ); $job['strategy'] = 'async'; self::log( $job, '⏱ Sync deadline reached — switching to async (cron-driven) for remaining phases.' ); self::persist_job( $job ); self::schedule_next_tick(); } // ───────────────────────────────────────────────────────────────────── // Phases — each returns ['status' => 'ok'|'in_progress'|'error', 'message' => string] // ───────────────────────────────────────────────────────────────────── // ── Phase factory ───────────────────────────────────────────────────────── private static function make_phase( string $name ): TMDO_Migration_Phase_Interface { $map = array( 'diagnose' => TMDO_Phase_Diagnose::class, 'backup' => TMDO_Phase_Backup::class, 'demote' => TMDO_Phase_Demote::class, 'install_schema' => TMDO_Phase_Install_Schema::class, 'backfill_bulk' => TMDO_Phase_Backfill_Bulk::class, 'backfill_unserialize' => TMDO_Phase_Backfill_Unserialize::class, 'promote_shadow' => TMDO_Phase_Promote_Shadow::class, 'verify_sample' => TMDO_Phase_Verify_Sample::class, 'promote_aeav' => TMDO_Phase_Promote_Aeav::class, 'cleanup' => TMDO_Phase_Cleanup::class, 'completed' => TMDO_Phase_Completed::class, ); if ( ! isset( $map[ $name ] ) ) { throw new \RuntimeException( "Unknown phase: {$name}" ); } return new $map[ $name ]( static::ENTITY_TYPE ); } // ───────────────────────────────────────────────────────────────────── // Helpers // ───────────────────────────────────────────────────────────────────── /** Returns all meta_keys registered as entity fields for `user`. */ private static function get_managed_keys(): array { $keys = array(); foreach ( TMDO_Entity_Registry::get_groups_for_type( static::ENTITY_TYPE ) as $group ) { $keys = array_merge( $keys, TMDO_Entity_Registry::get_group_keys( static::ENTITY_TYPE, $group ) ); } return array_values( array_unique( $keys ) ); } private static function group_has_json_field( string $group ): bool { foreach ( TMDO_Entity_Registry::get_group_fields( static::ENTITY_TYPE, $group ) as $field ) { if ( 'json' === ( $field['type'] ?? '' ) ) { return true; } } return false; } private static function count_eav_residue( array $keys = array() ): int { global $wpdb; if ( empty( $keys ) ) { $keys = self::get_managed_keys(); } if ( empty( $keys ) ) { return 0; } $placeholders = implode( ',', array_fill( 0, count( $keys ), '%s' ) ); return (int) $wpdb->get_var( $wpdb->prepare( "SELECT COUNT(*) FROM {$wpdb->usermeta} WHERE meta_key IN ({$placeholders})", ...$keys ) ); } private static function next_phase( string $current ): string { $idx = array_search( $current, self::PHASES, true ); if ( false === $idx || $idx + 1 >= count( self::PHASES ) ) { return 'completed'; } return self::PHASES[ $idx + 1 ]; } private static function log( array &$job, string $message ): void { $line = '[' . gmdate( 'H:i:s' ) . '] ' . $message; $job['log'][] = $line; if ( count( $job['log'] ) > self::MAX_LOG_LINES ) { $job['log'] = array_slice( $job['log'], -self::MAX_LOG_LINES ); } $job['updated_at'] = time(); } private static function persist_job( array $job ): void { update_option( self::OPT_JOB, $job, false ); set_transient( self::TRANS_PROGRESS, self::project_status( $job ), 60 ); } /** Public-safe projection of job state for the polling UI. */ private static function project_status( array $job ): array { return array( 'job_id' => $job['job_id'] ?? '', 'state' => $job['state'] ?? 'idle', 'phase' => $job['phase'] ?? 'idle', 'phase_index' => $job['phase_index'] ?? 0, 'phase_total' => count( self::PHASES ) - 1, 'overall_progress' => $job['overall_progress'] ?? 0, 'log' => $job['log'] ?? array(), 'metrics' => $job['metrics'] ?? array(), 'strategy' => $job['strategy'] ?? 'sync', 'started_at' => $job['started_at'] ?? 0, 'updated_at' => $job['updated_at'] ?? 0, 'completed_at' => $job['completed_at'] ?? null, 'cancelled_at' => $job['cancelled_at'] ?? null, 'errors' => $job['errors'] ?? array(), 'backup_path' => isset( $job['backup_path'] ) ? basename( $job['backup_path'] ) : '', ); } private static function acquire_lock(): bool { if ( false === add_option( self::OPT_LOCK, time(), '', false ) ) { $existing = (int) get_option( self::OPT_LOCK, 0 ); if ( $existing > 0 && time() - $existing > self::LOCK_TTL_SEC ) { delete_option( self::OPT_LOCK ); return add_option( self::OPT_LOCK, time(), '', false ); } return false; } return true; } private static function release_lock(): void { delete_option( self::OPT_LOCK ); } private static function schedule_next_tick(): void { if ( ! wp_next_scheduled( self::CRON_HOOK ) ) { wp_schedule_single_event( time() + 1, self::CRON_HOOK ); } } private static function handle_failure( array &$job, \Throwable $e ): void { $job['state'] = 'failed'; $job['errors'][] = array( 'phase' => $job['phase'] ?? 'unknown', 'message' => $e->getMessage(), 'at' => time(), ); self::log( $job, sprintf( '✗ %s failed: %s', $job['phase'] ?? '?', $e->getMessage() ) ); // Auto-rollback: only if mode is currently aeav_only AND we haven't reached cleanup yet. try { $current = TMDO_Mode_Manager::get( static::ENTITY_TYPE ); $phase = $job['phase'] ?? ''; $pre_cleanup = ! in_array( $phase, array( 'cleanup', 'completed' ), true ); if ( $pre_cleanup && TMDO_Mode_Manager::MODE_AEAV_ONLY === $current ) { TMDO_Mode_Manager::set( static::ENTITY_TYPE, TMDO_Mode_Manager::MODE_SHADOW_READ ); TMDO_Mode_Manager::set( static::ENTITY_TYPE, TMDO_Mode_Manager::MODE_DUAL_WRITE ); self::log( $job, '↩ Auto-rolled back to dual_write (reads safe)' ); } } catch ( \Throwable $inner ) { TMDO_Logger::warning( 'migration_rollback_failed', array( 'error' => $inner->getMessage() ) ); } TMDO_Logger::warning( 'migration_phase_failed', array( 'phase' => $job['phase'] ?? '?', 'job' => $job['job_id'] ?? '?', 'error' => $e->getMessage(), ) ); self::persist_job( $job ); self::release_lock(); self::bust_attention_cache(); } /** * Cron callback — wires CRON_HOOK to tick(). */ public static function cron_tick(): void { self::tick(); } }