maybe_create_table(); } /** * Get full table name with WordPress prefix. * * @return string */ public static function get_table_name() { global $wpdb; return $wpdb->prefix . self::TABLE_NAME; } /** * Create or upgrade table if needed. * * Checks stored version against current version and runs CREATE TABLE IF NOT EXISTS. * * @return void */ private function maybe_create_table() { $stored_version = get_option( self::DB_VERSION_OPTION, '0' ); if ( version_compare( $stored_version, self::DB_VERSION, '>=' ) ) { return; } global $wpdb; $table_name = self::get_table_name(); $charset_collate = $wpdb->get_charset_collate(); $sql = "CREATE TABLE {$table_name} ( id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, provider VARCHAR(32) NOT NULL DEFAULT 'postfix', newsletter_id BIGINT UNSIGNED NULL, campaign_name VARCHAR(255) NULL, from_email VARCHAR(255) NOT NULL, from_name VARCHAR(255) NULL, to_email VARCHAR(255) NOT NULL, to_name VARCHAR(255) NULL, reply_to VARCHAR(255) NULL, subject VARCHAR(998) NOT NULL, body_html MEDIUMTEXT NULL, body_text MEDIUMTEXT NULL, headers_json JSON NULL, created_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6), available_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6), locked_at DATETIME(6) NULL, sent_at DATETIME(6) NULL, locked_by VARCHAR(64) NULL, attempts INT UNSIGNED NOT NULL DEFAULT 0, last_error TEXT NULL, message_id VARCHAR(255) NULL, priority TINYINT NOT NULL DEFAULT 0, PRIMARY KEY (id), KEY idx_pick (sent_at, locked_at, available_at, priority, id), KEY idx_locked (locked_at), KEY idx_sent (sent_at), KEY idx_newsletter (newsletter_id, sent_at), KEY idx_created (created_at), KEY idx_to (to_email) ) {$charset_collate} ENGINE=InnoDB;"; require_once ABSPATH . 'wp-admin/includes/upgrade.php'; dbDelta( $sql ); update_option( self::DB_VERSION_OPTION, self::DB_VERSION ); } /** * Enqueue a single email message. * * @param Email_Message $message Email message object. * @param int $priority Priority (higher = processed first). * * @return int|false Inserted ID on success, false on failure. */ public function enqueue( Email_Message $message, $priority = 0 ) { global $wpdb; $table = self::get_table_name(); $data = array( 'provider' => $message->provider ?? 'postfix', 'newsletter_id' => $message->newsletter_id, 'campaign_name' => $message->campaign_name, 'from_email' => $message->from, 'from_name' => $message->from_name, 'to_email' => $message->to, 'to_name' => $message->to_name, 'reply_to' => $message->reply_to ?? null, 'subject' => $message->subject, 'body_html' => $message->body, 'body_text' => $message->body_text, 'headers_json' => ! empty( $message->headers ) ? wp_json_encode( $message->headers ) : null, 'priority' => absint( $priority ), ); $format = array( '%s', // provider. '%d', // newsletter_id. '%s', // campaign_name. '%s', // from_email. '%s', // from_name. '%s', // to_email. '%s', // to_name. '%s', // reply_to. '%s', // subject. '%s', // body_html. '%s', // body_text. '%s', // headers_json. '%d', // priority. ); $result = $wpdb->insert( $table, $data, $format ); return $result ? $wpdb->insert_id : false; } /** * Enqueue multiple email messages in bulk. * * @param Email_Message[] $messages Array of email message objects. * @param int $priority Priority for all messages. * * @return int[] Array of inserted IDs. */ public function enqueue_bulk( array $messages, $priority = 0 ) { if ( empty( $messages ) ) { return array(); } global $wpdb; $table = self::get_table_name(); $ids = array(); // Build multi-row INSERT. $values = array(); $placeholders = array(); foreach ( $messages as $message ) { $values[] = $message->provider ?? 'postfix'; $values[] = $message->newsletter_id; $values[] = $message->campaign_name; $values[] = $message->from; $values[] = $message->from_name; $values[] = $message->to; $values[] = $message->to_name; $values[] = $message->reply_to ?? null; $values[] = $message->subject; $values[] = $message->body; $values[] = $message->body_text; $values[] = ! empty( $message->headers ) ? wp_json_encode( $message->headers ) : null; $values[] = absint( $priority ); $placeholders[] = '(%s, %d, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %d)'; } $query = "INSERT INTO {$table} ( provider, newsletter_id, campaign_name, from_email, from_name, to_email, to_name, reply_to, subject, body_html, body_text, headers_json, priority ) VALUES " . implode( ', ', $placeholders ); // phpcs:ignore WordPress.DB.PreparedSQL.NotPrepared -- False positive: query is properly prepared with placeholders. $result = $wpdb->query( $wpdb->prepare( $query, $values ) ); if ( $result ) { // Get IDs of inserted rows (assumes sequential IDs). $first_id = $wpdb->insert_id; $ids = range( $first_id, $first_id + count( $messages ) - 1 ); } return $ids; } /** * Get and lock a batch of emails for processing. * * Uses transaction-based locking to prevent race conditions. * * @param int $limit Maximum emails to lock. * @param string $worker_id Unique worker identifier. * * @return object[] Array of email objects. */ public function get_and_lock_batch( $limit, $worker_id ) { global $wpdb; $table = self::get_table_name(); // Release stale locks first. $this->release_stale_locks(); // Start transaction for atomic lock acquisition. $wpdb->query( 'START TRANSACTION' ); // Select emails that are available and lock them. $items = $wpdb->get_results( $wpdb->prepare( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "SELECT * FROM {$table} WHERE sent_at IS NULL AND locked_at IS NULL AND available_at <= NOW(6) AND attempts < %d ORDER BY priority DESC, id ASC LIMIT %d FOR UPDATE", self::MAX_ATTEMPTS, $limit ) ); if ( ! empty( $items ) ) { $ids = wp_list_pluck( $items, 'id' ); $placeholders = implode( ',', array_fill( 0, count( $ids ), '%d' ) ); $wpdb->query( $wpdb->prepare( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name and placeholders are safe. "UPDATE {$table} SET locked_at = NOW(6), locked_by = %s WHERE id IN ({$placeholders})", array_merge( array( $worker_id ), $ids ) ) ); } $wpdb->query( 'COMMIT' ); return $items; } /** * Mark an email as successfully sent. * * @param int $id Email ID. * @param string|null $message_id Optional Message-ID header value. * * @return bool Success status. */ public function mark_sent( $id, $message_id = null ) { global $wpdb; $table = self::get_table_name(); $data = array( 'sent_at' => current_time( 'mysql', true ), 'message_id' => $message_id, ); $format = array( '%s', '%s' ); return false !== $wpdb->update( $table, $data, array( 'id' => $id ), $format, array( '%d' ) ); } /** * Mark an email as failed and increment retry counter. * * @param int $id Email ID. * @param string $error Error message. * * @return bool Success status. */ public function mark_failed( $id, $error ) { global $wpdb; $table = self::get_table_name(); // Get current attempts. $current = $wpdb->get_row( $wpdb->prepare( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "SELECT attempts FROM {$table} WHERE id = %d", $id ) ); if ( ! $current ) { return false; } $new_attempts = absint( $current->attempts ) + 1; // Calculate exponential backoff: 5 minutes, 15 minutes, 30 minutes. $delay_minutes = min( 5 * pow( 2, $new_attempts - 1 ), 30 ); $data = array( 'attempts' => $new_attempts, 'last_error' => $error, 'locked_at' => null, 'locked_by' => null, 'available_at' => gmdate( 'Y-m-d H:i:s.u', strtotime( "+{$delay_minutes} minutes" ) ), ); $format = array( '%d', '%s', '%s', '%s', '%s' ); return false !== $wpdb->update( $table, $data, array( 'id' => $id ), $format, array( '%d' ) ); } /** * Release lock on a specific email. * * @param int $id Email ID. * * @return bool Success status. */ public function release_lock( $id ) { global $wpdb; $table = self::get_table_name(); $data = array( 'locked_at' => null, 'locked_by' => null, ); return false !== $wpdb->update( $table, $data, array( 'id' => $id ), array( '%s', '%s' ), array( '%d' ) ); } /** * Release all locks older than timeout. * * @param int $timeout_seconds Timeout in seconds. * * @return int Number of locks released. */ public function release_stale_locks( $timeout_seconds = self::LOCK_TIMEOUT ) { global $wpdb; $table = self::get_table_name(); $result = $wpdb->query( $wpdb->prepare( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "UPDATE {$table} SET locked_at = NULL, locked_by = NULL WHERE locked_at IS NOT NULL AND locked_at < DATE_SUB(NOW(6), INTERVAL %d SECOND) AND sent_at IS NULL", $timeout_seconds ) ); return absint( $result ); } /** * Get real-time queue statistics. * * @return array{pending: int, locked: int, sent_24h: int, failed_24h: int} */ public function get_queue_stats() { global $wpdb; $table = self::get_table_name(); $stats = $wpdb->get_row( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "SELECT SUM(CASE WHEN sent_at IS NULL AND locked_at IS NULL AND attempts < " . self::MAX_ATTEMPTS . " THEN 1 ELSE 0 END) as pending, SUM(CASE WHEN locked_at IS NOT NULL THEN 1 ELSE 0 END) as locked, SUM(CASE WHEN sent_at IS NOT NULL AND sent_at >= DATE_SUB(NOW(), INTERVAL 24 HOUR) THEN 1 ELSE 0 END) as sent_24h, SUM(CASE WHEN attempts >= " . self::MAX_ATTEMPTS . " AND sent_at IS NULL THEN 1 ELSE 0 END) as failed_24h FROM {$table}", ARRAY_A ); return array( 'pending' => absint( $stats['pending'] ?? 0 ), 'locked' => absint( $stats['locked'] ?? 0 ), 'sent_24h' => absint( $stats['sent_24h'] ?? 0 ), 'failed_24h' => absint( $stats['failed_24h'] ?? 0 ), ); } /** * Get queue statistics for a specific campaign. * * @param int|null $newsletter_id Newsletter ID (null = all campaigns). * * @return array{total: int, sent: int, failed: int, pending: int} */ public function get_campaign_queue_stats( $newsletter_id = null ) { global $wpdb; $table = self::get_table_name(); $where = $newsletter_id ? $wpdb->prepare( 'WHERE newsletter_id = %d', $newsletter_id ) : ''; $stats = $wpdb->get_row( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name and WHERE clause are safe. "SELECT COUNT(*) as total, SUM(CASE WHEN sent_at IS NOT NULL THEN 1 ELSE 0 END) as sent, SUM(CASE WHEN attempts >= " . self::MAX_ATTEMPTS . " AND sent_at IS NULL THEN 1 ELSE 0 END) as failed, SUM(CASE WHEN sent_at IS NULL AND attempts < " . self::MAX_ATTEMPTS . " THEN 1 ELSE 0 END) as pending FROM {$table} {$where}", ARRAY_A ); return array( 'total' => absint( $stats['total'] ?? 0 ), 'sent' => absint( $stats['sent'] ?? 0 ), 'failed' => absint( $stats['failed'] ?? 0 ), 'pending' => absint( $stats['pending'] ?? 0 ), ); } /** * Get recent emails for admin UI. * * @param int $limit Maximum results. * @param int $offset Offset for pagination. * * @return object[] Array of email objects. */ public function get_recent_emails( $limit = 50, $offset = 0 ) { global $wpdb; $table = self::get_table_name(); return $wpdb->get_results( $wpdb->prepare( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "SELECT * FROM {$table} ORDER BY created_at DESC LIMIT %d OFFSET %d", $limit, $offset ) ); } /** * Get emails for a specific campaign. * * @param int $newsletter_id Newsletter ID. * @param int $limit Maximum results. * * @return object[] Array of email objects. */ public function get_campaign_emails( $newsletter_id, $limit = 50 ) { global $wpdb; $table = self::get_table_name(); return $wpdb->get_results( $wpdb->prepare( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "SELECT * FROM {$table} WHERE newsletter_id = %d ORDER BY created_at DESC LIMIT %d", $newsletter_id, $limit ) ); } /** * Count total emails in queue (for pagination). * * @return int Total count. */ public function count_total_emails() { global $wpdb; $table = self::get_table_name(); return absint( $wpdb->get_var( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "SELECT COUNT(*) FROM {$table}" ) ); } /** * Delete old sent/failed emails. * * @param int $days Number of days to keep. * * @return int Number of items deleted. */ public function cleanup_old_items( $days = 30 ) { global $wpdb; $table = self::get_table_name(); $result = $wpdb->query( $wpdb->prepare( // phpcs:ignore WordPress.DB.PreparedSQL.InterpolatedNotPrepared -- Table name is safe. "DELETE FROM {$table} WHERE sent_at IS NOT NULL AND sent_at < DATE_SUB(NOW(), INTERVAL %d DAY)", $days ) ); return absint( $result ); } /** * Delete all queue records for a campaign. * * Used when user wants to clear campaign history but preserve statistics. * * @param int $newsletter_id Newsletter ID. * * @return int Number of items deleted. */ public function delete_campaign_history( $newsletter_id ) { global $wpdb; $table = self::get_table_name(); $result = $wpdb->delete( $table, array( 'newsletter_id' => $newsletter_id ), array( '%d' ) ); return absint( $result ); } }