feat(org): add incremental employee sync via upstream change feed

- POST /api/aihr/org/sync-changes(superadmin/hr_operator):拉取上游
  /sync/changes,按稳定 employee.id 回源 /employees/{id},只替换受影响
  员工,不再依赖每次全量覆盖;dryRun 仍须显式,首次必须传 sinceTime
- extPartyId 优先级翻转为 employee.id 优先:上游 employee_number 仅
  7/3417 匹配既有主体,稳定 employee_id 匹配 3417/3417
- 脱敏手机号保留本地既有有效值;上游显式清空则同步清除;手机号字段
  整体缺失拒绝写入
- 范围字段明确不一致视为迁出,写入时同事务删除本地旧记录(含登录
  身份);范围字段缺失保守跳过不删;回源详情主体与 resource_id 不一致
  在任何范围动作前拒绝
- 任职快照拉取失败在写入路径 fail-closed,不用员工字段兜底降级覆盖
- 全量同步:脱敏手机号在任何写入模式(含 replaceExisting=false)默认
  拒绝,allowPartialReplace=true 时脱敏员工保留本地手机号
- 测试 38 个(org 域),全模块 690 通过
This commit is contained in:
2026-07-25 18:15:20 +08:00
parent 75469b850b
commit 30ef6a8f6b
7 changed files with 877 additions and 36 deletions
@@ -3,6 +3,8 @@ package org.dromara.aihr.controller;
import cn.dev33.satoken.annotation.SaCheckRole;
import cn.dev33.satoken.annotation.SaMode;
import lombok.RequiredArgsConstructor;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncResponse;
import org.dromara.aihr.domain.AihrOrgSyncDto.OrgSnapshotResponse;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncResponse;
@@ -33,6 +35,12 @@ public class AihrOrgSyncController {
return R.ok(governanceService.syncCurrentTenant(request));
}
@SaCheckRole(value = {TenantConstants.SUPER_ADMIN_ROLE_KEY, HR_OPERATOR_ROLE}, mode = SaMode.OR)
@PostMapping("/sync-changes")
public R<ChangeSyncResponse> syncChanges(@RequestBody(required = true) ChangeSyncRequest request) {
return R.ok(governanceService.syncChangesCurrentTenant(request));
}
@SaCheckRole(value = {TenantConstants.SUPER_ADMIN_ROLE_KEY, HR_OPERATOR_ROLE}, mode = SaMode.OR)
@GetMapping("/snapshot")
public R<OrgSnapshotResponse> snapshot(@RequestParam(required = false) String keyword,
@@ -35,6 +35,28 @@ public final class AihrOrgSyncDto {
) {
}
public record ChangeSyncRequest(
Boolean dryRun,
String sinceTime,
String cursor,
Integer pageSize,
Integer maxPages
) {
}
public record ChangeSyncResponse(
boolean dryRun,
String source,
String sinceTime,
String nextCursor,
int changeCount,
int employeeCount,
int syncedEmployeeCount,
int syncedRowCount,
List<String> warnings
) {
}
/** A non-sensitive, server-side view of the global external organization directory. */
public record ExternalDirectorySnapshot(List<ExternalDirectoryNode> nodes, List<String> warnings) {
}
@@ -4,6 +4,8 @@ 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.ChangeSyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncResponse;
import org.dromara.aihr.domain.AihrOrgSyncDto.OrgPersonRow;
import org.dromara.aihr.domain.AihrOrgSyncDto.OrgProjectOption;
import org.dromara.aihr.domain.AihrOrgSyncDto.OrgSnapshotResponse;
@@ -29,6 +31,8 @@ import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.time.Duration;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.time.format.DateTimeParseException;
import java.util.ArrayList;
import java.util.HashSet;
@@ -36,7 +40,9 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.TreeMap;
import java.util.TreeSet;
import java.util.UUID;
import java.util.stream.Stream;
@Service
@RequiredArgsConstructor
@@ -48,6 +54,7 @@ public class AihrOrgSyncService {
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 static final DateTimeFormatter CHANGE_TIME_FORMAT = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
private static final String[] PHONE_FIELDS = {"phone", "phone_number", "phoneNumber", "mobile", "mobile_phone",
"mobilePhone", "mobile_tel", "mobileTel", "cellphone", "cellPhone", "contact_phone",
"contactPhone", "contact_mobile", "contactMobile", "phonenumber", "telephone", "tel", "手机号", "手机", "联系电话"};
@@ -113,6 +120,7 @@ public class AihrOrgSyncService {
Map<String, List<AssignmentInfo>> assignmentGroups = assignmentGroups(assignmentItems);
List<OrgRow> rows = new ArrayList<>();
Set<String> maskedEmployeeIds = new TreeSet<>();
int skipped = 0;
int maskedPhone = 0;
int phoneLinked = 0;
@@ -123,6 +131,7 @@ public class AihrOrgSyncService {
} else {
if (employeeRows.get(0).personPhone().isBlank() && hasUnusablePhoneCandidate(employee)) {
maskedPhone++;
maskedEmployeeIds.add(employeeRows.get(0).extPartyId());
}
if (!employeeRows.get(0).personPhone().isBlank()) {
phoneLinked++;
@@ -158,11 +167,20 @@ public class AihrOrgSyncService {
if (!dryRun && duplicateMemberships > 0) {
throw new IllegalArgumentException("外部员工快照存在重复项目成员关系,已阻止写入组织人员快照;请先修复上游身份数据");
}
if (!dryRun && maskedPhone > 0 && !allowPartialReplace) {
throw new IllegalArgumentException("外部员工快照存在 " + maskedPhone
+ " 条脱敏/不可用手机号,已阻止写入组织人员快照;请先 dry-run,确认后显式传 allowPartialReplace=true");
}
if (!dryRun && replaceExisting && !allowPartialReplace
&& hasUnsafeReplaceData(employeeItems.size(), mappedEmployees, skipped, phoneLinked, maskedPhone, suspectText)) {
throw new IllegalArgumentException("外部员工快照存在不完整或疑似异常数据,已阻止覆盖本地组织人员快照;请先 dry-run,确认后显式传 allowPartialReplace=true");
}
if (!dryRun && !rows.isEmpty()) {
if (maskedPhone > 0) {
// 操作员已显式确认部分替换;脱敏值本身不含信息,不得覆盖本地既有有效手机号。
Map<String, String> localPhones = existingPhones(maskedEmployeeIds);
rows = rows.stream().map(row -> preserveExistingPhone(row, localPhones)).toList();
}
saveRows(rows, replaceExisting);
}
@@ -182,6 +200,124 @@ public class AihrOrgSyncService {
);
}
public ChangeSyncResponse syncChanges(ChangeSyncRequest request, String groupId, String companyId,
String departmentId) {
ChangeSyncRequest req = request == null
? new ChangeSyncRequest(null, null, null, null, null)
: request;
if (req.dryRun() == null) {
throw new IllegalArgumentException("增量组织同步必须明确传 dryRun=true 预检或 dryRun=false 写入");
}
String cursor = limited(req.cursor(), 2000, "cursor");
String sinceTime = changeSinceTime(req.sinceTime(), cursor);
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 = req.dryRun();
if (!dryRun) {
requireSnapshotTable();
}
String baseUrl = normalizeBaseUrl(configuredBaseUrl);
if (baseUrl.isBlank()) {
throw new IllegalArgumentException("请先配置 AIHR_ORG_SYNC_BASE_URL,值为外部开放平台 /api/open/v1 前缀");
}
String token = accessToken(baseUrl);
ChangeBatch changes = fetchEmployeeChanges(baseUrl, token, sinceTime, cursor, pageSize, maxPages);
List<String> warnings = new ArrayList<>();
if (changes.missingChangedFields()) {
warnings.add("上游增量事件未提供 changed_fields,已按 resource_id 回源员工详情");
}
if (changes.employeeIds().isEmpty()) {
return new ChangeSyncResponse(dryRun, baseUrl, sinceTime, changes.nextCursor(),
changes.changeCount(), 0, 0, 0, List.copyOf(warnings));
}
SyncRequest scope = new SyncRequest(true, false, pageSize, maxPages,
groupId, companyId, departmentId, false);
// 写入路径必须以任职快照为准重建项目成员行;拉取失败时 fail-closed,
// 不允许用员工字段兜底生成降级行后覆盖既有项目范围。
List<JsonNode> assignmentItems = dryRun
? fetchOptional(baseUrl, token, scope, "employee_project_assignment", pageSize, maxPages, warnings)
: fetchSnapshot(baseUrl, token, scope, "employee_project_assignment", pageSize, maxPages);
Map<String, List<AssignmentInfo>> assignments = assignmentGroups(assignmentItems);
Map<String, PositionMapping> positionMappings = positionMappings();
Map<String, List<OrgRow>> changedRows = new TreeMap<>();
Set<String> outOfScopeIds = new TreeSet<>();
Set<String> preservablePhoneIds = new TreeSet<>();
int maskedPhones = 0;
int clearedPhones = 0;
for (String resourceId : changes.employeeIds()) {
JsonNode employee = fetchEmployee(baseUrl, token, resourceId);
// 任何范围判断或删除之前,先确认回源详情就是事件主体本身;
// 上游错返其他员工时 fail-closed,绝不据此删除本地记录。
String detailId = firstNonBlank(text(employee, "employee_id", "employeeId", "id", "user_id", "userId"));
if (!resourceId.equals(detailId)) {
throw new IllegalStateException("上游员工详情主体与资源 ID 不一致,已拒绝写入");
}
ScopeMatch scopeMatch = scopeMatch(employee, groupId, companyId, departmentId);
if (scopeMatch == ScopeMatch.OUTSIDE) {
// 员工已明确迁出当前租户绑定范围:本地旧记录(含手机号登录身份)必须在写入时清除,
// 避免离职/迁出人员继续通过旧快照登录。
outOfScopeIds.add(resourceId);
warnings.add("员工 " + resourceId + " 已迁出当前租户绑定范围,写入时将清除本地旧记录");
continue;
}
if (scopeMatch == ScopeMatch.UNKNOWN) {
// 范围字段缺失时无法区分"迁出"与"上游漏字段",保守跳过、绝不删除。
warnings.add("员工 " + resourceId + " 缺少范围字段,无法确认归属,已跳过");
continue;
}
List<OrgRow> rows = orgRows(employee, Map.of(), Map.of(), Map.of(), assignments, positionMappings);
if (rows.isEmpty()) {
warnings.add("员工 " + resourceId + " 缺少稳定 employee.id,已跳过");
continue;
}
if (!resourceId.equals(rows.get(0).extPartyId())) {
throw new IllegalStateException("增量员工主体与资源 ID 不一致,已拒绝写入");
}
if (countDuplicateMemberships(rows) > 0) {
throw new IllegalStateException("增量员工存在重复项目成员关系,已拒绝写入");
}
if (rows.get(0).personPhone().isBlank()) {
if (hasUnusablePhoneCandidate(employee)) {
maskedPhones++;
preservablePhoneIds.add(resourceId);
} else if (!hasPhoneField(employee)) {
throw new IllegalStateException(
"上游员工详情缺少手机号字段,无法区分清空与脱敏,已拒绝写入");
} else {
// 上游显式清空:不得保留旧手机号,避免已失效身份继续登录。
clearedPhones++;
}
}
changedRows.put(resourceId, rows);
}
if (maskedPhones > 0) {
warnings.add("上游有 " + maskedPhones + " 名增量员工仅返回脱敏手机号;写入时保留本地既有有效手机号");
}
if (clearedPhones > 0) {
warnings.add("上游有 " + clearedPhones + " 名增量员工手机号已清空;写入时同步清除本地手机号,对应移动端登录身份将失效");
}
int rowCount = changedRows.values().stream().mapToInt(List::size).sum();
if (!dryRun && (!changedRows.isEmpty() || !outOfScopeIds.isEmpty())) {
Map<String, String> existingPhones = existingPhones(preservablePhoneIds);
Map<String, List<OrgRow>> safeRows = new TreeMap<>();
changedRows.forEach((employeeId, rows) ->
safeRows.put(employeeId, rows.stream().map(row -> preserveExistingPhone(row, existingPhones)).toList()));
saveChangedRows(safeRows, outOfScopeIds);
}
return new ChangeSyncResponse(
dryRun,
baseUrl,
sinceTime,
changes.nextCursor(),
changes.changeCount(),
changedRows.size(),
dryRun ? 0 : changedRows.size(),
dryRun ? 0 : rowCount,
List.copyOf(warnings)
);
}
/**
* Loads the platform-wide directory from the existing global connector.
* This deliberately has no caller supplied scope: binding scope is resolved
@@ -432,6 +568,101 @@ public class AihrOrgSyncService {
throw new IllegalStateException("外部 " + resourceType + " 快照超过最大分页 " + maxPages);
}
private ChangeBatch fetchEmployeeChanges(String baseUrl, String token, String sinceTime, String initialCursor,
int pageSize, int maxPages) {
Set<String> employeeIds = new HashSet<>();
String cursor = clean(initialCursor);
String nextCursor = cursor;
int changeCount = 0;
boolean missingChangedFields = false;
for (int page = 1; page <= maxPages; page++) {
Map<String, String> query = new TreeMap<>();
query.put("resource_type", "employee");
query.put("limit", String.valueOf(pageSize));
if (cursor.isBlank()) {
query.put("since_time", sinceTime);
} else {
query.put("cursor", cursor);
}
JsonNode data = requestJson(baseUrl, "GET", "/sync/changes", query, "", token).path("data");
JsonNode items = data.path("items");
if (!items.isArray()) {
throw new IllegalStateException("外部员工增量响应缺少 data.items");
}
for (JsonNode item : items) {
if (!"employee".equals(text(item, "resource_type", "resourceType"))) {
continue;
}
String resourceId = text(item, "resource_id", "resourceId");
if (!resourceId.isBlank()) {
employeeIds.add(resourceId);
changeCount++;
JsonNode changedFields = item.path("changed_fields");
missingChangedFields |= !changedFields.isArray() || changedFields.isEmpty();
}
}
nextCursor = firstNonBlank(text(data, "cursor", "next_cursor", "nextCursor"), nextCursor);
boolean hasMore = data.path("has_more").asBoolean(data.path("hasMore").asBoolean(false));
if (!hasMore) {
return new ChangeBatch(Set.copyOf(employeeIds), changeCount, nextCursor, missingChangedFields);
}
if (nextCursor.isBlank() || nextCursor.equals(cursor)) {
throw new IllegalStateException("外部员工增量分页游标缺失或未推进");
}
cursor = nextCursor;
}
throw new IllegalStateException("外部员工增量超过最大分页 " + maxPages);
}
private JsonNode fetchEmployee(String baseUrl, String token, String resourceId) {
JsonNode data = requestJson(baseUrl, "GET", "/employees/" + url(resourceId), Map.of(), "", token).path("data");
JsonNode item = data.path("item");
if (item.isMissingNode() && data.isObject()) {
item = data;
}
if (!item.isObject() || text(item, "id", "employee_id", "employeeId").isBlank()) {
throw new IllegalStateException("外部员工详情响应缺少 data.item");
}
return item;
}
/**
* 范围归属三态:约束字段存在且不一致才是"明确迁出";上游漏字段返回 UNKNOWN,
* 由调用方保守跳过,避免把字段缺失误判成迁出而删除有效记录。
*/
private static ScopeMatch scopeMatch(JsonNode employee, String groupId, String companyId, String departmentId) {
ScopeMatch result = ScopeMatch.INSIDE;
result = result.merge(scopeFieldMatch(clean(groupId), text(employee, "group_id", "groupId")));
result = result.merge(scopeFieldMatch(clean(companyId), text(employee, "company_id", "companyId")));
result = result.merge(scopeFieldMatch(clean(departmentId),
text(employee, "department_id", "departmentId", "dept_id", "deptId")));
return result;
}
private static ScopeMatch scopeFieldMatch(String constraint, String actual) {
if (constraint.isBlank()) {
return ScopeMatch.INSIDE;
}
if (actual.isBlank()) {
return ScopeMatch.UNKNOWN;
}
return constraint.equals(actual) ? ScopeMatch.INSIDE : ScopeMatch.OUTSIDE;
}
private enum ScopeMatch {
INSIDE, OUTSIDE, UNKNOWN;
private ScopeMatch merge(ScopeMatch other) {
if (this == OUTSIDE || other == OUTSIDE) {
return OUTSIDE;
}
if (this == UNKNOWN || other == UNKNOWN) {
return UNKNOWN;
}
return INSIDE;
}
}
private String accessToken(String baseUrl) {
String token = clean(configuredAccessToken);
if (!token.isBlank()) {
@@ -516,32 +747,88 @@ public class AihrOrgSyncService {
if (replaceExisting) {
jdbcTemplate.update("delete from aihr_org_snapshot where tenant_id = ?", tenantId());
}
jdbcTemplate.batchUpdate("""
insert into aihr_org_snapshot
(tenant_id, project_code, project_name, dept_name, ext_party_id, person_phone, person_name,
position_name, position_level, employment_status, hire_date, snapshot_date, create_time)
values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, now())
on duplicate key update
project_name = values(project_name),
dept_name = values(dept_name),
person_phone = values(person_phone),
person_name = values(person_name),
position_name = values(position_name),
position_level = values(position_level),
employment_status = values(employment_status),
hire_date = values(hire_date),
snapshot_date = values(snapshot_date),
create_time = now()
""", rows.stream().map(row -> row.args(tenantId())).toList());
upsertRows(rows);
});
}
private void saveChangedRows(Map<String, List<OrgRow>> rowsByEmployee, Set<String> deleteOnlyIds) {
transactionTemplate.executeWithoutResult(status -> {
String tenantId = tenantId();
Stream.concat(rowsByEmployee.keySet().stream(), deleteOnlyIds.stream()).forEach(employeeId ->
jdbcTemplate.update(
"delete from aihr_org_snapshot where tenant_id = ? and ext_party_id = ?",
tenantId, employeeId));
upsertRows(rowsByEmployee.values().stream().flatMap(List::stream).toList());
});
}
private void upsertRows(List<OrgRow> rows) {
if (rows.isEmpty()) {
return;
}
jdbcTemplate.batchUpdate("""
insert into aihr_org_snapshot
(tenant_id, project_code, project_name, dept_name, ext_party_id, person_phone, person_name,
position_name, position_level, employment_status, hire_date, snapshot_date, create_time)
values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, now())
on duplicate key update
project_name = values(project_name),
dept_name = values(dept_name),
person_phone = values(person_phone),
person_name = values(person_name),
position_name = values(position_name),
position_level = values(position_level),
employment_status = values(employment_status),
hire_date = values(hire_date),
snapshot_date = values(snapshot_date),
create_time = now()
""", rows.stream().map(row -> row.args(tenantId())).toList());
}
private Map<String, String> existingPhones(Set<String> employeeIds) {
if (employeeIds.isEmpty()) {
return Map.of();
}
String placeholders = String.join(",", java.util.Collections.nCopies(employeeIds.size(), "?"));
List<Object> args = new ArrayList<>();
args.add(tenantId());
args.addAll(employeeIds);
List<ExistingPhone> rows = jdbcTemplate.query("""
select ext_party_id, person_phone
from aihr_org_snapshot
where tenant_id = ? and ext_party_id in (""" + placeholders + ")",
(rs, rowNum) -> new ExistingPhone(rs.getString("ext_party_id"), rs.getString("person_phone")),
args.toArray());
Map<String, String> phones = new TreeMap<>();
for (ExistingPhone row : rows) {
String phone = normalizeMobilePhone(row.phone());
if (!phone.isBlank()) {
phones.merge(row.employeeId(), phone, (left, right) -> left.equals(right) ? left : "");
}
}
return phones;
}
private static OrgRow preserveExistingPhone(OrgRow row, Map<String, String> existingPhones) {
if (!row.personPhone().isBlank()) {
return row;
}
String phone = clean(existingPhones.get(row.extPartyId()));
if (phone.isBlank()) {
return row;
}
return new OrgRow(
row.projectCode(), row.projectName(), row.deptName(), row.extPartyId(), phone, row.personName(),
row.positionName(), row.positionLevel(), row.positionMapped(), row.employmentStatus(), row.hireDate(),
row.snapshotDate(), row.suspectText());
}
private OrgRow orgRow(JsonNode employee, Map<String, CompanyInfo> companies, Map<String, DepartmentInfo> departments,
Map<String, ProjectInfo> projects, Map<String, AssignmentInfo> assignments,
Map<String, PositionMapping> positionMappings) {
String employeeNumber = firstNonBlank(text(employee, "employee_number", "employeeNumber", "employeeNo"));
String employeeId = firstNonBlank(text(employee, "employee_id", "employeeId", "id", "user_id", "userId"));
String extPartyId = firstNonBlank(employeeNumber, employeeId);
String extPartyId = firstNonBlank(employeeId, employeeNumber);
if (extPartyId.isBlank()) {
return null;
}
@@ -1038,6 +1325,16 @@ public class AihrOrgSyncService {
return false;
}
/** 手机号字段是否在上游响应中显式出现(空串也算出现,表示上游明确清空)。 */
static boolean hasPhoneField(JsonNode employee) {
for (String field : PHONE_FIELDS) {
if (employee.has(field) && !employee.path(field).isNull()) {
return true;
}
}
return false;
}
private static boolean hasQuestionMark(String value) {
return clean(value).contains("?");
}
@@ -1147,6 +1444,29 @@ public class AihrOrgSyncService {
return value == null ? "" : value.trim();
}
private static String limited(String value, int maxLength, String field) {
String text = clean(value);
if (text.length() > maxLength) {
throw new IllegalArgumentException(field + " 长度不能超过 " + maxLength);
}
return text;
}
private static String changeSinceTime(String value, String cursor) {
String text = limited(value, 19, "sinceTime");
if (!clean(cursor).isBlank()) {
return text;
}
if (text.isBlank()) {
throw new IllegalArgumentException("首次增量同步必须提供 sinceTime,格式为 yyyy-MM-dd HH:mm:ss");
}
try {
return LocalDateTime.parse(text, CHANGE_TIME_FORMAT).format(CHANGE_TIME_FORMAT);
} catch (DateTimeParseException e) {
throw new IllegalArgumentException("sinceTime 格式必须为 yyyy-MM-dd HH:mm:ss");
}
}
private static LocalDate parseDate(String value) {
String text = clean(value);
if (text.isBlank()) {
@@ -1194,6 +1514,13 @@ public class AihrOrgSyncService {
private record PositionMapping(String name, String level) {
}
private record ChangeBatch(Set<String> employeeIds, int changeCount, String nextCursor,
boolean missingChangedFields) {
}
private record ExistingPhone(String employeeId, String phone) {
}
private record OrgRow(
String projectCode,
String projectName,
@@ -1,6 +1,8 @@
package org.dromara.aihr.service;
import lombok.RequiredArgsConstructor;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncResponse;
import org.dromara.aihr.domain.AihrOrgSyncDto.ExternalDirectoryNode;
import org.dromara.aihr.domain.AihrOrgSyncDto.ExternalDirectorySnapshot;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncRequest;
@@ -181,6 +183,26 @@ public class AihrTenantOrgGovernanceService {
return response;
}
public ChangeSyncResponse syncChangesCurrentTenant(ChangeSyncRequest requested) {
String tenantId = currentTenantId();
BindingRow binding = activeBindingForTenant(tenantId);
if (binding == null) {
throw new ServiceException("TENANT_ORG_BINDING_REQUIRED", HttpStatus.CONFLICT);
}
if (AGGREGATE_ONLY.equals(binding.mode())) {
throw new ServiceException("AGGREGATE_ONLY_CANNOT_SYNC", HttpStatus.CONFLICT);
}
SyncScope scope = syncScope(binding);
ChangeSyncResponse response = orgSyncService.syncChanges(
requested, scope.groupId(), scope.companyId(), scope.departmentId());
if (!response.dryRun()) {
jdbcTemplate.update(
"update aihr_tenant_org_binding set last_sync_at = now(), update_time = now() where id = ?",
binding.id());
}
return response;
}
static boolean isAncestorOrSame(String ancestorPath, String descendantPath) {
return hasText(ancestorPath) && hasText(descendantPath)
&& (ancestorPath.equals(descendantPath) || descendantPath.startsWith(ancestorPath + "/"));
@@ -221,24 +243,25 @@ public class AihrTenantOrgGovernanceService {
}
private SyncRequest scopedRequest(SyncRequest requested, BindingRow binding) {
if (binding == null || !hasText(binding.rootExternalId()) || !hasText(binding.rootNodeType())) {
throw new ServiceException("TENANT_ORG_BINDING_REQUIRED", HttpStatus.CONFLICT);
}
String groupId = null;
String companyId = null;
String departmentId = null;
String sourceExternalId = AihrOrgSyncService.rawDirectoryId(binding.rootExternalId(), binding.rootNodeType());
switch (binding.rootNodeType()) {
case "GROUP" -> groupId = sourceExternalId;
case "COMPANY" -> companyId = sourceExternalId;
case "DEPARTMENT" -> departmentId = sourceExternalId;
default -> throw new ServiceException("ORG_SCOPE_TYPE_NOT_SYNCABLE", HttpStatus.CONFLICT);
}
SyncScope scope = syncScope(binding);
SyncRequest source = requested == null
? new SyncRequest(null, null, null, null, null, null, null, null)
: requested;
return new SyncRequest(source.dryRun(), source.replaceExisting(), source.pageSize(), source.maxPages(),
groupId, companyId, departmentId, source.allowPartialReplace());
scope.groupId(), scope.companyId(), scope.departmentId(), source.allowPartialReplace());
}
private SyncScope syncScope(BindingRow binding) {
if (binding == null || !hasText(binding.rootExternalId()) || !hasText(binding.rootNodeType())) {
throw new ServiceException("TENANT_ORG_BINDING_REQUIRED", HttpStatus.CONFLICT);
}
String sourceExternalId = AihrOrgSyncService.rawDirectoryId(binding.rootExternalId(), binding.rootNodeType());
return switch (binding.rootNodeType()) {
case "GROUP" -> new SyncScope(sourceExternalId, null, null);
case "COMPANY" -> new SyncScope(null, sourceExternalId, null);
case "DEPARTMENT" -> new SyncScope(null, null, sourceExternalId);
default -> throw new ServiceException("ORG_SCOPE_TYPE_NOT_SYNCABLE", HttpStatus.CONFLICT);
};
}
private List<ScopeConflict> conflicts(BindingInput input, DirectoryRow root, List<BindingRow> active,
@@ -512,6 +535,9 @@ public class AihrTenantOrgGovernanceService {
private record BindingInput(String tenantId, String rootExternalId, String mode, String status) {
}
private record SyncScope(String groupId, String companyId, String departmentId) {
}
private record TenantRow(String tenantId, String companyName, String status) {
}
}
@@ -2,25 +2,32 @@ package org.dromara.aihr.service;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.sun.net.httpserver.HttpServer;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncResponse;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncRequest;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Tag;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.RowMapper;
import org.springframework.test.util.ReflectionTestUtils;
import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.TransactionTemplate;
import java.lang.reflect.Proxy;
import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.function.Consumer;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verifyNoInteractions;
@@ -148,7 +155,29 @@ public class AihrOrgSyncServiceTest {
assertEquals(2, rows.size());
assertEquals(List.of("FW001", "YSF001"), rows.stream()
.map(row -> String.valueOf(ReflectionTestUtils.getField(row, "projectCode"))).toList());
assertTrue(rows.stream().allMatch(row -> "YC001".equals(ReflectionTestUtils.getField(row, "extPartyId"))));
// 主体身份固定为稳定 employee.id,不再被易变的 employee_number 抢占
assertTrue(rows.stream().allMatch(row -> "EMP-1".equals(ReflectionTestUtils.getField(row, "extPartyId"))));
}
@Test
public void stableEmployeeIdTakesPriorityOverEmployeeNumberForIdentity() throws Exception {
ObjectMapper mapper = new ObjectMapper();
AihrOrgSyncService service = new AihrOrgSyncService(
new ObjectMapper(), mock(JdbcTemplate.class), mock(TransactionTemplate.class));
var employee = mapper.readTree("""
{"id":"1000","employee_number":"YC001","name":"测试员工","phone":"13900001111"}
""");
Object row = ReflectionTestUtils.invokeMethod(service, "orgRow", employee,
Map.of(), Map.of(), Map.of(), Map.of(), Map.of());
assertEquals("1000", ReflectionTestUtils.getField(row, "extPartyId"));
Object noIdRow = ReflectionTestUtils.invokeMethod(service, "orgRow",
mapper.readTree("{\"employee_number\":\"YC001\",\"name\":\"测试员工\"}"),
Map.of(), Map.of(), Map.of(), Map.of(), Map.of());
assertEquals("YC001", ReflectionTestUtils.getField(noIdRow, "extPartyId"));
}
@Test
@@ -328,6 +357,416 @@ public class AihrOrgSyncServiceTest {
}
}
@Test
public void changeSyncRequiresExplicitDryRunAndSinceTime() {
AihrOrgSyncService service = new AihrOrgSyncService(
new ObjectMapper(), mock(JdbcTemplate.class), mock(TransactionTemplate.class));
IllegalArgumentException missingMode = assertThrows(IllegalArgumentException.class,
() -> service.syncChanges(new ChangeSyncRequest(null, "2026-07-25 16:00:00", null, null, null),
null, null, null));
assertTrue(missingMode.getMessage().contains("dryRun"));
IllegalArgumentException missingSince = assertThrows(IllegalArgumentException.class,
() -> service.syncChanges(new ChangeSyncRequest(true, null, null, null, null), null, null, null));
assertTrue(missingSince.getMessage().contains("sinceTime"));
IllegalArgumentException badSince = assertThrows(IllegalArgumentException.class,
() -> service.syncChanges(new ChangeSyncRequest(true, "2026/07/25", null, null, null), null, null, null));
assertTrue(badSince.getMessage().contains("yyyy-MM-dd HH:mm:ss"));
}
@Test
public void changeSyncDryRunReadsChangesWithoutTouchingStorage() throws Exception {
HttpServer server = changeSyncServer(200);
try {
JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
TransactionTemplate transactionTemplate = mock(TransactionTemplate.class);
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, transactionTemplate);
ChangeSyncResponse response = service.syncChanges(
new ChangeSyncRequest(true, "2026-07-25 16:00:00", null, 10, 5), "1", null, null);
assertTrue(response.dryRun());
assertEquals(1, response.changeCount());
assertEquals(1, response.employeeCount());
assertEquals(0, response.syncedEmployeeCount());
assertEquals(0, response.syncedRowCount());
assertEquals("cursor-1", response.nextCursor());
assertTrue(response.warnings().stream().anyMatch(warning -> warning.contains("changed_fields")));
verifyNoInteractions(jdbcTemplate, transactionTemplate);
} finally {
server.stop(0);
}
}
@Test
public void changeSyncWritePreservesExistingPhoneWhenUpstreamMasks() throws Exception {
HttpServer server = changeSyncServer(200);
try {
ChangeSyncJdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
TransactionTemplate transactionTemplate = mock(TransactionTemplate.class);
doAnswer(invocation -> {
invocation.getArgument(0, Consumer.class).accept(mock(TransactionStatus.class));
return null;
}).when(transactionTemplate).executeWithoutResult(any());
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, transactionTemplate);
ChangeSyncResponse response = service.syncChanges(
new ChangeSyncRequest(false, "2026-07-25 16:00:00", null, 10, 5), "1", null, null);
assertFalse(response.dryRun());
assertEquals(1, response.syncedEmployeeCount());
assertEquals(1, response.syncedRowCount());
assertTrue(response.warnings().stream().anyMatch(warning -> warning.contains("脱敏手机号")));
assertEquals(List.of("1000"), jdbcTemplate.deletedEmployeeIds);
assertEquals(1, jdbcTemplate.upsertedRows.size());
Object[] row = jdbcTemplate.upsertedRows.get(0);
assertEquals("1000", row[4]);
assertEquals("13900001111", row[5]);
assertEquals("生活顾问", row[7]);
} finally {
server.stop(0);
}
}
@Test
public void changeSyncSkipsEmployeesOutsideTenantScope() throws Exception {
HttpServer server = changeSyncServer(200);
try {
AihrOrgSyncService service = changeSyncService(server,
mock(JdbcTemplate.class), mock(TransactionTemplate.class));
ChangeSyncResponse response = service.syncChanges(
new ChangeSyncRequest(true, "2026-07-25 16:00:00", null, 10, 5), "2", null, null);
assertEquals(1, response.changeCount());
assertEquals(0, response.employeeCount());
assertTrue(response.warnings().stream().anyMatch(warning -> warning.contains("已迁出当前租户绑定范围")));
} finally {
server.stop(0);
}
}
@Test
public void changeSyncWriteDeletesRowsForEmployeesMovedOutOfScope() throws Exception {
HttpServer server = changeSyncServer(200);
try {
ChangeSyncJdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
TransactionTemplate transactionTemplate = executingTransaction();
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, transactionTemplate);
ChangeSyncResponse response = service.syncChanges(
new ChangeSyncRequest(false, "2026-07-25 16:00:00", null, 10, 5), "2", null, null);
assertEquals(List.of("1000"), jdbcTemplate.deletedEmployeeIds);
assertTrue(jdbcTemplate.upsertedRows.isEmpty());
assertTrue(response.warnings().stream().anyMatch(warning -> warning.contains("已迁出当前租户绑定范围")));
} finally {
server.stop(0);
}
}
@Test
public void changeSyncSkipsWithoutDeletingWhenScopeFieldsMissing() throws Exception {
HttpServer server = changeSyncServer(200,
"{\"id\":\"1000\",\"name\":\"测试员工\",\"position_name\":\"生活顾问\",\"phone\":\"13900001111\"}");
try {
ChangeSyncJdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, executingTransaction());
ChangeSyncResponse response = service.syncChanges(
new ChangeSyncRequest(false, "2026-07-25 16:00:00", null, 10, 5), "1", null, null);
assertEquals(0, response.employeeCount());
assertTrue(response.warnings().stream().anyMatch(warning -> warning.contains("无法确认归属")));
assertTrue(jdbcTemplate.deletedEmployeeIds.isEmpty());
assertTrue(jdbcTemplate.upsertedRows.isEmpty());
} finally {
server.stop(0);
}
}
@Test
public void changeSyncWriteClearsLocalPhoneWhenUpstreamExplicitlyClears() throws Exception {
HttpServer server = changeSyncServer(200,
"{\"id\":\"1000\",\"group_id\":\"1\",\"name\":\"测试员工\",\"position_name\":\"生活顾问\",\"phone\":\"\"}");
try {
ChangeSyncJdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, executingTransaction());
ChangeSyncResponse response = service.syncChanges(
new ChangeSyncRequest(false, "2026-07-25 16:00:00", null, 10, 5), "1", null, null);
assertEquals(1, response.syncedEmployeeCount());
assertTrue(response.warnings().stream().anyMatch(warning -> warning.contains("已清空")));
assertEquals(1, jdbcTemplate.upsertedRows.size());
assertEquals("", jdbcTemplate.upsertedRows.get(0)[5]);
} finally {
server.stop(0);
}
}
@Test
public void changeSyncRejectsDetailFromDifferentEmployeeBeforeAnyScopeAction() throws Exception {
HttpServer server = changeSyncServer(200,
"{\"id\":\"2000\",\"group_id\":\"9\",\"name\":\"错误员工\",\"phone\":\"13900002222\"}");
try {
ChangeSyncJdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, executingTransaction());
IllegalStateException error = assertThrows(IllegalStateException.class, () -> service.syncChanges(
new ChangeSyncRequest(false, "2026-07-25 16:00:00", null, 10, 5), "1", null, null));
assertTrue(error.getMessage().contains("主体与资源 ID 不一致"));
assertTrue(jdbcTemplate.deletedEmployeeIds.isEmpty());
assertTrue(jdbcTemplate.upsertedRows.isEmpty());
} finally {
server.stop(0);
}
}
@Test
public void changeSyncRejectsWhenPhoneFieldMissingEntirely() throws Exception {
HttpServer server = changeSyncServer(200,
"{\"id\":\"1000\",\"group_id\":\"1\",\"name\":\"测试员工\",\"position_name\":\"生活顾问\"}");
try {
AihrOrgSyncService service = changeSyncService(server,
new ChangeSyncJdbcTemplate(), executingTransaction());
IllegalStateException error = assertThrows(IllegalStateException.class, () -> service.syncChanges(
new ChangeSyncRequest(false, "2026-07-25 16:00:00", null, 10, 5), "1", null, null));
assertTrue(error.getMessage().contains("缺少手机号字段"));
} finally {
server.stop(0);
}
}
@Test
public void fullSyncRejectsMaskedPhonesEvenWithoutReplace() throws Exception {
HttpServer server = maskedEmployeeSnapshotServer();
try {
AihrOrgSyncService service = changeSyncService(server,
new ChangeSyncJdbcTemplate(), executingTransaction());
IllegalArgumentException error = assertThrows(IllegalArgumentException.class, () -> service.sync(
new SyncRequest(false, false, 10, 5, null, null, null, false)));
assertTrue(error.getMessage().contains("脱敏"));
} finally {
server.stop(0);
}
}
@Test
public void fullSyncPartialReplacePreservesLocalPhoneBehindMaskedUpstream() throws Exception {
HttpServer server = maskedEmployeeSnapshotServer();
try {
ChangeSyncJdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, executingTransaction());
var response = service.sync(new SyncRequest(false, false, 10, 5, null, null, null, true));
assertEquals(1, response.syncedCount());
assertEquals(1, jdbcTemplate.upsertedRows.size());
assertEquals("13900001111", jdbcTemplate.upsertedRows.get(0)[5]);
} finally {
server.stop(0);
}
}
@Test
public void fullSyncReplaceWithPartialReplaceAlsoPreservesMaskedPhones() throws Exception {
HttpServer server = maskedEmployeeSnapshotServer();
try {
ChangeSyncJdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, executingTransaction());
var response = service.sync(new SyncRequest(false, true, 10, 5, null, null, null, true));
assertEquals(1, response.syncedCount());
assertEquals(1, jdbcTemplate.upsertedRows.size());
assertEquals("13900001111", jdbcTemplate.upsertedRows.get(0)[5]);
} finally {
server.stop(0);
}
}
private static HttpServer maskedEmployeeSnapshotServer() throws Exception {
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
server.createContext("/api/open/v1/sync/snapshot", exchange -> {
boolean employee = exchange.getRequestURI().getRawQuery().contains("resource_type=employee");
String items = employee
? "[{\"id\":\"1000\",\"name\":\"测试员工\",\"phone\":\"139****1111\"}]"
: "[]";
writeJson(exchange, "{\"data\":{\"items\":" + items + ",\"total\":" + (employee ? 1 : 0)
+ ",\"has_more\":false}}");
});
server.start();
return server;
}
private static TransactionTemplate executingTransaction() {
TransactionTemplate transactionTemplate = mock(TransactionTemplate.class);
doAnswer(invocation -> {
invocation.getArgument(0, Consumer.class).accept(mock(TransactionStatus.class));
return null;
}).when(transactionTemplate).executeWithoutResult(any());
return transactionTemplate;
}
@Test
public void changeSyncRejectsDuplicateProjectMemberships() throws Exception {
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
server.createContext("/api/open/v1/sync/changes", exchange -> writeJson(exchange,
"{\"data\":{\"items\":[{\"event_type\":\"employee.updated\",\"resource_type\":\"employee\","
+ "\"resource_id\":\"1000\",\"changed_fields\":[]}],\"cursor\":\"cursor-1\",\"has_more\":false}}"));
server.createContext("/api/open/v1/employees/1000", exchange -> writeJson(exchange,
"{\"data\":{\"item\":{\"id\":\"1000\",\"group_id\":\"1\",\"name\":\"测试员工\",\"phone\":\"13900001111\"}}}"));
server.createContext("/api/open/v1/sync/snapshot", exchange -> writeJson(exchange,
"{\"data\":{\"items\":["
+ "{\"employee_id\":\"1000\",\"project_code\":\"FW001\",\"project_name\":\"翡翠湾\","
+ "\"project_position_name\":\"生活顾问\",\"status\":\"active\"},"
+ "{\"employee_id\":\"1000\",\"project_code\":\"FW001\",\"project_name\":\"翡翠湾\","
+ "\"project_position_name\":\"管家\",\"status\":\"active\"}"
+ "],\"has_more\":false}}"));
server.start();
try {
AihrOrgSyncService service = changeSyncService(server,
mock(JdbcTemplate.class), mock(TransactionTemplate.class));
IllegalStateException error = assertThrows(IllegalStateException.class, () -> service.syncChanges(
new ChangeSyncRequest(true, "2026-07-25 16:00:00", null, 10, 5), "1", null, null));
assertTrue(error.getMessage().contains("重复项目成员关系"));
} finally {
server.stop(0);
}
}
@Test
public void changeSyncWriteFailsClosedWhenAssignmentSnapshotFails() throws Exception {
HttpServer server = changeSyncServer(500);
try {
JdbcTemplate jdbcTemplate = new ChangeSyncJdbcTemplate();
TransactionTemplate transactionTemplate = mock(TransactionTemplate.class);
AihrOrgSyncService service = changeSyncService(server, jdbcTemplate, transactionTemplate);
assertThrows(IllegalStateException.class, () -> service.syncChanges(
new ChangeSyncRequest(false, "2026-07-25 16:00:00", null, 10, 5), "1", null, null));
ChangeSyncResponse dryRun = service.syncChanges(
new ChangeSyncRequest(true, "2026-07-25 16:00:00", null, 10, 5), "1", null, null);
assertTrue(dryRun.warnings().stream().anyMatch(warning -> warning.contains("快照拉取失败")));
} finally {
server.stop(0);
}
}
private static HttpServer changeSyncServer(int assignmentStatus) throws Exception {
return changeSyncServer(assignmentStatus,
"{\"id\":\"1000\",\"group_id\":\"1\",\"name\":\"测试员工\","
+ "\"position_name\":\"生活顾问\",\"phone\":\"139****1111\"}");
}
private static HttpServer changeSyncServer(int assignmentStatus, String employeeJson) throws Exception {
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
server.createContext("/api/open/v1/sync/changes", exchange -> {
String query = exchange.getRequestURI().getRawQuery();
if (query == null || !query.contains("limit=") || query.contains("page_size=")) {
exchange.sendResponseHeaders(400, -1);
exchange.close();
return;
}
writeJson(exchange,
"{\"data\":{\"items\":[{\"event_type\":\"employee.updated\",\"resource_type\":\"employee\","
+ "\"resource_id\":\"1000\",\"changed_fields\":[]}],\"cursor\":\"cursor-1\",\"has_more\":false}}");
});
server.createContext("/api/open/v1/employees/1000", exchange -> writeJson(exchange,
"{\"data\":{\"item\":" + employeeJson + "}}"));
server.createContext("/api/open/v1/sync/snapshot", exchange -> {
if (assignmentStatus != 200) {
exchange.sendResponseHeaders(assignmentStatus, -1);
exchange.close();
return;
}
writeJson(exchange,
"{\"data\":{\"items\":[{\"employee_id\":\"1000\",\"project_code\":\"FW001\",\"project_name\":\"翡翠湾\","
+ "\"project_position_name\":\"生活顾问\",\"status\":\"active\"}],\"has_more\":false}}");
});
server.start();
return server;
}
private static void writeJson(com.sun.net.httpserver.HttpExchange exchange, String body) throws java.io.IOException {
byte[] bytes = body.getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "application/json; charset=utf-8");
exchange.sendResponseHeaders(200, bytes.length);
exchange.getResponseBody().write(bytes);
exchange.close();
}
private static AihrOrgSyncService changeSyncService(HttpServer server, JdbcTemplate jdbcTemplate,
TransactionTemplate transactionTemplate) {
AihrOrgSyncService service = new AihrOrgSyncService(new ObjectMapper(), jdbcTemplate, transactionTemplate);
ReflectionTestUtils.setField(service, "configuredBaseUrl",
"http://127.0.0.1:" + server.getAddress().getPort() + "/api/open/v1");
ReflectionTestUtils.setField(service, "configuredAccessToken", "change-sync-test-token");
ReflectionTestUtils.setField(service, "configuredClientId", "");
ReflectionTestUtils.setField(service, "configuredClientSecret", "");
ReflectionTestUtils.setField(service, "configuredSigningSecret", "");
return service;
}
private static final class ChangeSyncJdbcTemplate extends JdbcTemplate {
private final List<String> deletedEmployeeIds = new ArrayList<>();
private final List<Object[]> upsertedRows = new ArrayList<>();
@Override
public <T> T queryForObject(String sql, Class<T> requiredType) {
return requiredType.cast(1);
}
@Override
public <T> T queryForObject(String sql, Class<T> requiredType, Object... args) {
return requiredType.cast(1);
}
@Override
public <T> List<T> query(String sql, RowMapper<T> rowMapper, Object... args) {
List<T> rows = new ArrayList<>();
if (sql.contains("person_phone")) {
try {
rows.add(rowMapper.mapRow((java.sql.ResultSet) Proxy.newProxyInstance(
ChangeSyncJdbcTemplate.class.getClassLoader(),
new Class<?>[] {java.sql.ResultSet.class},
(proxy, method, methodArgs) -> switch (method.getName()) {
case "getString" -> "ext_party_id".equals(methodArgs[0]) ? "1000" : "13900001111";
default -> throw new UnsupportedOperationException(method.getName());
}), 0));
} catch (java.sql.SQLException e) {
throw new IllegalStateException(e);
}
}
return rows;
}
@Override
public int update(String sql, Object... args) {
if (sql.contains("delete from aihr_org_snapshot") && args.length > 1) {
deletedEmployeeIds.add(String.valueOf(args[1]));
}
return 1;
}
@Override
public int[] batchUpdate(String sql, List<Object[]> batchArgs) {
upsertedRows.addAll(batchArgs);
return new int[batchArgs.size()];
}
}
private static final class SnapshotJdbcTemplate extends JdbcTemplate {
private final List<String> queries = new ArrayList<>();
@@ -1,5 +1,7 @@
package org.dromara.aihr.service;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.ChangeSyncResponse;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncRequest;
import org.dromara.aihr.domain.AihrOrgSyncDto.SyncResponse;
import org.dromara.common.core.exception.ServiceException;
@@ -19,6 +21,7 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
@@ -79,6 +82,20 @@ class AihrTenantOrgGovernanceServiceTest {
assertEquals(null, request.getValue().departmentId());
}
@Test
void changeSyncUsesTheBindingRoot() {
ChangeSyncRequest request = new ChangeSyncRequest(true, "2026-07-25 16:00:00", null, 20, 1);
AihrOrgSyncService orgSyncService = mock(AihrOrgSyncService.class);
when(orgSyncService.syncChanges(any(), any(), any(), any())).thenReturn(
new ChangeSyncResponse(true, "source", request.sinceTime(), "cursor", 1, 1, 0, 0, List.of()));
AihrTenantOrgGovernanceService service = new AihrTenantOrgGovernanceService(
new BindingJdbcTemplate("EXCLUSIVE"), mock(TransactionTemplate.class), orgSyncService);
service.syncChangesCurrentTenant(request);
verify(orgSyncService).syncChanges(eq(request), eq(null), eq("C-1"), eq(null));
}
private static final class EmptyJdbcTemplate extends JdbcTemplate {
@Override
public <T> List<T> query(String sql, RowMapper<T> rowMapper, Object... args) {
+5 -3
View File
@@ -54,7 +54,9 @@ portless
Qdrant 默认本地无需配置;远端或自定义 collection 可用 `AIHR_QDRANT_URL`、`AIHR_QDRANT_COLLECTION`、`AIHR_QDRANT_API_KEY` 覆盖。服务端资料导入根目录可用 `AIHR_IMPORT_ROOT` 或 `-Daihr.import.root` 覆盖。组织人员同步可用 `AIHR_ORG_SYNC_BASE_URL` 指向外部开放平台 `/api/open/v1` 前缀,并配置 `AIHR_ORG_SYNC_ACCESS_TOKEN` 或 `AIHR_ORG_SYNC_CLIENT_ID`/`AIHR_ORG_SYNC_CLIENT_SECRET`;业务请求会用 client secret 生成 HMAC-SHA256 hex 签名。移动端手机号登录的短信模板 ID、阿里云 AccessKey、Secret 和短信签名都通过环境变量注入;本地放根目录 `.env.local`,`scripts/dev-backend.sh` 会自动加载。若开放平台凭证放在 `backend/.env`,启动脚本也会加载该文件,并把 `client_id`/`client_secret` 映射为组织同步实际读取的 `AIHR_ORG_SYNC_CLIENT_ID`/`AIHR_ORG_SYNC_CLIENT_SECRET`。
组织同步写入前先运行只读预检:`node scripts/verify-demo-questions.mjs --org-dry-run`。dry-run 不检查或变更本地快照表结构,也不写数据库;输出只包含人数、手机号覆盖、脱敏数、疑似乱码数和警告,不输出员工姓名。2026-07-15 当前开放平台数据已满足写入条件:3417 名员工中 3392 人可手机号映射,疑似乱码为 0;仍有 25 人手机号不可用,确认后可用 `allowPartialReplace=true` 覆盖写入。
组织同步写入前先运行只读预检:`node scripts/verify-demo-questions.mjs --org-dry-run`。dry-run 不检查或变更本地快照表结构,也不写数据库;输出只包含人数、手机号覆盖、脱敏数、疑似乱码数和警告,不输出员工姓名。2026-07-25 只读预检返回 3424 名员工、3398 个脱敏手机号和 1 条重复项目成员关系;当前上游优先字段 `employee_number` 仅匹配生产既有主体 `7/3417`,稳定 `employee_id` 匹配 `3417/3417`。全量快照仍存在脱敏手机号和重复关系,不得执行非 dry-run 请求,也不得用 `allowPartialReplace=true` 绕过。
员工岗位等字段的日常变化不要执行不安全的全量覆盖;部署支持增量同步的后端后,使用受 `superadmin/hr_operator` 保护的 `POST /api/aihr/org/sync-changes`。首次传 `{"dryRun":true,"sinceTime":"YYYY-MM-DD HH:mm:ss"}` 预检,再以相同 `sinceTime` 和 `dryRun:false` 写入;接口主动拉 `/sync/changes` 并按稳定 `employee.id` 回源详情,只替换命中的员工,脱敏手机号保留本地有效值,任职快照拉取失败则拒绝写库。
正式试点预检必须指定当前批次租户和时间窗,例如:`AIHR_PILOT_TENANT_ID=000000 AIHR_PILOT_START_DATE=2026-07-07 AIHR_PILOT_END_DATE=2026-07-10 AIHR_PILOT_STRICT=true ./scripts/demo-check.sh`。脚本只接受安全租户编号和 `YYYY-MM-DD` 日期,只统计目标租户窗口内完成的训练、校准和 SOP 评审;人员先按唯一手机号映射到在职组织快照,完训口径为每人至少 10 次已完成对练。
@@ -127,8 +129,8 @@ AIHR_AI_SPEECH_ENABLED=true
- `aihr_practice_mysql8.sql` 已导入;移动端员工训练记录落 `aihr_practice_session`,用于训练历史、主管待复盘列表和能力画像聚合
- `aihr_interview_result_mysql8.sql` 已纳入 reset 脚本;AI 面试评分完成后结果落 `aihr_interview_result`
- `aihr_candidate_material_mysql8.sql` 已纳入;候选人端补充资料文件写 `sys_oss`/MinIO,关系落 `aihr_candidate_material`
- `aihr_org_snapshot_mysql8.sql` 已纳入 reset 脚本;组织人员本地 seed(2 个住宅项目 22 人,项目经理/主管/一线三层)支撑演示,外部开放组织系统配置完成后用 `POST /api/aihr/org/sync` 拉取 `company/department/employee` 快照并覆盖本地 `aihr_org_snapshot`。2026-07-15 生产已完成该配置和覆盖同步,线上快照为 3417 名员工。
- 组织同步生产默认关闭 `aihr.org-sync.store-display-fields`,不把外部姓名/部门写入或返回组织人员展示快照;开发环境显式打开该开关仅用于 Demo。项目范围、岗位和外部主体 ID仍用于权限与身份映射。
- `aihr_org_snapshot_mysql8.sql` 已纳入 reset 脚本;组织人员本地 seed(2 个住宅项目 22 人,项目经理/主管/一线三层)支撑演示,外部开放组织系统配置完成后用 `POST /api/aihr/org/sync` 拉取 `company/department/employee` 快照。2026-07-25 生产快照为 3417 名员工、2938 人在职、3392 个手机号映射;在职姓名已按稳定 `employee_id` 定向补齐,本次未做全量覆盖。
- 源码 `application-prod.yml` 默认关闭 `aihr.org-sync.store-display-fields`;生产当前通过外部配置显式开启,用于主管团队、派发和考试对象显示已授权姓名。关闭该开关会让接口回退为编号/“员工”,修改前必须确认隐私口径与业务验收影响。项目范围、岗位和外部主体 ID仍用于权限与身份映射。
- 本地组织 seed 带演示手机号;从 `aihr_org_snapshot` 按角色查询测试账号,分别验证员工端岗位识别和主管端项目范围,不在文档保存具体号码
- 生产移动端回归同样由执行者从 `aihr_org_snapshot` 只读选择账号,不等待人工提供手机号:限定目标岗位、`employment_status='active'`、11 位 `person_phone`,并确认手机号与 `ext_party_id` 双向唯一。账号值只注入当前浏览器进程,不输出到终端、截图、报告或仓库。试点固定码开启时,登录顺序仍为先请求 `/resource/sms/code`、再提交固定码;直接提交固定码会得到 `Captcha invalid`,不应误判为账号或功能故障。
- SOP 知识库支持 `.txt/.md/.markdown/.pdf/.doc/.docx/.xls/.xlsx/.ppt/.pptx` 上传到 MinIO 后解析入库,接口为 `POST /api/knowledge/doc/upload`,单文件上限 100MB;管理端上传请求单独放宽到 180s,PDF 解析/归类/向量化较慢时不要改全局 axios 超时。