| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371 |
- const { createHash } = require('crypto');
- const { spawnSync } = require('child_process');
- // prd 临时补偿任务:入职账户已创建、但主部门同步失败时,重试设置 EIAM 主部门并回写宜搭。
- const APP_TYPE = 'APP_MQP6UV1H5S5BDOU68TQD';
- const FORM_UUID = 'FORM-3794DBD2C83E42C5863853A80A3C0ACB5K6Q';
- const HOST = '120.55.113.155';
- const TABLE_FIELD = 'tableField_mrcw2nqc';
- const PHONE_FIELD = 'numberField_mrcw2nqe';
- const NAME_FIELD = 'textField_mrcw2nqz';
- const ORG_FIELD = 'selectField_mrmyvzrt_id';
- const STATUS_FIELD = 'selectField_mrcxjcq8';
- const TOTAL_FIELD = 'numberField_mrcxjcq9';
- const SUCCESS_FIELD = 'numberField_mrcxjcqe';
- const FAILED_FIELD = 'numberField_mrcxjcqf';
- const FAILURE_MESSAGE_FIELD = 'textareaField_mrn3kh8m';
- const REMOTE_PYTHON = String.raw`
- import json
- import sys
- import urllib.parse
- import urllib.request
- from pathlib import Path
- from typing import Any, Dict, List, Set, Tuple
- import yaml
- CONFIG_PATH = Path('/home/server/benteler/application-prod.yml')
- def request_json(request: urllib.request.Request) -> Tuple[Dict[str, Any], int]:
- with urllib.request.urlopen(request, timeout=30) as response:
- body = response.read().decode('utf-8')
- return json.loads(body), response.getcode()
- def get_access_token(config: Dict[str, Any]) -> str:
- eiam = config['eiam']
- url = '{}/v2/{}/{}/oauth2/token'.format(
- eiam['baseUrl'], eiam['instanceId'], eiam['applicationId'])
- body = urllib.parse.urlencode({
- 'grant_type': 'client_credentials',
- 'client_id': eiam['clientId'],
- 'client_secret': eiam['clientSecret'],
- }).encode('utf-8')
- payload, _ = request_json(urllib.request.Request(url, data=body, method='POST'))
- token = payload.get('access_token') or payload.get('accessToken')
- if not token:
- raise RuntimeError('EIAM token response did not contain access token')
- return token
- def list_users(config: Dict[str, Any], token: str) -> Tuple[List[Dict[str, Any]], str]:
- eiam = config['eiam']
- base = '{}/v2/{}/{}/users'.format(
- eiam['baseUrl'], eiam['instanceId'], eiam['applicationId'])
- users = []
- page_number = 1
- while True:
- query = urllib.parse.urlencode({'pageNumber': page_number, 'pageSize': 100})
- request = urllib.request.Request(
- '{}?{}'.format(base, query),
- headers={'Authorization': 'Bearer {}'.format(token)},
- method='GET')
- payload, _ = request_json(request)
- page = payload.get('data') or []
- users.extend(item for item in page if isinstance(item, dict))
- total = int(payload.get('totalCount') or len(users))
- if not page or len(users) >= total:
- return users, base
- page_number += 1
- def identifiers(user: Dict[str, Any]) -> Set[str]:
- values = (user.get('username'), user.get('phoneNumber'), user.get('userExternalId'))
- return set(str(value) for value in values if value not in (None, ''))
- def get_user(base: str, token: str, user_id: str) -> Dict[str, Any]:
- request = urllib.request.Request(
- '{}/{}'.format(base, urllib.parse.quote(str(user_id), safe='')),
- headers={'Authorization': 'Bearer {}'.format(token)},
- method='GET')
- payload, _ = request_json(request)
- return payload
- def set_primary_org(base: str, token: str, user_id: str, org_id: str) -> int:
- url = '{}/{}/actions/setUserPrimaryOrganizationalUnit'.format(
- base, urllib.parse.quote(str(user_id), safe=''))
- body = json.dumps({'organizationalUnitId': org_id}).encode('utf-8')
- request = urllib.request.Request(
- url,
- data=body,
- headers={
- 'Authorization': 'Bearer {}'.format(token),
- 'Content-Type': 'application/json',
- },
- method='POST')
- with urllib.request.urlopen(request, timeout=30) as response:
- response.read()
- return response.getcode()
- def main() -> None:
- tasks = json.loads(sys.stdin.readline())
- config = yaml.safe_load(CONFIG_PATH.read_text(encoding='utf-8'))
- token = get_access_token(config)
- users, base = list_users(config, token)
- results = []
- for task in tasks:
- matched = next((user for user in users if task['username'] in identifiers(user)), None)
- if not matched:
- results.append({'key': task['key'], 'success': False, 'message': 'EIAM account not found'})
- continue
- detail = get_user(base, token, matched.get('userId'))
- current_org = detail.get('primaryOrganizationalUnitId')
- if current_org == task['organizationalUnitId']:
- results.append({'key': task['key'], 'success': True, 'changed': False,
- 'userId': matched.get('userId'), 'httpStatus': None})
- continue
- try:
- status = set_primary_org(base, token, matched.get('userId'), task['organizationalUnitId'])
- results.append({'key': task['key'], 'success': True, 'changed': True,
- 'userId': matched.get('userId'), 'httpStatus': status})
- except Exception as error:
- results.append({'key': task['key'], 'success': False,
- 'message': '{}: {}'.format(type(error).__name__, str(error))[:240]})
- print(json.dumps({'results': results}, ensure_ascii=False))
- if __name__ == '__main__':
- main()
- `;
- /**
- * 解析命令行参数。
- * @returns {{execute: boolean, scanFailed: boolean, processIds: string[]}} 参数
- */
- function parseArgs() {
- const args = process.argv.slice(2);
- const parsed = { execute: false, scanFailed: false, processIds: [] };
- for (let index = 0; index < args.length; index += 1) {
- if (args[index] === '--execute') parsed.execute = true;
- else if (args[index] === '--scan-failed') parsed.scanFailed = true;
- else if (args[index] === '--process-id' && args[index + 1]) parsed.processIds.push(args[++index]);
- else throw new Error(`unknown or incomplete argument: ${args[index]}`);
- }
- if (!parsed.scanFailed && parsed.processIds.length === 0) parsed.scanFailed = true;
- return parsed;
- }
- /**
- * 从 OpenYida CLI 输出解析 JSON。
- * @param {string} output CLI 输出
- * @returns {Record<string, unknown>} JSON 响应
- */
- function parsePayload(output) {
- const start = output.indexOf('{\n "content"');
- if (start < 0) throw new Error('cannot locate OpenYida JSON response');
- return JSON.parse(output.slice(start));
- }
- /**
- * 分页查询全部入职流程。
- * @returns {Array<Record<string, unknown>>} 流程实例
- */
- function queryProcesses() {
- const instances = [];
- for (let page = 1; ; page += 1) {
- const result = spawnSync('openyida', [
- 'data', 'query', 'process', APP_TYPE, FORM_UUID,
- '--page', String(page), '--size', '100',
- ], { encoding: 'utf8' });
- if (result.status !== 0) throw new Error(result.stderr || result.stdout);
- const content = parsePayload(result.stdout).content || {};
- instances.push(...(content.data || []));
- if (instances.length >= Number(content.totalCount || instances.length)) return instances;
- }
- }
- /**
- * 将字段值归一化为单个字符串。
- * @param {unknown} value 字段值
- * @returns {string} 单值
- */
- function firstText(value) {
- if (Array.isArray(value)) return value.length ? firstText(value[0]) : '';
- return value === null || value === undefined ? '' : String(value).trim();
- }
- /**
- * 移除查询响应中不可回写的 _value 派生字段。
- * @param {unknown} value 原始值
- * @returns {unknown} 可回写值
- */
- function stripDerived(value) {
- if (Array.isArray(value)) return value.map(stripDerived);
- if (!value || typeof value !== 'object') return value;
- return Object.fromEntries(Object.entries(value)
- .filter(([key]) => !key.endsWith('_value'))
- .map(([key, item]) => [key, stripDerived(item)]));
- }
- /**
- * 为去重任务生成不包含人员标识的键。
- * @param {string} username EIAM 账户名
- * @param {string} organizationalUnitId 目标部门 ID
- * @returns {string} 任务键
- */
- function taskKey(username, organizationalUnitId) {
- return createHash('sha256').update(`${username}\0${organizationalUnitId}`).digest('hex').slice(0, 16);
- }
- /**
- * 构建待执行的人员部门修复明细。
- * @param {Array<Record<string, unknown>>} instances 流程实例
- * @param {{scanFailed: boolean, processIds: string[]}} options 筛选参数
- * @returns {Array<Record<string, unknown>>} 修复明细
- */
- function buildCandidates(instances, options) {
- const requested = new Set(options.processIds);
- const found = new Set();
- const candidates = [];
- for (const instance of instances) {
- if (requested.size && !requested.has(instance.processInstanceId)) continue;
- found.add(instance.processInstanceId);
- const data = instance.data || {};
- const rows = Array.isArray(data[TABLE_FIELD]) ? data[TABLE_FIELD] : [];
- if (rows.length === 50) throw new Error(`${instance.processInstanceId}: subtable may be truncated at 50 rows`);
- rows.forEach((row, rowIndex) => {
- const status = firstText(row[STATUS_FIELD] || row[`${STATUS_FIELD}_id`]);
- if (options.scanFailed && status !== '失败') return;
- const username = firstText(row[PHONE_FIELD] ?? row[`${PHONE_FIELD}_value`]);
- const organizationalUnitId = firstText(data[ORG_FIELD]);
- if (!username || !organizationalUnitId.startsWith('ou_')) {
- throw new Error(`${instance.processInstanceId} row ${rowIndex + 1}: missing account or valid department ID`);
- }
- candidates.push({ instance, rowIndex, username, organizationalUnitId,
- key: taskKey(username, organizationalUnitId) });
- });
- }
- const missing = [...requested].filter((id) => !found.has(id));
- if (missing.length) throw new Error(`process instances not found: ${missing.join(', ')}`);
- return candidates;
- }
- /**
- * 阻止同一账户在同一批次被设置为多个主部门。
- * @param {Array<Record<string, unknown>>} candidates 修复明细
- * @returns {void}
- */
- function assertNoConflicts(candidates) {
- const byUser = new Map();
- for (const item of candidates) {
- if (!byUser.has(item.username)) byUser.set(item.username, new Set());
- byUser.get(item.username).add(item.organizationalUnitId);
- }
- const conflicts = [...byUser.entries()].filter(([, orgs]) => orgs.size > 1);
- if (conflicts.length) {
- throw new Error(`department conflicts: ${conflicts.map(([user, orgs]) =>
- `${mask(user)} => ${[...orgs].join(',')}`).join('; ')}`);
- }
- }
- /**
- * 遮蔽账户标识。
- * @param {string} value 账户标识
- * @returns {string} 遮蔽值
- */
- function mask(value) {
- return value.length >= 4 ? `***${value.slice(-4)}` : '***';
- }
- /**
- * 通过 SSH 在服务器内调用 EIAM,凭据只在远端进程内使用。
- * @param {Array<Record<string, unknown>>} tasks 去重后任务
- * @returns {Array<Record<string, unknown>>} EIAM 执行结果
- */
- function runRemoteTasks(tasks) {
- const encoded = Buffer.from(REMOTE_PYTHON, 'utf8').toString('base64');
- const command = `python3 -c "import base64;exec(base64.b64decode('${encoded}'))"`;
- const result = spawnSync('ssh', [
- '-o', 'BatchMode=yes', '-o', 'ConnectTimeout=8', '-p', '22', `root@${HOST}`, command,
- ], { encoding: 'utf8', input: `${JSON.stringify(tasks)}\n`, maxBuffer: 10 * 1024 * 1024 });
- if (result.status !== 0) throw new Error(result.stderr || result.stdout || `ssh exited ${result.status}`);
- return JSON.parse(result.stdout).results || [];
- }
- /**
- * 将 EIAM 执行结果回写到流程实例。
- * @param {Array<Record<string, unknown>>} candidates 修复明细
- * @param {Array<Record<string, unknown>>} remoteResults EIAM 结果
- * @returns {Array<Record<string, unknown>>} 回写摘要
- */
- function writeBack(candidates, remoteResults) {
- const resultByKey = new Map(remoteResults.map((item) => [item.key, item]));
- const byProcess = new Map();
- for (const candidate of candidates) {
- if (!byProcess.has(candidate.instance.processInstanceId)) byProcess.set(candidate.instance.processInstanceId, []);
- byProcess.get(candidate.instance.processInstanceId).push(candidate);
- }
- const summaries = [];
- for (const [processId, items] of byProcess) {
- summaries.push(updateOneProcess(processId, items, resultByKey));
- }
- return summaries;
- }
- /**
- * 回写单个流程的全量子表和统计。
- * @param {string} processId 流程实例 ID
- * @param {Array<Record<string, unknown>>} items 当前流程修复明细
- * @param {Map<string, Record<string, unknown>>} resultByKey EIAM 结果
- * @returns {Record<string, unknown>} 回写摘要
- */
- function updateOneProcess(processId, items, resultByKey) {
- const data = items[0].instance.data || {};
- const rows = stripDerived(data[TABLE_FIELD] || []);
- const errors = [];
- for (const item of items) {
- const remote = resultByKey.get(item.key);
- const success = Boolean(remote?.success);
- rows[item.rowIndex][STATUS_FIELD] = success ? '成功' : '失败';
- rows[item.rowIndex][`${STATUS_FIELD}_id`] = success ? '成功' : '失败';
- if (!success) errors.push(`${firstText(rows[item.rowIndex][NAME_FIELD]) || `第 ${item.rowIndex + 1} 行`} - ${remote?.message || '部门更新失败'}`);
- }
- const successCount = rows.filter((row) => firstText(row[STATUS_FIELD]) === '成功').length;
- const failedCount = rows.filter((row) => firstText(row[STATUS_FIELD]) === '失败').length;
- const patch = { [TABLE_FIELD]: rows, [TOTAL_FIELD]: rows.length,
- [SUCCESS_FIELD]: successCount, [FAILED_FIELD]: failedCount,
- [FAILURE_MESSAGE_FIELD]: failedCount === 0 ? '' : errors.join('\n') || data[FAILURE_MESSAGE_FIELD] || '' };
- const result = spawnSync('openyida', ['data', 'update', 'process', APP_TYPE,
- '--process-inst-id', processId, '--form-uuid', FORM_UUID,
- '--data-json', JSON.stringify(patch)], { encoding: 'utf8' });
- if (result.status !== 0) throw new Error(result.stderr || result.stdout);
- return { processInstanceId: processId, total: rows.length, success: successCount, failed: failedCount };
- }
- /**
- * 验证目标流程回写结果。
- * @param {string[]} processIds 目标流程 ID
- * @returns {Array<Record<string, unknown>>} 验证摘要
- */
- function verify(processIds) {
- const wanted = new Set(processIds);
- return queryProcesses().filter((instance) => wanted.has(instance.processInstanceId)).map((instance) => {
- const data = instance.data || {};
- const rows = Array.isArray(data[TABLE_FIELD]) ? data[TABLE_FIELD] : [];
- return { processInstanceId: instance.processInstanceId,
- total: Number(data[TOTAL_FIELD] ?? data[`${TOTAL_FIELD}_value`] ?? -1),
- success: Number(data[SUCCESS_FIELD] ?? data[`${SUCCESS_FIELD}_value`] ?? -1),
- failed: Number(data[FAILED_FIELD] ?? data[`${FAILED_FIELD}_value`] ?? -1),
- detailStatuses: rows.map((row) => firstText(row[STATUS_FIELD] || row[`${STATUS_FIELD}_id`])),
- failureMessageEmpty: !data[FAILURE_MESSAGE_FIELD] };
- });
- }
- const options = parseArgs();
- const instances = queryProcesses();
- const candidates = buildCandidates(instances, options);
- assertNoConflicts(candidates);
- const uniqueTasks = [...new Map(candidates.map((item) => [item.key, {
- key: item.key, username: item.username, organizationalUnitId: item.organizationalUnitId,
- }])).values()];
- process.stdout.write(`${JSON.stringify({ mode: options.execute ? 'execute' : 'dry-run',
- processes: new Set(candidates.map((item) => item.instance.processInstanceId)).size,
- rows: candidates.length, uniqueUsers: uniqueTasks.length,
- candidates: candidates.map((item) => ({ processInstanceId: item.instance.processInstanceId,
- row: item.rowIndex + 1, account: mask(item.username), organizationalUnitId: item.organizationalUnitId })) }, null, 2)}\n`);
- if (options.execute && uniqueTasks.length) {
- const remoteResults = runRemoteTasks(uniqueTasks);
- const writeback = writeBack(candidates, remoteResults);
- const verification = verify([...new Set(candidates.map((item) => item.instance.processInstanceId))]);
- process.stdout.write(`${JSON.stringify({ remoteResults, writeback, verification }, null, 2)}\n`);
- if (remoteResults.some((item) => !item.success)) process.exitCode = 2;
- }
|