retry-onboarding-departments.js 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371
  1. const { createHash } = require('crypto');
  2. const { spawnSync } = require('child_process');
  3. // prd 临时补偿任务:入职账户已创建、但主部门同步失败时,重试设置 EIAM 主部门并回写宜搭。
  4. const APP_TYPE = 'APP_MQP6UV1H5S5BDOU68TQD';
  5. const FORM_UUID = 'FORM-3794DBD2C83E42C5863853A80A3C0ACB5K6Q';
  6. const HOST = '120.55.113.155';
  7. const TABLE_FIELD = 'tableField_mrcw2nqc';
  8. const PHONE_FIELD = 'numberField_mrcw2nqe';
  9. const NAME_FIELD = 'textField_mrcw2nqz';
  10. const ORG_FIELD = 'selectField_mrmyvzrt_id';
  11. const STATUS_FIELD = 'selectField_mrcxjcq8';
  12. const TOTAL_FIELD = 'numberField_mrcxjcq9';
  13. const SUCCESS_FIELD = 'numberField_mrcxjcqe';
  14. const FAILED_FIELD = 'numberField_mrcxjcqf';
  15. const FAILURE_MESSAGE_FIELD = 'textareaField_mrn3kh8m';
  16. const REMOTE_PYTHON = String.raw`
  17. import json
  18. import sys
  19. import urllib.parse
  20. import urllib.request
  21. from pathlib import Path
  22. from typing import Any, Dict, List, Set, Tuple
  23. import yaml
  24. CONFIG_PATH = Path('/home/server/benteler/application-prod.yml')
  25. def request_json(request: urllib.request.Request) -> Tuple[Dict[str, Any], int]:
  26. with urllib.request.urlopen(request, timeout=30) as response:
  27. body = response.read().decode('utf-8')
  28. return json.loads(body), response.getcode()
  29. def get_access_token(config: Dict[str, Any]) -> str:
  30. eiam = config['eiam']
  31. url = '{}/v2/{}/{}/oauth2/token'.format(
  32. eiam['baseUrl'], eiam['instanceId'], eiam['applicationId'])
  33. body = urllib.parse.urlencode({
  34. 'grant_type': 'client_credentials',
  35. 'client_id': eiam['clientId'],
  36. 'client_secret': eiam['clientSecret'],
  37. }).encode('utf-8')
  38. payload, _ = request_json(urllib.request.Request(url, data=body, method='POST'))
  39. token = payload.get('access_token') or payload.get('accessToken')
  40. if not token:
  41. raise RuntimeError('EIAM token response did not contain access token')
  42. return token
  43. def list_users(config: Dict[str, Any], token: str) -> Tuple[List[Dict[str, Any]], str]:
  44. eiam = config['eiam']
  45. base = '{}/v2/{}/{}/users'.format(
  46. eiam['baseUrl'], eiam['instanceId'], eiam['applicationId'])
  47. users = []
  48. page_number = 1
  49. while True:
  50. query = urllib.parse.urlencode({'pageNumber': page_number, 'pageSize': 100})
  51. request = urllib.request.Request(
  52. '{}?{}'.format(base, query),
  53. headers={'Authorization': 'Bearer {}'.format(token)},
  54. method='GET')
  55. payload, _ = request_json(request)
  56. page = payload.get('data') or []
  57. users.extend(item for item in page if isinstance(item, dict))
  58. total = int(payload.get('totalCount') or len(users))
  59. if not page or len(users) >= total:
  60. return users, base
  61. page_number += 1
  62. def identifiers(user: Dict[str, Any]) -> Set[str]:
  63. values = (user.get('username'), user.get('phoneNumber'), user.get('userExternalId'))
  64. return set(str(value) for value in values if value not in (None, ''))
  65. def get_user(base: str, token: str, user_id: str) -> Dict[str, Any]:
  66. request = urllib.request.Request(
  67. '{}/{}'.format(base, urllib.parse.quote(str(user_id), safe='')),
  68. headers={'Authorization': 'Bearer {}'.format(token)},
  69. method='GET')
  70. payload, _ = request_json(request)
  71. return payload
  72. def set_primary_org(base: str, token: str, user_id: str, org_id: str) -> int:
  73. url = '{}/{}/actions/setUserPrimaryOrganizationalUnit'.format(
  74. base, urllib.parse.quote(str(user_id), safe=''))
  75. body = json.dumps({'organizationalUnitId': org_id}).encode('utf-8')
  76. request = urllib.request.Request(
  77. url,
  78. data=body,
  79. headers={
  80. 'Authorization': 'Bearer {}'.format(token),
  81. 'Content-Type': 'application/json',
  82. },
  83. method='POST')
  84. with urllib.request.urlopen(request, timeout=30) as response:
  85. response.read()
  86. return response.getcode()
  87. def main() -> None:
  88. tasks = json.loads(sys.stdin.readline())
  89. config = yaml.safe_load(CONFIG_PATH.read_text(encoding='utf-8'))
  90. token = get_access_token(config)
  91. users, base = list_users(config, token)
  92. results = []
  93. for task in tasks:
  94. matched = next((user for user in users if task['username'] in identifiers(user)), None)
  95. if not matched:
  96. results.append({'key': task['key'], 'success': False, 'message': 'EIAM account not found'})
  97. continue
  98. detail = get_user(base, token, matched.get('userId'))
  99. current_org = detail.get('primaryOrganizationalUnitId')
  100. if current_org == task['organizationalUnitId']:
  101. results.append({'key': task['key'], 'success': True, 'changed': False,
  102. 'userId': matched.get('userId'), 'httpStatus': None})
  103. continue
  104. try:
  105. status = set_primary_org(base, token, matched.get('userId'), task['organizationalUnitId'])
  106. results.append({'key': task['key'], 'success': True, 'changed': True,
  107. 'userId': matched.get('userId'), 'httpStatus': status})
  108. except Exception as error:
  109. results.append({'key': task['key'], 'success': False,
  110. 'message': '{}: {}'.format(type(error).__name__, str(error))[:240]})
  111. print(json.dumps({'results': results}, ensure_ascii=False))
  112. if __name__ == '__main__':
  113. main()
  114. `;
  115. /**
  116. * 解析命令行参数。
  117. * @returns {{execute: boolean, scanFailed: boolean, processIds: string[]}} 参数
  118. */
  119. function parseArgs() {
  120. const args = process.argv.slice(2);
  121. const parsed = { execute: false, scanFailed: false, processIds: [] };
  122. for (let index = 0; index < args.length; index += 1) {
  123. if (args[index] === '--execute') parsed.execute = true;
  124. else if (args[index] === '--scan-failed') parsed.scanFailed = true;
  125. else if (args[index] === '--process-id' && args[index + 1]) parsed.processIds.push(args[++index]);
  126. else throw new Error(`unknown or incomplete argument: ${args[index]}`);
  127. }
  128. if (!parsed.scanFailed && parsed.processIds.length === 0) parsed.scanFailed = true;
  129. return parsed;
  130. }
  131. /**
  132. * 从 OpenYida CLI 输出解析 JSON。
  133. * @param {string} output CLI 输出
  134. * @returns {Record<string, unknown>} JSON 响应
  135. */
  136. function parsePayload(output) {
  137. const start = output.indexOf('{\n "content"');
  138. if (start < 0) throw new Error('cannot locate OpenYida JSON response');
  139. return JSON.parse(output.slice(start));
  140. }
  141. /**
  142. * 分页查询全部入职流程。
  143. * @returns {Array<Record<string, unknown>>} 流程实例
  144. */
  145. function queryProcesses() {
  146. const instances = [];
  147. for (let page = 1; ; page += 1) {
  148. const result = spawnSync('openyida', [
  149. 'data', 'query', 'process', APP_TYPE, FORM_UUID,
  150. '--page', String(page), '--size', '100',
  151. ], { encoding: 'utf8' });
  152. if (result.status !== 0) throw new Error(result.stderr || result.stdout);
  153. const content = parsePayload(result.stdout).content || {};
  154. instances.push(...(content.data || []));
  155. if (instances.length >= Number(content.totalCount || instances.length)) return instances;
  156. }
  157. }
  158. /**
  159. * 将字段值归一化为单个字符串。
  160. * @param {unknown} value 字段值
  161. * @returns {string} 单值
  162. */
  163. function firstText(value) {
  164. if (Array.isArray(value)) return value.length ? firstText(value[0]) : '';
  165. return value === null || value === undefined ? '' : String(value).trim();
  166. }
  167. /**
  168. * 移除查询响应中不可回写的 _value 派生字段。
  169. * @param {unknown} value 原始值
  170. * @returns {unknown} 可回写值
  171. */
  172. function stripDerived(value) {
  173. if (Array.isArray(value)) return value.map(stripDerived);
  174. if (!value || typeof value !== 'object') return value;
  175. return Object.fromEntries(Object.entries(value)
  176. .filter(([key]) => !key.endsWith('_value'))
  177. .map(([key, item]) => [key, stripDerived(item)]));
  178. }
  179. /**
  180. * 为去重任务生成不包含人员标识的键。
  181. * @param {string} username EIAM 账户名
  182. * @param {string} organizationalUnitId 目标部门 ID
  183. * @returns {string} 任务键
  184. */
  185. function taskKey(username, organizationalUnitId) {
  186. return createHash('sha256').update(`${username}\0${organizationalUnitId}`).digest('hex').slice(0, 16);
  187. }
  188. /**
  189. * 构建待执行的人员部门修复明细。
  190. * @param {Array<Record<string, unknown>>} instances 流程实例
  191. * @param {{scanFailed: boolean, processIds: string[]}} options 筛选参数
  192. * @returns {Array<Record<string, unknown>>} 修复明细
  193. */
  194. function buildCandidates(instances, options) {
  195. const requested = new Set(options.processIds);
  196. const found = new Set();
  197. const candidates = [];
  198. for (const instance of instances) {
  199. if (requested.size && !requested.has(instance.processInstanceId)) continue;
  200. found.add(instance.processInstanceId);
  201. const data = instance.data || {};
  202. const rows = Array.isArray(data[TABLE_FIELD]) ? data[TABLE_FIELD] : [];
  203. if (rows.length === 50) throw new Error(`${instance.processInstanceId}: subtable may be truncated at 50 rows`);
  204. rows.forEach((row, rowIndex) => {
  205. const status = firstText(row[STATUS_FIELD] || row[`${STATUS_FIELD}_id`]);
  206. if (options.scanFailed && status !== '失败') return;
  207. const username = firstText(row[PHONE_FIELD] ?? row[`${PHONE_FIELD}_value`]);
  208. const organizationalUnitId = firstText(data[ORG_FIELD]);
  209. if (!username || !organizationalUnitId.startsWith('ou_')) {
  210. throw new Error(`${instance.processInstanceId} row ${rowIndex + 1}: missing account or valid department ID`);
  211. }
  212. candidates.push({ instance, rowIndex, username, organizationalUnitId,
  213. key: taskKey(username, organizationalUnitId) });
  214. });
  215. }
  216. const missing = [...requested].filter((id) => !found.has(id));
  217. if (missing.length) throw new Error(`process instances not found: ${missing.join(', ')}`);
  218. return candidates;
  219. }
  220. /**
  221. * 阻止同一账户在同一批次被设置为多个主部门。
  222. * @param {Array<Record<string, unknown>>} candidates 修复明细
  223. * @returns {void}
  224. */
  225. function assertNoConflicts(candidates) {
  226. const byUser = new Map();
  227. for (const item of candidates) {
  228. if (!byUser.has(item.username)) byUser.set(item.username, new Set());
  229. byUser.get(item.username).add(item.organizationalUnitId);
  230. }
  231. const conflicts = [...byUser.entries()].filter(([, orgs]) => orgs.size > 1);
  232. if (conflicts.length) {
  233. throw new Error(`department conflicts: ${conflicts.map(([user, orgs]) =>
  234. `${mask(user)} => ${[...orgs].join(',')}`).join('; ')}`);
  235. }
  236. }
  237. /**
  238. * 遮蔽账户标识。
  239. * @param {string} value 账户标识
  240. * @returns {string} 遮蔽值
  241. */
  242. function mask(value) {
  243. return value.length >= 4 ? `***${value.slice(-4)}` : '***';
  244. }
  245. /**
  246. * 通过 SSH 在服务器内调用 EIAM,凭据只在远端进程内使用。
  247. * @param {Array<Record<string, unknown>>} tasks 去重后任务
  248. * @returns {Array<Record<string, unknown>>} EIAM 执行结果
  249. */
  250. function runRemoteTasks(tasks) {
  251. const encoded = Buffer.from(REMOTE_PYTHON, 'utf8').toString('base64');
  252. const command = `python3 -c "import base64;exec(base64.b64decode('${encoded}'))"`;
  253. const result = spawnSync('ssh', [
  254. '-o', 'BatchMode=yes', '-o', 'ConnectTimeout=8', '-p', '22', `root@${HOST}`, command,
  255. ], { encoding: 'utf8', input: `${JSON.stringify(tasks)}\n`, maxBuffer: 10 * 1024 * 1024 });
  256. if (result.status !== 0) throw new Error(result.stderr || result.stdout || `ssh exited ${result.status}`);
  257. return JSON.parse(result.stdout).results || [];
  258. }
  259. /**
  260. * 将 EIAM 执行结果回写到流程实例。
  261. * @param {Array<Record<string, unknown>>} candidates 修复明细
  262. * @param {Array<Record<string, unknown>>} remoteResults EIAM 结果
  263. * @returns {Array<Record<string, unknown>>} 回写摘要
  264. */
  265. function writeBack(candidates, remoteResults) {
  266. const resultByKey = new Map(remoteResults.map((item) => [item.key, item]));
  267. const byProcess = new Map();
  268. for (const candidate of candidates) {
  269. if (!byProcess.has(candidate.instance.processInstanceId)) byProcess.set(candidate.instance.processInstanceId, []);
  270. byProcess.get(candidate.instance.processInstanceId).push(candidate);
  271. }
  272. const summaries = [];
  273. for (const [processId, items] of byProcess) {
  274. summaries.push(updateOneProcess(processId, items, resultByKey));
  275. }
  276. return summaries;
  277. }
  278. /**
  279. * 回写单个流程的全量子表和统计。
  280. * @param {string} processId 流程实例 ID
  281. * @param {Array<Record<string, unknown>>} items 当前流程修复明细
  282. * @param {Map<string, Record<string, unknown>>} resultByKey EIAM 结果
  283. * @returns {Record<string, unknown>} 回写摘要
  284. */
  285. function updateOneProcess(processId, items, resultByKey) {
  286. const data = items[0].instance.data || {};
  287. const rows = stripDerived(data[TABLE_FIELD] || []);
  288. const errors = [];
  289. for (const item of items) {
  290. const remote = resultByKey.get(item.key);
  291. const success = Boolean(remote?.success);
  292. rows[item.rowIndex][STATUS_FIELD] = success ? '成功' : '失败';
  293. rows[item.rowIndex][`${STATUS_FIELD}_id`] = success ? '成功' : '失败';
  294. if (!success) errors.push(`${firstText(rows[item.rowIndex][NAME_FIELD]) || `第 ${item.rowIndex + 1} 行`} - ${remote?.message || '部门更新失败'}`);
  295. }
  296. const successCount = rows.filter((row) => firstText(row[STATUS_FIELD]) === '成功').length;
  297. const failedCount = rows.filter((row) => firstText(row[STATUS_FIELD]) === '失败').length;
  298. const patch = { [TABLE_FIELD]: rows, [TOTAL_FIELD]: rows.length,
  299. [SUCCESS_FIELD]: successCount, [FAILED_FIELD]: failedCount,
  300. [FAILURE_MESSAGE_FIELD]: failedCount === 0 ? '' : errors.join('\n') || data[FAILURE_MESSAGE_FIELD] || '' };
  301. const result = spawnSync('openyida', ['data', 'update', 'process', APP_TYPE,
  302. '--process-inst-id', processId, '--form-uuid', FORM_UUID,
  303. '--data-json', JSON.stringify(patch)], { encoding: 'utf8' });
  304. if (result.status !== 0) throw new Error(result.stderr || result.stdout);
  305. return { processInstanceId: processId, total: rows.length, success: successCount, failed: failedCount };
  306. }
  307. /**
  308. * 验证目标流程回写结果。
  309. * @param {string[]} processIds 目标流程 ID
  310. * @returns {Array<Record<string, unknown>>} 验证摘要
  311. */
  312. function verify(processIds) {
  313. const wanted = new Set(processIds);
  314. return queryProcesses().filter((instance) => wanted.has(instance.processInstanceId)).map((instance) => {
  315. const data = instance.data || {};
  316. const rows = Array.isArray(data[TABLE_FIELD]) ? data[TABLE_FIELD] : [];
  317. return { processInstanceId: instance.processInstanceId,
  318. total: Number(data[TOTAL_FIELD] ?? data[`${TOTAL_FIELD}_value`] ?? -1),
  319. success: Number(data[SUCCESS_FIELD] ?? data[`${SUCCESS_FIELD}_value`] ?? -1),
  320. failed: Number(data[FAILED_FIELD] ?? data[`${FAILED_FIELD}_value`] ?? -1),
  321. detailStatuses: rows.map((row) => firstText(row[STATUS_FIELD] || row[`${STATUS_FIELD}_id`])),
  322. failureMessageEmpty: !data[FAILURE_MESSAGE_FIELD] };
  323. });
  324. }
  325. const options = parseArgs();
  326. const instances = queryProcesses();
  327. const candidates = buildCandidates(instances, options);
  328. assertNoConflicts(candidates);
  329. const uniqueTasks = [...new Map(candidates.map((item) => [item.key, {
  330. key: item.key, username: item.username, organizationalUnitId: item.organizationalUnitId,
  331. }])).values()];
  332. process.stdout.write(`${JSON.stringify({ mode: options.execute ? 'execute' : 'dry-run',
  333. processes: new Set(candidates.map((item) => item.instance.processInstanceId)).size,
  334. rows: candidates.length, uniqueUsers: uniqueTasks.length,
  335. candidates: candidates.map((item) => ({ processInstanceId: item.instance.processInstanceId,
  336. row: item.rowIndex + 1, account: mask(item.username), organizationalUnitId: item.organizationalUnitId })) }, null, 2)}\n`);
  337. if (options.execute && uniqueTasks.length) {
  338. const remoteResults = runRemoteTasks(uniqueTasks);
  339. const writeback = writeBack(candidates, remoteResults);
  340. const verification = verify([...new Set(candidates.map((item) => item.instance.processInstanceId))]);
  341. process.stdout.write(`${JSON.stringify({ remoteResults, writeback, verification }, null, 2)}\n`);
  342. if (remoteResults.some((item) => !item.success)) process.exitCode = 2;
  343. }