corp-transfer/scripts/migrate-corpid-sender.js

237 lines
6.7 KiB
JavaScript
Raw Permalink Normal View History

2026-04-03 10:26:45 +08:00
"use strict";
const axios = require("axios");
const { MongoClient } = require("mongodb");
const { EJSON } = require("bson");
const CONFIG = {
// 生产源库(请按真实生产配置核对)
sourceMongoUri: "mongodb://root:PRO%26youcan2025233@101.35.91.238:27017/admin",
// 通过 nginx 暴露的接收端地址(/jcpt/ 反代到 8081
receiverUrl: "https://crm.gykqyy.com/jcpt",
corpId: "wpLgjyawAAeRkCPQMp9-z5q-xEzK64nA",
token: "ca764511bdace06c64f22bb97dd19911bcc4f0c9a07b2211cbbc1d43e9271657",
chunkSize: 300,
includePublicDocs: true,
// 可选:指定库,空数组表示扫描全部业务库
dbAllowlist: ["admin", "corp", "public"],
};
const args = parseArgs(process.argv.slice(2));
const SOURCE_MONGO_URI =
args.sourceMongoUri ||
CONFIG.sourceMongoUri;
const RECEIVER_BASE_URL = args.receiverUrl || CONFIG.receiverUrl;
const MIGRATION_TOKEN = args.token || CONFIG.token;
const CORP_ID = args.corpId || CONFIG.corpId;
const CHUNK_SIZE = Number(args.chunkSize || CONFIG.chunkSize || 300);
const DB_ALLOWLIST = parseCsv(args.dbs, CONFIG.dbAllowlist);
const INCLUDE_PUBLIC_DOCS = parseBool(args.includePublicDocs, CONFIG.includePublicDocs);
if (!SOURCE_MONGO_URI) {
console.error("[sender] Missing source mongo uri");
process.exit(1);
}
if (!RECEIVER_BASE_URL) {
console.error("[sender] Missing receiver url");
process.exit(1);
}
if (!MIGRATION_TOKEN) {
console.error("[sender] Missing migration token");
process.exit(1);
}
if (!CORP_ID) {
console.error("[sender] Missing corpId");
process.exit(1);
}
const filter = {
$or: [
{ corpId: CORP_ID },
{ corpid: CORP_ID },
{ corp_id: CORP_ID },
{ wxCorpid: CORP_ID },
],
};
const PUBLIC_CORP_FIELDS = ["corpId", "corpid", "corp_id", "wxCorpid"];
function buildPublicFilter() {
return {
$and: PUBLIC_CORP_FIELDS.map((field) => ({
$or: [
{ [field]: { $exists: false } },
{ [field]: null },
{ [field]: "" },
],
})),
};
}
function buildMigrationFilter() {
if (!INCLUDE_PUBLIC_DOCS) return filter;
return { $or: [filter, buildPublicFilter()] };
}
const MIGRATION_FILTER = buildMigrationFilter();
async function requestReceiver(path, payload) {
const url = `${RECEIVER_BASE_URL.replace(/\/$/, "")}${path}`;
const res = await axios.post(url, payload, {
timeout: 60000,
headers: {
"x-migration-token": MIGRATION_TOKEN,
"content-type": "application/json",
},
});
if (!res.data || !res.data.success) {
throw new Error(`[sender] Receiver error on ${path}: ${JSON.stringify(res.data)}`);
}
return res.data;
}
function parseCsv(input, defaultList = []) {
if (!input || typeof input !== "string") return defaultList;
return input
.split(",")
.map((s) => s.trim())
.filter(Boolean);
}
function parseBool(input, defaultValue = false) {
if (input === undefined || input === null || input === "") return Boolean(defaultValue);
if (typeof input === "boolean") return input;
const normalized = String(input).trim().toLowerCase();
if (["1", "true", "yes", "y", "on"].includes(normalized)) return true;
if (["0", "false", "no", "n", "off"].includes(normalized)) return false;
return Boolean(defaultValue);
}
function parseArgs(argv) {
const out = {};
for (let i = 0; i < argv.length; i += 1) {
const a = argv[i];
if (a === "--corpId") out.corpId = argv[++i];
else if (a === "--sourceMongoUri") out.sourceMongoUri = argv[++i];
else if (a === "--receiverUrl") out.receiverUrl = argv[++i];
else if (a === "--token") out.token = argv[++i];
else if (a === "--chunkSize") out.chunkSize = argv[++i];
else if (a === "--dbs") out.dbs = argv[++i];
else if (a === "--includePublicDocs") out.includePublicDocs = true;
else if (a === "--noPublicDocs") out.includePublicDocs = false;
}
return out;
}
function shouldSkipDb(name) {
if (["admin", "config", "local"].includes(name)) return name !== "admin";
if (!DB_ALLOWLIST.length) return false;
return !DB_ALLOWLIST.includes(name);
}
async function migrateCollection(db, dbName, collectionName) {
const col = db.collection(collectionName);
const count = await col.countDocuments(MIGRATION_FILTER);
if (!count) return { migrated: 0, skipped: true };
const indexes = await col.indexes();
await requestReceiver("/migrate/indexes", {
dbName,
collectionName,
indexes: indexes.map((idx) => EJSON.serialize(idx, { relaxed: false })),
});
let migrated = 0;
let page = 0;
const cursor = col.find(MIGRATION_FILTER, { noCursorTimeout: true });
const buffer = [];
while (await cursor.hasNext()) {
const doc = await cursor.next();
buffer.push(EJSON.serialize(doc, { relaxed: false }));
if (buffer.length >= CHUNK_SIZE) {
page += 1;
await requestReceiver("/migrate/chunk", {
dbName,
collectionName,
docs: buffer.splice(0, buffer.length),
});
migrated += CHUNK_SIZE;
console.log(`[sender] ${dbName}.${collectionName} page=${page} migrated=${migrated}/${count}`);
}
}
if (buffer.length) {
page += 1;
const lastSize = buffer.length;
await requestReceiver("/migrate/chunk", {
dbName,
collectionName,
docs: buffer.splice(0, buffer.length),
});
migrated += lastSize;
console.log(`[sender] ${dbName}.${collectionName} page=${page} migrated=${migrated}/${count}`);
}
await cursor.close();
return { migrated, skipped: false };
}
async function main() {
const client = new MongoClient(SOURCE_MONGO_URI, {
maxPoolSize: 30,
minPoolSize: 2,
connectTimeoutMS: 30000,
socketTimeoutMS: 60000,
serverSelectionTimeoutMS: 30000,
});
const summary = {
corpId: CORP_ID,
dbs: 0,
collections: 0,
migratedDocs: 0,
};
try {
await client.connect();
await requestReceiver("/migrate/init", { corpId: CORP_ID });
const dbList = await client.db().admin().listDatabases();
for (const dbInfo of dbList.databases || []) {
const dbName = dbInfo.name;
if (shouldSkipDb(dbName)) continue;
const db = client.db(dbName);
const collections = await db.listCollections({}, { nameOnly: true }).toArray();
for (const item of collections) {
const collectionName = item.name;
if (collectionName.startsWith("system.")) continue;
summary.collections += 1;
const result = await migrateCollection(db, dbName, collectionName);
if (!result.skipped) {
summary.migratedDocs += result.migrated;
}
}
summary.dbs += 1;
}
await requestReceiver("/migrate/finish", { corpId: CORP_ID, summary });
console.log("[sender] Migration completed", summary);
} catch (err) {
console.error("[sender] Migration failed", err.message);
process.exitCode = 1;
} finally {
await client.close();
}
}
main();