Files
2meet-data-optimizer/includes/migration/class-tmdo-migration-orchestrator.php
T
wpdev 76c01e44df refactor: 全部 128 個生產檔加入 declare(strict_types=1)(PR-H)
對齊 A v3.2.0。型別強制會把隱式轉換變成 TypeError,所以一次全檔加入
並跑完整測試(unit 451 / integration 398 全綠,無迴歸)。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TbG1keQQ7XBa7qMQY16KCY
2026-07-31 06:13:33 +08:00

644 lines
22 KiB
PHP
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
<?php
/**
* One-click User Entity Migration Orchestrator.
*
* Drives the full demote → backfill → verify → promote → cleanup flow as a
* single state machine. Designed for both small (sync, <30s) and large
* (async via cron, minutes-to-hours) sites with auto-detection.
*
* Fast mode: bulk SQL pivot (INSERT...SELECT...GROUP BY) for text-only groups
* is 10-50× faster than per-row PHP loops. Per-row safe_unserialize fallback
* applies only to groups containing `json` fields (currently only `hp_user`).
*
* @package TMDO
* @since 2.8.0
*/
declare(strict_types=1);
// phpcs:ignore WPDO.AntiEAV -- platform migration tool: WPDO v1->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();
}
}