/** * Apply lobster profile stats / summary projection / segment seed / menu remove to all active tenant DBs. * * Usage: * node apply-lobster-profile-migration-all-tenants.js * node apply-lobster-profile-migration-all-tenants.js --tenant-id=33 */ const fs = require('fs') const path = require('path') const mysql = require('mysql2/promise') const MASTER = { host: process.env.DB_HOST || 'cq-cdb-8fjmemkb.sql.tencentcdb.com', port: Number(process.env.DB_PORT || 27220), user: process.env.DB_USER || 'root', password: process.env.DB_PASSWORD || 'Ylrz_1q2w3e4r5t6y', database: process.env.DB_NAME || 'ylrz_saas' } const SQL_DIR = path.join(__dirname, '../sql') const PROFILE_DDL = fs.readFileSync(path.join(SQL_DIR, 'lobster_user_profile.sql'), 'utf8') const STATS_DDL = fs.readFileSync(path.join(SQL_DIR, 'lobster_profile_stats.sql'), 'utf8') const SEGMENT_DDL = fs.readFileSync(path.join(__dirname, '../fs-agent/src/main/resources/db/migration/tenant/V20260625_02__lobster_user_segment.sql'), 'utf8') const SEGMENT_SQL = fs.readFileSync(path.join(SQL_DIR, 'V20260722_01__lobster_segment_threshold_config.sql'), 'utf8') const MENU_SQL = fs.readFileSync(path.join(SQL_DIR, 'V20260711_01__lobster_remove_conversation_summary_menu.sql'), 'utf8') const PROFILE_COLUMNS = [ ['channel_type', "VARCHAR(20) NOT NULL DEFAULT 'QW'"], ['latest_stage', "VARCHAR(64) DEFAULT NULL"], ['latest_customer_attitude', "VARCHAR(32) DEFAULT NULL"], ['latest_summary_text', "VARCHAR(512) DEFAULT NULL"], ['profile_aggregate_preview', "VARCHAR(420) DEFAULT NULL"], ['latest_summary_time', "DATETIME DEFAULT NULL"] ] const SUMMARY_COLUMNS = [ ['external_user_id', "varchar(128) DEFAULT NULL COMMENT '\u5916\u90e8\u7528\u6237ID'"], ['channel_type', "varchar(20) NOT NULL DEFAULT 'QW' COMMENT '\u6e20\u9053\u7c7b\u578b QW/WX/IM'"], ['summary_text', "text DEFAULT NULL COMMENT '\u6458\u8981\u6587\u672c\u5185\u5bb9'"], ['customer_attitude', "varchar(32) DEFAULT NULL COMMENT '\u5ba2\u6237\u6001\u5ea6'"], ['stage', "varchar(64) DEFAULT NULL COMMENT '\u5bf9\u8bdd\u9636\u6bb5'"] ] const PROFILE_INDEXES = [ ['idx_profile_latest_stage', '(company_id, latest_stage)'], ['idx_profile_latest_attitude', '(company_id, latest_customer_attitude)'], ['idx_profile_latest_summary_time', '(company_id, latest_summary_time)'] ] function parseArgs() { const tenantIds = [] for (const arg of process.argv.slice(2)) { if (arg.startsWith('--tenant-id=')) { tenantIds.push(Number(arg.split('=')[1])) } } return tenantIds } function parseJdbcUrl(jdbcUrl, dbName, user, password) { if (jdbcUrl && jdbcUrl.startsWith('jdbc:mysql://')) { const body = jdbcUrl.slice('jdbc:mysql://'.length) const slash = body.indexOf('/') const hostPort = slash >= 0 ? body.slice(0, slash) : body const rest = slash >= 0 ? body.slice(slash + 1) : dbName const q = rest.indexOf('?') const database = q >= 0 ? rest.slice(0, q) : rest const colon = hostPort.indexOf(':') return { host: colon >= 0 ? hostPort.slice(0, colon) : hostPort, port: colon >= 0 ? Number(hostPort.slice(colon + 1)) : MASTER.port, database, user: user || MASTER.user, password: password || MASTER.password } } return { host: MASTER.host, port: MASTER.port, database: dbName, user: user || MASTER.user, password: password || MASTER.password } } async function tableExists(conn, table) { const [rows] = await conn.query( 'SELECT COUNT(*) AS cnt FROM information_schema.TABLES WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ?', [table] ) return rows[0].cnt > 0 } async function columnExists(conn, table, column) { const [rows] = await conn.query( 'SELECT COUNT(*) AS cnt FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND COLUMN_NAME = ?', [table, column] ) return rows[0].cnt > 0 } async function indexExists(conn, table, indexName) { const [rows] = await conn.query( 'SELECT COUNT(*) AS cnt FROM information_schema.STATISTICS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND INDEX_NAME = ?', [table, indexName] ) return rows[0].cnt > 0 } async function addColumnIfMissing(conn, table, column, definition) { if (await columnExists(conn, table, column)) return false await conn.query(`ALTER TABLE \`${table}\` ADD COLUMN \`${column}\` ${definition}`) return true } async function addIndexIfMissing(conn, table, indexName, columns) { if (await indexExists(conn, table, indexName)) return false await conn.query(`ALTER TABLE \`${table}\` ADD INDEX \`${indexName}\` ${columns}`) return true } async function ensureProfileSchema(conn) { if (!(await tableExists(conn, 'lobster_user_profile'))) { await conn.query(PROFILE_DDL) console.log(' created lobster_user_profile') } await conn.query(STATS_DDL) for (const [col, def] of PROFILE_COLUMNS) { if (await addColumnIfMissing(conn, 'lobster_user_profile', col, def)) { console.log(` added lobster_user_profile.${col}`) } } for (const [idx, cols] of PROFILE_INDEXES) { if (await addIndexIfMissing(conn, 'lobster_user_profile', idx, cols)) { console.log(` added index ${idx}`) } } } async function ensureSummaryColumns(conn) { if (!(await tableExists(conn, 'lobster_conversation_summary'))) return false for (const [col, def] of SUMMARY_COLUMNS) { await addColumnIfMissing(conn, 'lobster_conversation_summary', col, def) } return true } async function backfillSummaryProjection(conn) { if (!(await tableExists(conn, 'lobster_user_profile'))) return if (!(await tableExists(conn, 'lobster_conversation_summary'))) return if (!(await columnExists(conn, 'lobster_conversation_summary', 'external_user_id'))) { console.log(' skip summary backfill: external_user_id missing') return } const hasStage = await columnExists(conn, 'lobster_conversation_summary', 'stage') const hasAttitude = await columnExists(conn, 'lobster_conversation_summary', 'customer_attitude') const hasSummaryText = await columnExists(conn, 'lobster_conversation_summary', 'summary_text') const hasSummaryContent = await columnExists(conn, 'lobster_conversation_summary', 'summary_content') const hasDelFlag = await columnExists(conn, 'lobster_conversation_summary', 'del_flag') const stageExpr = hasStage ? 's.stage' : 'NULL' const attitudeExpr = hasAttitude ? 's.customer_attitude' : 'NULL' let summaryExpr = 'NULL' if (hasSummaryText && hasSummaryContent) { summaryExpr = "COALESCE(NULLIF(TRIM(s.summary_text), ''), NULLIF(TRIM(s.summary_content), ''))" } else if (hasSummaryText) { summaryExpr = "NULLIF(TRIM(s.summary_text), '')" } else if (hasSummaryContent) { summaryExpr = "NULLIF(TRIM(s.summary_content), '')" } const hasProfileChannel = await columnExists(conn, 'lobster_user_profile', 'channel_type') const joinChannel = hasProfileChannel ? ' AND p.channel_type = src.channel_type' : '' const channelTypeExpr = hasProfileChannel ? 's.channel_type' : "'QW'" const delFilter = hasDelFlag ? '(del_flag IS NULL OR del_flag = 0)' : '1=1' const sql = ` UPDATE lobster_user_profile p INNER JOIN ( SELECT s.company_id, s.external_user_id, ${channelTypeExpr} AS channel_type, ${stageExpr} AS stage, ${attitudeExpr} AS customer_attitude, ${summaryExpr} AS summary_text, COALESCE(s.update_time, s.create_time) AS summary_time FROM lobster_conversation_summary s INNER JOIN ( SELECT company_id, external_user_id, channel_type, MAX(COALESCE(update_time, create_time)) AS max_time FROM lobster_conversation_summary WHERE ${delFilter} AND external_user_id IS NOT NULL AND company_id IS NOT NULL GROUP BY company_id, external_user_id, channel_type ) latest ON latest.company_id = s.company_id AND latest.external_user_id = s.external_user_id AND latest.channel_type = s.channel_type AND latest.max_time = COALESCE(s.update_time, s.create_time) WHERE ${delFilter} ) src ON p.company_id = src.company_id AND p.external_user_id = src.external_user_id${joinChannel} SET p.latest_stage = COALESCE(src.stage, p.latest_stage), p.latest_customer_attitude = COALESCE(src.customer_attitude, p.latest_customer_attitude), p.latest_summary_text = COALESCE(LEFT(src.summary_text, 512), p.latest_summary_text), p.profile_aggregate_preview = CASE WHEN src.stage IS NOT NULL AND src.summary_text IS NOT NULL AND TRIM(src.summary_text) != '' THEN LEFT(CONCAT('\u3010\u9636\u6bb5\u3011', src.stage, '\uff5c', TRIM(src.summary_text)), 420) WHEN src.stage IS NOT NULL THEN CONCAT('\u3010\u9636\u6bb5\u3011', src.stage) WHEN src.summary_text IS NOT NULL AND TRIM(src.summary_text) != '' THEN LEFT(TRIM(src.summary_text), 420) ELSE p.profile_aggregate_preview END, p.latest_summary_time = COALESCE(src.summary_time, p.latest_summary_time) WHERE p.deleted = 0 AND (p.latest_summary_text IS NULL OR p.latest_summary_time IS NULL)` const [result] = await conn.query(sql) console.log(` summary backfill affectedRows=${result.affectedRows || 0}`) } async function backfillStatsCache(conn) { if (!(await tableExists(conn, 'lobster_user_profile'))) return const companySql = ` INSERT INTO lobster_profile_stats (company_id, total_profiles, enriched_count, interaction_sum, avg_conversations, active_today, refresh_time) SELECT p.company_id, COUNT(1), IFNULL(SUM(CASE WHEN p.value_score > 10 OR p.total_purchase > 0 OR (p.internal_tags IS NOT NULL AND p.internal_tags != '') THEN 1 ELSE 0 END), 0), IFNULL(SUM(p.interaction_count), 0), IFNULL(ROUND(AVG(p.interaction_count), 2), 0), IFNULL(SUM(CASE WHEN DATE(p.last_active_time) = CURDATE() THEN 1 ELSE 0 END), 0), NOW() FROM lobster_user_profile p WHERE p.deleted = 0 GROUP BY p.company_id ON DUPLICATE KEY UPDATE total_profiles = VALUES(total_profiles), enriched_count = VALUES(enriched_count), interaction_sum = VALUES(interaction_sum), avg_conversations = VALUES(avg_conversations), active_today = VALUES(active_today), refresh_time = NOW()` const tenantSql = ` INSERT INTO lobster_profile_stats (company_id, total_profiles, enriched_count, interaction_sum, avg_conversations, active_today, refresh_time) SELECT 0, COUNT(1), IFNULL(SUM(CASE WHEN p.value_score > 10 OR p.total_purchase > 0 OR (p.internal_tags IS NOT NULL AND p.internal_tags != '') THEN 1 ELSE 0 END), 0), IFNULL(SUM(p.interaction_count), 0), IFNULL(ROUND(AVG(p.interaction_count), 2), 0), IFNULL(SUM(CASE WHEN DATE(p.last_active_time) = CURDATE() THEN 1 ELSE 0 END), 0), NOW() FROM lobster_user_profile p WHERE p.deleted = 0 ON DUPLICATE KEY UPDATE total_profiles = VALUES(total_profiles), enriched_count = VALUES(enriched_count), interaction_sum = VALUES(interaction_sum), avg_conversations = VALUES(avg_conversations), active_today = VALUES(active_today), refresh_time = NOW()` await conn.query(companySql) await conn.query(tenantSql) console.log(' stats cache backfill done') } async function ensureUserSegmentTable(conn) { if (!(await tableExists(conn, 'lobster_user_segment'))) { await conn.query(SEGMENT_DDL) console.log(' created lobster_user_segment') } } async function applySegmentSeed(conn) { await ensureUserSegmentTable(conn) await conn.query(SEGMENT_SQL) console.log(' segment seed applied') } async function removeSummaryMenu(conn) { const statements = MENU_SQL.split(';').map(s => s.trim()).filter(Boolean) for (const sql of statements) { await conn.query(sql) } console.log(' conversation-summary menu removed') } async function verifyTenant(conn) { return { lobster_profile_stats: await tableExists(conn, 'lobster_profile_stats'), latest_stage: await columnExists(conn, 'lobster_user_profile', 'latest_stage'), latest_summary_text: await columnExists(conn, 'lobster_user_profile', 'latest_summary_text'), interaction_sum: await columnExists(conn, 'lobster_profile_stats', 'interaction_sum'), builtinSegmentRows: (await tableExists(conn, 'lobster_user_segment')) ? (await conn.query("SELECT COUNT(*) AS cnt FROM lobster_user_segment WHERE deleted = 0 AND segment_code IN ('NEW','ACTIVE','SLEEP','DORMANT','CHURN')"))[0][0].cnt : 0, statsCacheRows: (await tableExists(conn, 'lobster_profile_stats')) ? (await conn.query('SELECT COUNT(*) AS cnt FROM lobster_profile_stats'))[0][0].cnt : 0 } } async function applyTenant(tenant) { const conn = await mysql.createConnection({ host: tenant.host, port: tenant.port, user: tenant.user, password: tenant.password, database: tenant.database, multipleStatements: true, charset: 'utf8mb4' }) try { console.log(`\n[RUN] tenantId=${tenant.id} db=${tenant.database}`) await ensureProfileSchema(conn) await ensureSummaryColumns(conn) await backfillSummaryProjection(conn) await backfillStatsCache(conn) await applySegmentSeed(conn) await removeSummaryMenu(conn) const checks = await verifyTenant(conn) console.log(` verify: ${JSON.stringify(checks)}`) } finally { await conn.end() } } async function main() { const filterIds = parseArgs() const masterConn = await mysql.createConnection({ ...MASTER, charset: 'utf8mb4' }) const tenants = [] try { const [rows] = await masterConn.query( 'SELECT id, db_name, db_url, db_account, db_pwd, status FROM tenant_info ORDER BY id' ) for (const row of rows) { if (row.status !== 1) continue if (filterIds.length > 0 && !filterIds.includes(Number(row.id))) continue tenants.push({ id: row.id, status: row.status, ...parseJdbcUrl(row.db_url, row.db_name, row.db_account, row.db_pwd) }) } } finally { await masterConn.end() } console.log(`Active tenants to migrate: ${tenants.length}`) let ok = 0 let fail = 0 for (const tenant of tenants) { try { await applyTenant(tenant) ok++ } catch (err) { fail++ console.error(`[FAIL] tenantId=${tenant.id} db=${tenant.database}: ${err.message}`) } } console.log(`\nDone. ok=${ok} fail=${fail}`) if (fail > 0) process.exit(1) } main().catch(err => { console.error(err) process.exit(1) })