feat(aihr): add open org snapshot sync endpoint

新增 POST /api/aihr/org/sync,从开放组织系统 /open/v1/sync/snapshot
拉取公司/部门/员工快照刷新 aihr_org_snapshot;支持 dryRun 预检、
Bearer/client-credentials 鉴权与可选 HMAC 签名;外部未配置时保留
SQL seed。附接口契约设计 v1 与 API/DEV_SETUP 配置说明。
This commit is contained in:
2026-07-07 11:05:41 +08:00
parent f84f5c6530
commit f5137a2b05
6 changed files with 2259 additions and 3 deletions
@@ -0,0 +1,24 @@
package org.dromara.aihr.controller;
import lombok.RequiredArgsConstructor;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncResponse;
import org.dromara.aihr.service.AihrOrgSyncService;
import org.dromara.common.core.domain.R;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RequiredArgsConstructor
@RestController
@RequestMapping("/api/aihr/org")
public class AihrOrgSyncController {
private final AihrOrgSyncService orgSyncService;
@PostMapping("/sync")
public R<SyncResponse> sync(@RequestBody(required = false) SyncRequest request) {
return R.ok(orgSyncService.sync(request));
}
}
@@ -0,0 +1,33 @@
package org.dromara.aihr.domain;
import java.util.List;
public final class AihrOrgSyncDto {
private AihrOrgSyncDto() {
}
public record SyncRequest(
Boolean dryRun,
Boolean replaceExisting,
Integer pageSize,
Integer maxPages,
String groupId,
String companyId,
String departmentId
) {
}
public record SyncResponse(
boolean dryRun,
boolean replaceExisting,
String source,
int companyCount,
int departmentCount,
int employeeCount,
int syncedCount,
int skippedCount,
List<String> warnings
) {
}
}
@@ -0,0 +1,470 @@
package org.dromara.aihr.service;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncResponse;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.dao.DataAccessException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.support.TransactionTemplate;
import javax.crypto.Mac;
import javax.crypto.spec.SecretKeySpec;
import java.net.URI;
import java.net.URLEncoder;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.time.Duration;
import java.time.LocalDate;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Base64;
import java.util.List;
import java.util.Map;
import java.util.TreeMap;
import java.util.UUID;
@Service
@RequiredArgsConstructor
@Slf4j
public class AihrOrgSyncService {
private static final String TENANT_ID = "000000";
private static final int DEFAULT_PAGE_SIZE = 200;
private static final int MAX_PAGE_SIZE = 500;
private static final int DEFAULT_MAX_PAGES = 200;
private static final Duration HTTP_TIMEOUT = Duration.ofSeconds(30);
private final ObjectMapper objectMapper;
private final JdbcTemplate jdbcTemplate;
private final TransactionTemplate transactionTemplate;
@Value("${aihr.org-sync.base-url:${AIHR_ORG_SYNC_BASE_URL:}}")
private String configuredBaseUrl;
@Value("${aihr.org-sync.access-token:${AIHR_ORG_SYNC_ACCESS_TOKEN:}}")
private String configuredAccessToken;
@Value("${aihr.org-sync.client-id:${AIHR_ORG_SYNC_CLIENT_ID:}}")
private String configuredClientId;
@Value("${aihr.org-sync.client-secret:${AIHR_ORG_SYNC_CLIENT_SECRET:}}")
private String configuredClientSecret;
@Value("${aihr.org-sync.signing-secret:${AIHR_ORG_SYNC_SIGNING_SECRET:}}")
private String configuredSigningSecret;
public SyncResponse sync(SyncRequest request) {
SyncRequest req = request == null ? new SyncRequest(null, null, null, null, null, null, null) : request;
String baseUrl = normalizeBaseUrl(configuredBaseUrl);
if (baseUrl.isBlank()) {
throw new IllegalArgumentException("请先配置 AIHR_ORG_SYNC_BASE_URL,值为外部开放平台 /open/v1 前缀");
}
requireSnapshotTable();
int pageSize = clamp(req.pageSize(), DEFAULT_PAGE_SIZE, 1, MAX_PAGE_SIZE);
int maxPages = clamp(req.maxPages(), DEFAULT_MAX_PAGES, 1, DEFAULT_MAX_PAGES);
boolean dryRun = Boolean.TRUE.equals(req.dryRun());
boolean replaceExisting = req.replaceExisting() == null || Boolean.TRUE.equals(req.replaceExisting());
String token = accessToken(baseUrl);
List<String> warnings = new ArrayList<>();
List<JsonNode> companyItems = fetchOptional(baseUrl, token, req, "company", pageSize, maxPages, warnings);
List<JsonNode> departmentItems = fetchOptional(baseUrl, token, req, "department", pageSize, maxPages, warnings);
List<JsonNode> employeeItems = fetchSnapshot(baseUrl, token, req, "employee", pageSize, maxPages);
Map<String, CompanyInfo> companies = companyMap(companyItems);
Map<String, DepartmentInfo> departments = departmentMap(departmentItems);
List<OrgRow> rows = new ArrayList<>();
int skipped = 0;
for (JsonNode employee : employeeItems) {
OrgRow row = orgRow(employee, companies, departments);
if (row == null) {
skipped++;
} else {
rows.add(row);
}
}
if (!dryRun && replaceExisting && rows.isEmpty()) {
throw new IllegalArgumentException("外部员工快照为空,已阻止覆盖本地组织人员快照");
}
if (!dryRun && !rows.isEmpty()) {
saveRows(rows, replaceExisting);
}
return new SyncResponse(
dryRun,
replaceExisting,
baseUrl,
companyItems.size(),
departmentItems.size(),
employeeItems.size(),
dryRun ? 0 : rows.size(),
skipped,
List.copyOf(warnings)
);
}
private List<JsonNode> fetchOptional(String baseUrl, String token, SyncRequest req, String resourceType,
int pageSize, int maxPages, List<String> warnings) {
try {
return fetchSnapshot(baseUrl, token, req, resourceType, pageSize, maxPages);
} catch (Exception e) {
String warning = "外部 " + resourceType + " 快照拉取失败,已用员工字段兜底: " + e.getMessage();
warnings.add(warning);
log.warn(warning);
return List.of();
}
}
private List<JsonNode> fetchSnapshot(String baseUrl, String token, SyncRequest req, String resourceType,
int pageSize, int maxPages) {
List<JsonNode> items = new ArrayList<>();
for (int page = 1; page <= maxPages; page++) {
Map<String, String> query = new TreeMap<>();
query.put("resource_type", resourceType);
query.put("page", String.valueOf(page));
query.put("page_size", String.valueOf(pageSize));
putIfNotBlank(query, "group_id", req.groupId());
putIfNotBlank(query, "company_id", req.companyId());
putIfNotBlank(query, "department_id", req.departmentId());
JsonNode data = requestJson(baseUrl, "GET", "/sync/snapshot", query, "", token).path("data");
JsonNode pageItems = data.path("items");
if (pageItems.isMissingNode() && data.isArray()) {
pageItems = data;
}
if (!pageItems.isArray()) {
throw new IllegalStateException("外部 " + resourceType + " 快照响应缺少 data.items");
}
for (JsonNode item : pageItems) {
items.add(item);
}
int total = data.path("total").asInt(-1);
boolean hasMore = data.path("has_more").asBoolean(data.path("hasMore").asBoolean(false));
if (!hasMore && (total >= 0 ? page * pageSize >= total : pageItems.size() < pageSize)) {
return items;
}
}
throw new IllegalStateException("外部 " + resourceType + " 快照超过最大分页 " + maxPages);
}
private String accessToken(String baseUrl) {
String token = clean(configuredAccessToken);
if (!token.isBlank()) {
return token;
}
String clientId = clean(configuredClientId);
String clientSecret = clean(configuredClientSecret);
if (clientId.isBlank() || clientSecret.isBlank()) {
return "";
}
try {
String body = objectMapper.writeValueAsString(Map.of(
"grant_type", "client_credentials",
"client_id", clientId,
"client_secret", clientSecret
));
JsonNode root = requestJson(baseUrl, "POST", "/auth/token", Map.of(), body, "");
String fetchedToken = clean(root.path("data").path("access_token").asText(root.path("access_token").asText("")));
if (fetchedToken.isBlank()) {
throw new IllegalStateException("外部令牌接口未返回 access_token");
}
return fetchedToken;
} catch (Exception e) {
throw new IllegalStateException("外部组织同步令牌获取失败: " + e.getMessage(), e);
}
}
private JsonNode requestJson(String baseUrl, String method, String path, Map<String, String> query, String body, String token) {
try {
String queryString = queryString(query);
URI uri = URI.create(baseUrl + path + (queryString.isBlank() ? "" : "?" + queryString));
HttpRequest.Builder builder = HttpRequest.newBuilder(uri)
.timeout(HTTP_TIMEOUT)
.header("Accept", "application/json");
if (!clean(token).isBlank()) {
builder.header("Authorization", bearer(token));
}
String requestBody = body == null ? "" : body;
if ("POST".equals(method)) {
builder.header("Content-Type", "application/json").POST(HttpRequest.BodyPublishers.ofString(requestBody));
} else {
builder.GET();
}
sign(builder, method, uri.getRawPath(), queryString, requestBody);
HttpResponse<String> response = HttpClient.newBuilder()
.connectTimeout(HTTP_TIMEOUT)
.build()
.send(builder.build(), HttpResponse.BodyHandlers.ofString(StandardCharsets.UTF_8));
if (response.statusCode() < 200 || response.statusCode() >= 300) {
throw new IllegalStateException("HTTP " + response.statusCode() + ": " + truncate(response.body(), 200));
}
JsonNode root = objectMapper.readTree(response.body());
if (root.has("success") && !root.path("success").asBoolean()) {
throw new IllegalStateException(root.path("message").asText("外部接口返回失败"));
}
return root;
} catch (Exception e) {
throw new IllegalStateException(e.getMessage(), e);
}
}
private void sign(HttpRequest.Builder builder, String method, String path, String queryString, String body) throws Exception {
String signingSecret = clean(configuredSigningSecret);
String clientId = clean(configuredClientId);
if (signingSecret.isBlank() || clientId.isBlank()) {
return;
}
String timestamp = Instant.now().toString();
String nonce = UUID.randomUUID().toString();
String plain = method + "\n" + path + "\n" + queryString + "\n" + sha256(body) + "\n" + timestamp + "\n" + nonce;
Mac mac = Mac.getInstance("HmacSHA256");
mac.init(new SecretKeySpec(signingSecret.getBytes(StandardCharsets.UTF_8), "HmacSHA256"));
builder.header("X-Client-Id", clientId)
.header("X-Timestamp", timestamp)
.header("X-Nonce", nonce)
.header("X-Signature", Base64.getEncoder().encodeToString(mac.doFinal(plain.getBytes(StandardCharsets.UTF_8))));
}
private void saveRows(List<OrgRow> rows, boolean replaceExisting) {
transactionTemplate.executeWithoutResult(status -> {
if (replaceExisting) {
jdbcTemplate.update("delete from aihr_org_snapshot where tenant_id = ?", TENANT_ID);
}
jdbcTemplate.batchUpdate("""
insert into aihr_org_snapshot
(tenant_id, project_code, project_name, dept_name, ext_party_id, person_name,
position_name, position_level, employment_status, snapshot_date, create_time)
values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, now())
on duplicate key update
project_code = values(project_code),
project_name = values(project_name),
dept_name = values(dept_name),
person_name = values(person_name),
position_name = values(position_name),
position_level = values(position_level),
employment_status = values(employment_status),
snapshot_date = values(snapshot_date),
create_time = now()
""", rows.stream().map(OrgRow::args).toList());
});
}
private OrgRow orgRow(JsonNode employee, Map<String, CompanyInfo> companies, Map<String, DepartmentInfo> departments) {
String extPartyId = firstNonBlank(text(employee, "employee_number", "employeeNo", "employee_id", "employeeId", "id", "user_id"));
if (extPartyId.isBlank()) {
return null;
}
String departmentId = firstNonBlank(text(employee, "department_id", "departmentId", "dept_id", "deptId"));
DepartmentInfo dept = departments.get(departmentId);
String companyId = firstNonBlank(text(employee, "company_id", "companyId"), dept == null ? "" : dept.companyId());
CompanyInfo company = companies.get(companyId);
String deptName = firstNonBlank(text(employee, "department_name", "departmentName", "dept_name", "deptName"), dept == null ? "" : dept.name());
String projectCode = firstNonBlank(
text(employee, "project_code", "projectCode", "company_code", "companyCode"),
company == null ? "" : company.code(),
companyId,
dept == null ? "" : dept.code(),
"ORG"
);
String projectName = firstNonBlank(
text(employee, "project_name", "projectName", "company_name", "companyName"),
company == null ? "" : company.name(),
deptName,
"组织架构"
);
String positionName = firstNonBlank(text(employee, "position_name", "positionName", "job_title", "jobTitle", "title"), "员工");
String positionLevel = firstNonBlank(text(employee, "position_level", "positionLevel", "job_level", "jobLevel"), level(positionName));
return new OrgRow(
projectCode,
projectName,
deptName,
extPartyId,
firstNonBlank(text(employee, "name", "person_name", "personName", "employee_name", "employeeName"), "-"),
positionName,
positionLevel,
status(text(employee, "status", "employment_status", "employmentStatus")),
LocalDate.now()
);
}
private Map<String, CompanyInfo> companyMap(List<JsonNode> items) {
Map<String, CompanyInfo> map = new TreeMap<>();
for (JsonNode item : items) {
String id = firstNonBlank(text(item, "id", "company_id", "companyId"));
if (!id.isBlank()) {
map.put(id, new CompanyInfo(
firstNonBlank(text(item, "code", "company_code", "companyCode"), id),
firstNonBlank(text(item, "name", "company_name", "companyName"), id)
));
}
}
return map;
}
private Map<String, DepartmentInfo> departmentMap(List<JsonNode> items) {
Map<String, DepartmentInfo> map = new TreeMap<>();
for (JsonNode item : items) {
String id = firstNonBlank(text(item, "id", "department_id", "departmentId", "dept_id", "deptId"));
if (!id.isBlank()) {
map.put(id, new DepartmentInfo(
firstNonBlank(text(item, "code", "department_code", "departmentCode", "dept_code", "deptCode"), id),
firstNonBlank(text(item, "name", "department_name", "departmentName", "dept_name", "deptName"), id),
firstNonBlank(text(item, "company_id", "companyId"), "")
));
}
}
return map;
}
private void requireSnapshotTable() {
try {
Integer count = jdbcTemplate.queryForObject("""
SELECT COUNT(*)
FROM information_schema.TABLES
WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'aihr_org_snapshot'
""", Integer.class);
if (count == null || count == 0) {
throw new IllegalArgumentException("请先执行 aihr_org_snapshot_mysql8.sql 初始化组织人员快照表");
}
} catch (DataAccessException e) {
throw new IllegalStateException("组织人员快照表检查失败: " + e.getMessage(), e);
}
}
private static String status(String value) {
String text = clean(value).toLowerCase();
if (text.isBlank() || text.equals("1") || text.equals("active") || text.equals("enabled")
|| text.equals("normal") || text.equals("在职") || text.equals("正常")) {
return "active";
}
if (text.equals("0") || text.equals("departed") || text.equals("left") || text.equals("inactive")
|| text.equals("disabled") || text.equals("离职")) {
return "departed";
}
return text.length() > 20 ? text.substring(0, 20) : text;
}
private static String level(String positionName) {
String text = clean(positionName);
if (text.contains("项目经理")) {
return "项目经理";
}
if (text.contains("主管") || text.contains("经理") || text.contains("负责人") || text.contains("组长")) {
return "主管";
}
return "一线";
}
private static String text(JsonNode node, String... fields) {
for (String field : fields) {
JsonNode child = node.path(field);
if (child.isValueNode()) {
String value = clean(child.asText());
if (!value.isBlank()) {
return value;
}
}
}
return "";
}
private static void putIfNotBlank(Map<String, String> query, String key, String value) {
String clean = clean(value);
if (!clean.isBlank()) {
query.put(key, clean);
}
}
private static int clamp(Integer value, int fallback, int min, int max) {
int result = value == null ? fallback : value;
return Math.max(min, Math.min(max, result));
}
private static String bearer(String token) {
String clean = clean(token);
return clean.regionMatches(true, 0, "Bearer ", 0, 7) ? clean : "Bearer " + clean;
}
private static String normalizeBaseUrl(String value) {
String clean = clean(value);
while (clean.endsWith("/")) {
clean = clean.substring(0, clean.length() - 1);
}
return clean;
}
private static String queryString(Map<String, String> query) {
List<String> parts = new ArrayList<>();
for (Map.Entry<String, String> entry : query.entrySet()) {
parts.add(url(entry.getKey()) + "=" + url(entry.getValue()));
}
return String.join("&", parts);
}
private static String url(String value) {
return URLEncoder.encode(value, StandardCharsets.UTF_8).replace("+", "%20");
}
private static String sha256(String value) throws Exception {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] bytes = digest.digest((value == null ? "" : value).getBytes(StandardCharsets.UTF_8));
StringBuilder hex = new StringBuilder(bytes.length * 2);
for (byte b : bytes) {
hex.append(String.format("%02x", b));
}
return hex.toString();
}
private static String firstNonBlank(String... values) {
for (String value : values) {
String clean = clean(value);
if (!clean.isBlank()) {
return clean;
}
}
return "";
}
private static String clean(String value) {
return value == null ? "" : value.trim();
}
private static String truncate(String value, int maxChars) {
String text = clean(value);
return text.length() <= maxChars ? text : text.substring(0, maxChars);
}
private record CompanyInfo(String code, String name) {
}
private record DepartmentInfo(String code, String name, String companyId) {
}
private record OrgRow(
String projectCode,
String projectName,
String deptName,
String extPartyId,
String personName,
String positionName,
String positionLevel,
String employmentStatus,
LocalDate snapshotDate
) {
Object[] args() {
return new Object[] {
TENANT_ID, projectCode, projectName, deptName, extPartyId, personName,
positionName, positionLevel, employmentStatus, snapshotDate
};
}
}
}