diff --git a/web-app/src/main/java/com/zcloud/primeport/service/DaHuaVehicleAccessRecordSyncApplicationService.java b/web-app/src/main/java/com/zcloud/primeport/service/DaHuaVehicleAccessRecordSyncApplicationService.java new file mode 100644 index 0000000..3d03f51 --- /dev/null +++ b/web-app/src/main/java/com/zcloud/primeport/service/DaHuaVehicleAccessRecordSyncApplicationService.java @@ -0,0 +1,434 @@ +package com.zcloud.primeport.service; + +import cn.hutool.json.JSONUtil; +import com.zcloud.primeport.domain.gateway.DaHuaGateway; +import com.zcloud.primeport.domain.gateway.DaHuaResourceRepositoryGateway; +import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordRepositoryGateway; +import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordSyncGateway; +import com.zcloud.primeport.domain.model.CorpInfoSnapshotE; +import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE; +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE; +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE; +import com.zcloud.primeport.domain.model.VehicleApplySnapshotE; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; + +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneId; +import java.time.format.DateTimeFormatter; +import java.time.format.DateTimeParseException; +import java.text.Normalizer; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Set; +import java.util.function.Consumer; + +@Service +@RequiredArgsConstructor +public class DaHuaVehicleAccessRecordSyncApplicationService + implements DaHuaVehicleAccessRecordSyncGateway { + + private static final DateTimeFormatter DATE_TIME_FORMATTER = + DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); + private static final int MAX_EXCEPTION_DETAILS = 100; + private static final int MAX_PAGE_COUNT = 10000; + + private final DaHuaGateway daHuaGateway; + private final DaHuaResourceRepositoryGateway resourceRepositoryGateway; + private final DaHuaVehicleAccessRecordRepositoryGateway recordRepositoryGateway; + + @Override + public DaHuaVehicleAccessRecordSyncResultE syncIncremental(int pageSize, int delayMinutes, + int overlapMinutes, + int initialLookbackMinutes) { + validatePageSize(pageSize); + LocalDateTime end = LocalDateTime.now().minusMinutes(Math.max(delayMinutes, 0)); + LocalDateTime checkpoint = recordRepositoryGateway.findSyncCheckpoint(); + LocalDateTime start = checkpoint == null + ? end.minusMinutes(Math.max(initialLookbackMinutes, 1)) + : checkpoint.minusMinutes(Math.max(overlapMinutes, 0)); + if (!start.isBefore(end)) { + return new DaHuaVehicleAccessRecordSyncResultE(); + } + DaHuaVehicleAccessRecordSyncResultE result = sync(format(start), format(end), pageSize); + recordRepositoryGateway.saveSyncCheckpoint(end); + return result; + } + + @Override + public DaHuaVehicleAccessRecordSyncResultE sync(String startTime, String endTime, int pageSize) { + LocalDateTime start = parseRequired(startTime, "startTime"); + LocalDateTime end = parseRequired(endTime, "endTime"); + if (start.isAfter(end)) { + throw new IllegalArgumentException("startTime cannot be after endTime"); + } + validatePageSize(pageSize); + SyncIndex index = buildIndex(); + Map synchronizedRecords = new HashMap<>(); + DaHuaVehicleAccessRecordSyncResultE result = new DaHuaVehicleAccessRecordSyncResultE(); + syncWindow("enterTimeStrLeft", "enterTimeStrRight", format(start), format(end), + pageSize, index, synchronizedRecords, result); + syncWindow("exitTimeStrLeft", "exitTimeStrRight", format(start), format(end), + pageSize, index, synchronizedRecords, result); + if (result.getRecordInvalid() > 0) { + throw new IllegalStateException("Dahua vehicle access record sync contains " + + result.getRecordInvalid() + " invalid record(s): " + + String.join("; ", result.getExceptionDetails())); + } + return result; + } + + private void syncWindow(String leftField, String rightField, String start, String end, + int pageSize, SyncIndex index, + Map synchronizedRecords, + DaHuaVehicleAccessRecordSyncResultE result) { + Set seenPageSignatures = new HashSet<>(); + for (int pageNum = 1; pageNum <= MAX_PAGE_COUNT; pageNum++) { + Map params = new LinkedHashMap<>(); + params.put("pageNum", pageNum); + params.put("pageSize", pageSize); + params.put(leftField, start); + params.put(rightField, end); + Map response = daHuaGateway.queryVehicleAccessRecord(params); + if (response == null || !isSuccessful(response)) { + throw new IllegalStateException("Dahua vehicle access query failed: code=" + + (response == null ? null : response.get("code")) + ", errMsg=" + + (response == null ? "empty response" : response.get("errMsg"))); + } + Map data = toMap(response.get("data")); + List> pageData = toRecordList(data.get("pageData")); + result.setPageCount(result.getPageCount() + 1); + if (pageData.isEmpty()) { + return; + } + String pageSignature = JSONUtil.toJsonStr(pageData); + if (!seenPageSignatures.add(pageSignature)) { + throw new IllegalStateException("Dahua vehicle access query returned a repeated page: " + + pageNum); + } + persistPage(pageData, index, synchronizedRecords, result); + } + throw new IllegalStateException("Dahua vehicle access query exceeded maximum page count: " + + MAX_PAGE_COUNT); + } + + private void persistPage(List> pageData, SyncIndex index, + Map synchronizedRecords, + DaHuaVehicleAccessRecordSyncResultE result) { + Set accessIds = new LinkedHashSet<>(); + Set plates = new LinkedHashSet<>(); + for (Map source : pageData) { + String accessId = blankToNull(firstString(source, "accessId", "id")); + if (accessId != null && !synchronizedRecords.containsKey(accessId)) { + accessIds.add(accessId); + } + addPlate(plates, firstString(source, "carNum")); + addPlate(plates, firstString(source, "exitCarNum")); + } + for (DaHuaVehicleAccessRecordE existing : recordRepositoryGateway.listByDahuaAccessIds(accessIds)) { + synchronizedRecords.put(existing.getDahuaAccessId(), existing); + } + for (VehicleApplySnapshotE snapshot : nullSafe(resourceRepositoryGateway.listVehicleApplySnapshots(plates))) { + index.vehicleSnapshotsByPlate.computeIfAbsent(normalizePlate(snapshot.getLicenceNo()), + ignored -> new ArrayList<>()).add(snapshot); + } + + for (Map source : pageData) { + result.setRecordTotal(result.getRecordTotal() + 1); + try { + String accessId = blankToNull(firstString(source, "accessId", "id")); + if (accessId == null) { + throw new IllegalArgumentException("accessId is missing"); + } + DaHuaVehicleAccessRecordE record = synchronizedRecords.get(accessId); + boolean created = record == null; + if (created) { + record = new DaHuaVehicleAccessRecordE(); + } + mapRecord(record, source, accessId, index); + if (created) { + recordRepositoryGateway.add(record); + result.setRecordCreated(result.getRecordCreated() + 1); + } else { + recordRepositoryGateway.update(record); + result.setRecordUpdated(result.getRecordUpdated() + 1); + } + synchronizedRecords.put(accessId, record); + } catch (Exception e) { + result.setRecordInvalid(result.getRecordInvalid() + 1); + if (result.getExceptionDetails().size() < MAX_EXCEPTION_DETAILS) { + result.getExceptionDetails().add("Vehicle access record sync failed: " + e.getMessage()); + } + } + } + } + + private void mapRecord(DaHuaVehicleAccessRecordE record, Map source, + String accessId, SyncIndex index) { + String departmentName = blankToNull(firstString(source, "departmentName", "name")); + Long departmentId = firstLong(source, "departmentId"); + String orgCode = blankToNull(firstString(source, "orgCode")); + DaHuaDepartmentCorpMappingE mapping = resolveMapping(index, departmentId, orgCode, departmentName); + VehicleApplySnapshotE vehicle = resolveVehicle(index, source); + CorpInfoSnapshotE corp = mapping == null || mapping.getCorpId() == null + ? null : index.corpsById.get(mapping.getCorpId()); + if (mapping != null) { + setIfNotNull(record::setCorpId, mapping.getCorpId()); + setIfNotNull(record::setCorpName, + blankToNull(corp == null ? mapping.getCorpName() : corp.getCorpName())); + setIfNotNull(record::setPortArea, corp == null ? null : corp.getPortArea()); + } else if (vehicle != null) { + corp = vehicle.getVehicleCorpId() == null ? null : index.corpsById.get(vehicle.getVehicleCorpId()); + setIfNotNull(record::setCorpId, vehicle.getVehicleCorpId()); + setIfNotNull(record::setCorpName, + blankToNull(corp == null ? vehicle.getVehicleCorpName() : corp.getCorpName())); + setIfNotNull(record::setPortArea, corp == null ? null : corp.getPortArea()); + } + + record.setDahuaAccessId(accessId); + setIfNotNull(record::setDahuaDepartmentId, departmentId != null ? departmentId + : mapping == null ? (vehicle == null ? null : vehicle.getVehicleDepartmentId()) + : mapping.getDahuaDepartmentId()); + setIfNotNull(record::setDahuaDepartmentName, departmentName != null ? departmentName + : blankToNull(mapping == null ? (vehicle == null ? null : vehicle.getVehicleDepartmentName()) + : mapping.getDahuaDepartmentName())); + setIfNotNull(record::setDahuaOrgCode, + orgCode != null ? orgCode : blankToNull(mapping == null ? null : mapping.getDahuaOrgCode())); + + setIfNotNull(record::setParkingLotCode, blankToNull(firstString(source, "parkingLotCode"))); + setIfNotNull(record::setParkingLotName, blankToNull(firstString(source, "parkingLot"))); + setIfNotNull(record::setCarType, firstInteger(source, "carType")); + setIfNotNull(record::setCarTypeName, blankToNull(firstString(source, "carTypeStr"))); + setIfNotNull(record::setCarNum, blankToNull(firstString(source, "carNum"))); + setIfNotNull(record::setExitCarNum, blankToNull(firstString(source, "exitCarNum"))); + + setIfNotNull(record::setEnterItcChannelCode, blankToNull(firstString(source, "enterItcDevChnid"))); + setIfNotNull(record::setEnterSluiceChannelCode, blankToNull(firstString(source, "enterSluiceDevChnid"))); + setIfNotNull(record::setEnterSluiceChannelName, blankToNull(firstString(source, "enterSluiceDevChnname"))); + setIfNotNull(record::setEnterTime, firstDateTime(source, "enterTimeStr", "enterTime")); + setIfNotNull(record::setEnterImageUrl, blankToNull(firstString(source, "enterImg"))); + + setIfNotNull(record::setExitItcChannelCode, blankToNull(firstString(source, "exitItcDevChnid"))); + setIfNotNull(record::setExitSluiceChannelCode, blankToNull(firstString(source, "exitSluiceDevChnid"))); + setIfNotNull(record::setExitSluiceChannelName, blankToNull(firstString(source, "exitSluiceDevChnname"))); + setIfNotNull(record::setExitTime, firstDateTime(source, "exitTimeStr", "exitTime")); + setIfNotNull(record::setExitImageUrl, blankToNull(firstString(source, "exitImg"))); + record.setSourceType("PULL"); + record.setRawData(JSONUtil.toJsonStr(source)); + record.setLastSyncTime(LocalDateTime.now()); + record.setDeleteEnum("FALSE"); + } + + private SyncIndex buildIndex() { + SyncIndex index = new SyncIndex(); + for (CorpInfoSnapshotE corp : nullSafe(resourceRepositoryGateway.listActiveCorps())) { + index.corpsById.put(corp.getId(), corp); + } + for (DaHuaDepartmentCorpMappingE mapping : nullSafe(resourceRepositoryGateway.listMappings())) { + if (mapping == null || "TRUE".equalsIgnoreCase(mapping.getDeleteEnum()) + || (mapping.getBindStatus() != null && mapping.getBindStatus() != 1)) { + continue; + } + if (mapping.getDahuaDepartmentId() != null) { + index.mappingsByDepartmentId.put(mapping.getDahuaDepartmentId(), mapping); + } + if (blankToNull(mapping.getDahuaOrgCode()) != null) { + index.mappingsByOrgCode.put(mapping.getDahuaOrgCode(), mapping); + } + String name = normalizeName(mapping.getDahuaDepartmentName()); + if (name != null) { + index.mappingsByName.computeIfAbsent(name, ignored -> new ArrayList<>()).add(mapping); + } + } + return index; + } + + private DaHuaDepartmentCorpMappingE resolveMapping(SyncIndex index, Long departmentId, + String orgCode, String departmentName) { + DaHuaDepartmentCorpMappingE mapping = departmentId == null ? null + : index.mappingsByDepartmentId.get(departmentId); + if (mapping == null && orgCode != null) { + mapping = index.mappingsByOrgCode.get(orgCode); + } + if (mapping == null) { + List matches = index.mappingsByName.get(normalizeName(departmentName)); + mapping = matches != null && matches.size() == 1 ? matches.get(0) : null; + } + return mapping; + } + + private VehicleApplySnapshotE resolveVehicle(SyncIndex index, Map source) { + VehicleApplySnapshotE enter = uniqueVehicle(index.vehicleSnapshotsByPlate, + firstString(source, "carNum")); + VehicleApplySnapshotE exit = uniqueVehicle(index.vehicleSnapshotsByPlate, + firstString(source, "exitCarNum")); + return enter != null ? enter : exit; + } + + private VehicleApplySnapshotE uniqueVehicle(Map> snapshots, + String licenceNo) { + List matches = snapshots.get(normalizePlate(licenceNo)); + if (matches == null || matches.isEmpty()) { + return null; + } + Set corpIds = new HashSet<>(); + for (VehicleApplySnapshotE match : matches) { + if (match.getVehicleCorpId() != null) { + corpIds.add(match.getVehicleCorpId()); + } + } + return corpIds.size() <= 1 ? matches.get(0) : null; + } + + private void addPlate(Set plates, String plate) { + if (blankToNull(plate) != null) { + plates.add(plate); + String normalized = normalizePlate(plate); + if (normalized != null) { + plates.add(normalized); + } + } + } + + private boolean isSuccessful(Map response) { + if (Boolean.TRUE.equals(response.get("passed")) || Boolean.TRUE.equals(response.get("success"))) { + return true; + } + String code = String.valueOf(response.get("code")); + return "0".equals(code) || "200".equals(code); + } + + private LocalDateTime firstDateTime(Map source, String textField, String epochField) { + Object value = source.get(textField); + if (value == null) { + value = source.get(epochField); + } + if (value == null) { + return null; + } + String text = blankToNull(value.toString()); + if (text == null) { + return null; + } + if (value instanceof Number || text.matches("\\d{10,}")) { + long epoch = Long.parseLong(text); + Instant instant = text.length() <= 10 + ? Instant.ofEpochSecond(epoch) : Instant.ofEpochMilli(epoch); + return instant + .atZone(ZoneId.systemDefault()).toLocalDateTime(); + } + try { + return LocalDateTime.parse(text, DATE_TIME_FORMATTER); + } catch (DateTimeParseException e) { + throw new IllegalArgumentException(textField + " must use yyyy-MM-dd HH:mm:ss", e); + } + } + + private LocalDateTime parseRequired(String value, String field) { + if (blankToNull(value) == null) { + throw new IllegalArgumentException(field + " cannot be blank"); + } + try { + return LocalDateTime.parse(value.trim(), DATE_TIME_FORMATTER); + } catch (DateTimeParseException e) { + throw new IllegalArgumentException(field + " must use yyyy-MM-dd HH:mm:ss", e); + } + } + + private String format(LocalDateTime value) { + return value.format(DATE_TIME_FORMATTER); + } + + private void validatePageSize(int pageSize) { + if (pageSize < 1 || pageSize > 1000) { + throw new IllegalArgumentException("pageSize must be between 1 and 1000"); + } + } + + private Object firstValue(Map source, String... fields) { + for (String field : fields) { + if (source.get(field) != null) { + return source.get(field); + } + } + return null; + } + + private String firstString(Map source, String... fields) { + Object value = firstValue(source, fields); + return value == null ? null : value.toString(); + } + + private Long firstLong(Map source, String... fields) { + Object value = firstValue(source, fields); + return value == null ? null : Long.valueOf(value.toString()); + } + + private Integer firstInteger(Map source, String... fields) { + Object value = firstValue(source, fields); + return value == null ? null : Integer.valueOf(value.toString()); + } + + @SuppressWarnings("unchecked") + private Map toMap(Object value) { + if (value instanceof Map) { + return new LinkedHashMap<>((Map) value); + } + return value == null ? new LinkedHashMap<>() : JSONUtil.parseObj(JSONUtil.toJsonStr(value)); + } + + @SuppressWarnings("unchecked") + private List> toRecordList(Object value) { + if (!(value instanceof Iterable)) { + return Collections.emptyList(); + } + List> result = new ArrayList<>(); + for (Object item : (Iterable) value) { + result.add(item instanceof Map ? new LinkedHashMap<>((Map) item) + : JSONUtil.parseObj(JSONUtil.toJsonStr(item))); + } + return result; + } + + private Collection nullSafe(Collection values) { + return values == null ? Collections.emptyList() : values; + } + + private void setIfNotNull(Consumer setter, T value) { + if (value != null) { + setter.accept(value); + } + } + + private String normalizeName(String value) { + String normalized = blankToNull(value); + return normalized == null ? null : Normalizer.normalize(normalized, Normalizer.Form.NFKC) + .replaceAll("\\s+", "").toLowerCase(Locale.ROOT); + } + + private String normalizePlate(String value) { + String normalized = blankToNull(value); + return normalized == null ? null : normalized.replaceAll("\\s+", "").toUpperCase(Locale.ROOT); + } + + private String blankToNull(String value) { + return value == null || value.trim().isEmpty() ? null : value.trim(); + } + + private static class SyncIndex { + private final Map corpsById = new HashMap<>(); + private final Map mappingsByDepartmentId = new HashMap<>(); + private final Map mappingsByOrgCode = new HashMap<>(); + private final Map> mappingsByName = new HashMap<>(); + private final Map> vehicleSnapshotsByPlate = new HashMap<>(); + } +} diff --git a/web-app/src/test/java/com/zcloud/primeport/service/DaHuaVehicleAccessRecordSyncApplicationServiceTest.java b/web-app/src/test/java/com/zcloud/primeport/service/DaHuaVehicleAccessRecordSyncApplicationServiceTest.java new file mode 100644 index 0000000..4b7fe37 --- /dev/null +++ b/web-app/src/test/java/com/zcloud/primeport/service/DaHuaVehicleAccessRecordSyncApplicationServiceTest.java @@ -0,0 +1,319 @@ +package com.zcloud.primeport.service; + +import com.zcloud.primeport.domain.gateway.DaHuaGateway; +import com.zcloud.primeport.domain.gateway.DaHuaResourceRepositoryGateway; +import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordRepositoryGateway; +import com.zcloud.primeport.domain.model.CorpInfoSnapshotE; +import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE; +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE; +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE; +import com.zcloud.primeport.domain.model.VehicleApplySnapshotE; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; + +import static org.junit.jupiter.api.Assertions.assertEquals; +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.anyCollection; +import static org.mockito.ArgumentMatchers.anyMap; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class DaHuaVehicleAccessRecordSyncApplicationServiceTest { + + @Mock + private DaHuaGateway daHuaGateway; + @Mock + private DaHuaResourceRepositoryGateway resourceRepositoryGateway; + @Mock + private DaHuaVehicleAccessRecordRepositoryGateway recordRepositoryGateway; + + private DaHuaVehicleAccessRecordSyncApplicationService applicationService; + + @BeforeEach + void setUp() { + applicationService = new DaHuaVehicleAccessRecordSyncApplicationService( + daHuaGateway, resourceRepositoryGateway, recordRepositoryGateway); + when(resourceRepositoryGateway.listActiveCorps()).thenReturn(Collections.emptyList()); + when(resourceRepositoryGateway.listMappings()).thenReturn(Collections.emptyList()); + lenient().when(resourceRepositoryGateway.listVehicleApplySnapshots(anyCollection())) + .thenReturn(Collections.emptyList()); + lenient().when(recordRepositoryGateway.listByDahuaAccessIds(anyCollection())) + .thenReturn(Collections.emptyList()); + } + + @Test + void shouldQueryBothWindowsAndUpdateTheSameAccessIdIdempotently() { + Map enter = record("A-1"); + enter.put("departmentId", 10L); + enter.put("enterSluiceDevChnid", "CH-ENTER"); + enter.put("enterTimeStr", "2026-08-01 10:00:00"); + Map exit = record("A-1"); + exit.put("departmentId", 10L); + exit.put("exitSluiceDevChnid", "CH-EXIT"); + exit.put("exitTimeStr", "2026-08-01 18:00:00"); + + when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> { + Map params = invocation.getArgument(0); + if (Integer.valueOf(1).equals(params.get("pageNum")) + && params.containsKey("enterTimeStrLeft")) { + return response(Collections.singletonList(enter)); + } + if (Integer.valueOf(1).equals(params.get("pageNum")) + && params.containsKey("exitTimeStrLeft")) { + return response(Collections.singletonList(exit)); + } + return response(Collections.emptyList()); + }); + + DaHuaVehicleAccessRecordSyncResultE result = applicationService.sync( + "2026-08-01 00:00:00", "2026-08-02 00:00:00", 100); + + assertEquals(4, result.getPageCount()); + assertEquals(2, result.getRecordTotal()); + assertEquals(1, result.getRecordCreated()); + assertEquals(1, result.getRecordUpdated()); + verify(recordRepositoryGateway).add(any(DaHuaVehicleAccessRecordE.class)); + ArgumentCaptor updateCaptor = + ArgumentCaptor.forClass(DaHuaVehicleAccessRecordE.class); + verify(recordRepositoryGateway).update(updateCaptor.capture()); + assertEquals("A-1", updateCaptor.getValue().getDahuaAccessId()); + assertEquals("2026-08-01 18:00:00", + updateCaptor.getValue().getExitTime().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))); + assertEquals("2026-08-01 10:00:00", + updateCaptor.getValue().getEnterTime().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))); + assertEquals("CH-ENTER", updateCaptor.getValue().getEnterSluiceChannelCode()); + verify(daHuaGateway, times(4)).queryVehicleAccessRecord(anyMap()); + } + + @Test + void shouldFetchEveryPageForEachTimeWindow() { + Map first = record("A-1"); + Map second = record("A-2"); + when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> { + Map params = invocation.getArgument(0); + Integer pageNum = (Integer) params.get("pageNum"); + if (params.containsKey("enterTimeStrLeft")) { + return pageNum == 1 ? response(Collections.singletonList(first)) + : pageNum == 2 ? response(Collections.singletonList(second)) + : response(Collections.emptyList()); + } + return response(Collections.emptyList()); + }); + + DaHuaVehicleAccessRecordSyncResultE result = applicationService.sync( + "2026-08-01 00:00:00", "2026-08-02 00:00:00", 1); + + assertEquals(4, result.getPageCount()); + assertEquals(2, result.getRecordTotal()); + assertEquals(2, result.getRecordCreated()); + verify(daHuaGateway).queryVehicleAccessRecord(params("enterTimeStrLeft", 1)); + verify(daHuaGateway).queryVehicleAccessRecord(params("enterTimeStrLeft", 2)); + verify(daHuaGateway).queryVehicleAccessRecord(params("enterTimeStrLeft", 3)); + verify(daHuaGateway).queryVehicleAccessRecord(params("exitTimeStrLeft", 1)); + } + + @Test + void shouldResolveCorpByDepartmentAndKeepChannelData() { + CorpInfoSnapshotE personCorp = corp(100L, "人员企业", 1); + when(resourceRepositoryGateway.listActiveCorps()) + .thenReturn(Collections.singletonList(personCorp)); + + DaHuaDepartmentCorpMappingE mapping = new DaHuaDepartmentCorpMappingE(); + mapping.setDahuaDepartmentId(10L); + mapping.setDahuaDepartmentName("部门"); + mapping.setCorpId(100L); + mapping.setCorpName("人员企业"); + mapping.setBindStatus(1); + mapping.setDeleteEnum("FALSE"); + when(resourceRepositoryGateway.listMappings()).thenReturn(Collections.singletonList(mapping)); + + when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> { + Map params = invocation.getArgument(0); + return Integer.valueOf(1).equals(params.get("pageNum")) + && params.containsKey("enterTimeStrLeft") + ? response(Collections.singletonList(recordWithDepartmentAndChannel("A-1"))) + : response(Collections.emptyList()); + }); + + applicationService.sync("2026-08-01 00:00:00", "2026-08-02 00:00:00", 100); + + ArgumentCaptor captor = + ArgumentCaptor.forClass(DaHuaVehicleAccessRecordE.class); + verify(recordRepositoryGateway).add(captor.capture()); + DaHuaVehicleAccessRecordE saved = captor.getValue(); + assertEquals(Long.valueOf(100L), saved.getCorpId()); + assertEquals("人员企业", saved.getCorpName()); + assertEquals(Integer.valueOf(1), saved.getPortArea()); + assertEquals("CH-ENTER", saved.getEnterSluiceChannelCode()); + } + + @Test + void shouldResolveCorpByVehiclePlateWhenDepartmentIsUnmapped() { + VehicleApplySnapshotE vehicle = new VehicleApplySnapshotE(); + vehicle.setLicenceNo("冀A12345"); + vehicle.setVehicleCorpId(300L); + vehicle.setVehicleCorpName("车辆企业"); + vehicle.setVehicleDepartmentId(30L); + vehicle.setVehicleDepartmentName("车辆部门"); + when(resourceRepositoryGateway.listVehicleApplySnapshots(anyCollection())) + .thenReturn(Collections.singletonList(vehicle)); + when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> { + Map params = invocation.getArgument(0); + return Integer.valueOf(1).equals(params.get("pageNum")) + && params.containsKey("enterTimeStrLeft") + ? response(Collections.singletonList(recordWithPlate("A-1", "冀A12345"))) + : response(Collections.emptyList()); + }); + + applicationService.sync("2026-08-01 00:00:00", "2026-08-02 00:00:00", 100); + + ArgumentCaptor captor = + ArgumentCaptor.forClass(DaHuaVehicleAccessRecordE.class); + verify(recordRepositoryGateway).add(captor.capture()); + assertEquals(Long.valueOf(300L), captor.getValue().getCorpId()); + assertEquals("车辆企业", captor.getValue().getCorpName()); + assertEquals(Long.valueOf(30L), captor.getValue().getDahuaDepartmentId()); + assertEquals("车辆部门", captor.getValue().getDahuaDepartmentName()); + } + + @Test + void shouldFailForMissingAccessIdAndRejectOversizedPage() { + when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> { + Map params = invocation.getArgument(0); + return params.containsKey("enterTimeStrLeft") + && Integer.valueOf(1).equals(params.get("pageNum")) + ? response(Collections.singletonList(record(null))) + : response(Collections.emptyList()); + }); + + IllegalStateException exception = assertThrows(IllegalStateException.class, + () -> applicationService.sync( + "2026-08-01 00:00:00", "2026-08-02 00:00:00", 100)); + assertTrue(exception.getMessage().contains("accessId is missing")); + verify(recordRepositoryGateway, never()).add(any(DaHuaVehicleAccessRecordE.class)); + assertThrows(IllegalArgumentException.class, () -> applicationService.sync( + "2026-08-01 00:00:00", "2026-08-02 00:00:00", 1001)); + } + + @Test + @SuppressWarnings("unchecked") + void shouldAdvanceCheckpointOnlyAfterSuccessfulIncrementalSync() { + LocalDateTime checkpoint = LocalDateTime.of(2026, 8, 1, 12, 0, 0); + when(recordRepositoryGateway.findSyncCheckpoint()).thenReturn(checkpoint); + when(daHuaGateway.queryVehicleAccessRecord(anyMap())) + .thenReturn(response(Collections.emptyList())); + + applicationService.syncIncremental(100, 5, 10, 1440); + + ArgumentCaptor> paramsCaptor = ArgumentCaptor.forClass(Map.class); + verify(daHuaGateway, times(2)).queryVehicleAccessRecord(paramsCaptor.capture()); + assertEquals("2026-08-01 11:50:00", + paramsCaptor.getAllValues().get(0).get("enterTimeStrLeft")); + + ArgumentCaptor checkpointCaptor = ArgumentCaptor.forClass(LocalDateTime.class); + verify(recordRepositoryGateway).saveSyncCheckpoint(checkpointCaptor.capture()); + assertEquals(checkpointCaptor.getValue().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")), + paramsCaptor.getAllValues().get(0).get("enterTimeStrRight")); + } + + @Test + void shouldNotAdvanceCheckpointWhenIncrementalSyncContainsInvalidRecord() { + when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> { + Map params = invocation.getArgument(0); + return params.containsKey("enterTimeStrLeft") + && Integer.valueOf(1).equals(params.get("pageNum")) + ? response(Collections.singletonList(record(null))) + : response(Collections.emptyList()); + }); + + assertThrows(IllegalStateException.class, + () -> applicationService.syncIncremental(100, 5, 10, 1440)); + + verify(recordRepositoryGateway, never()).saveSyncCheckpoint(any(LocalDateTime.class)); + } + + @Test + void shouldFailWhenDahuaRepeatsTheSamePage() { + when(daHuaGateway.queryVehicleAccessRecord(anyMap())) + .thenReturn(response(Collections.singletonList(record("A-1")))); + + IllegalStateException exception = assertThrows(IllegalStateException.class, + () -> applicationService.sync( + "2026-08-01 00:00:00", "2026-08-02 00:00:00", 100)); + + assertTrue(exception.getMessage().contains("repeated page")); + verify(recordRepositoryGateway, times(1)).add(any(DaHuaVehicleAccessRecordE.class)); + } + + private Map params(String field, int pageNum) { + Map params = new HashMap<>(); + params.put(field, "2026-08-01 00:00:00"); + params.put("pageNum", pageNum); + params.put("pageSize", 1); + return org.mockito.ArgumentMatchers.argThat(value -> + value.equals(params) || (value.containsKey(field) + && Integer.valueOf(pageNum).equals(value.get("pageNum")))); + } + + private CorpInfoSnapshotE corp(Long id, String name, Integer portArea) { + CorpInfoSnapshotE corp = new CorpInfoSnapshotE(); + corp.setId(id); + corp.setCorpName(name); + corp.setPortArea(portArea); + return corp; + } + + private Map record(String id) { + Map value = new HashMap<>(); + if (id != null) { + value.put("accessId", id); + } + value.put("carNum", "冀A12345"); + value.put("enterTimeStr", "2026-08-01 10:00:00"); + return value; + } + + private Map recordWithDepartmentAndChannel(String id) { + Map value = record(id); + value.put("departmentId", 10L); + value.put("departmentName", "部门"); + value.put("enterSluiceDevChnid", "CH-ENTER"); + return value; + } + + private Map recordWithPlate(String id, String plate) { + Map value = record(id); + value.put("carNum", plate); + return value; + } + + private Map response(Collection> records) { + Map data = new HashMap<>(); + data.put("pageData", new ArrayList<>(records)); + Map response = new HashMap<>(); + response.put("passed", true); + response.put("success", true); + response.put("code", "0"); + response.put("data", data); + return response; + } +} diff --git a/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaGateway.java b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaGateway.java index b4c6751..baa8b6c 100644 --- a/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaGateway.java +++ b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaGateway.java @@ -31,9 +31,11 @@ public interface DaHuaGateway { Map queryAccessRecord(DaHuaAccessRecordCmd cmd); + Map queryVehicleAccessRecord(Map params); + String uploadPersonImg(String imageBase64); Map saveTempVehicle(DaHuaTempVehicleSaveCmd cmd); Long findDahuaDepartmentIdByCorpId(Long corpId); -} \ No newline at end of file +} diff --git a/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaResourceRepositoryGateway.java b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaResourceRepositoryGateway.java index cb47c3e..b3a3604 100644 --- a/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaResourceRepositoryGateway.java +++ b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaResourceRepositoryGateway.java @@ -4,7 +4,9 @@ import com.zcloud.primeport.domain.model.CorpInfoSnapshotE; import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE; import com.zcloud.primeport.domain.model.DaHuaDeviceE; import com.zcloud.primeport.domain.model.DaHuaDeviceChannelE; +import com.zcloud.primeport.domain.model.VehicleApplySnapshotE; +import java.util.Collection; import java.util.List; public interface DaHuaResourceRepositoryGateway { @@ -25,6 +27,8 @@ public interface DaHuaResourceRepositoryGateway { List listDeviceChannels(); + List listVehicleApplySnapshots(Collection licenceNos); + void addDeviceChannel(DaHuaDeviceChannelE channel); void updateDeviceChannel(DaHuaDeviceChannelE channel); diff --git a/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaVehicleAccessRecordRepositoryGateway.java b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaVehicleAccessRecordRepositoryGateway.java new file mode 100644 index 0000000..88eb114 --- /dev/null +++ b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaVehicleAccessRecordRepositoryGateway.java @@ -0,0 +1,20 @@ +package com.zcloud.primeport.domain.gateway; + +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE; + +import java.time.LocalDateTime; +import java.util.Collection; +import java.util.List; + +public interface DaHuaVehicleAccessRecordRepositoryGateway { + + List listByDahuaAccessIds(Collection accessIds); + + LocalDateTime findSyncCheckpoint(); + + void saveSyncCheckpoint(LocalDateTime checkpointTime); + + void add(DaHuaVehicleAccessRecordE record); + + void update(DaHuaVehicleAccessRecordE record); +} diff --git a/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaVehicleAccessRecordSyncGateway.java b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaVehicleAccessRecordSyncGateway.java new file mode 100644 index 0000000..a5bf4df --- /dev/null +++ b/web-domain/src/main/java/com/zcloud/primeport/domain/gateway/DaHuaVehicleAccessRecordSyncGateway.java @@ -0,0 +1,11 @@ +package com.zcloud.primeport.domain.gateway; + +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE; + +public interface DaHuaVehicleAccessRecordSyncGateway { + + DaHuaVehicleAccessRecordSyncResultE syncIncremental(int pageSize, int delayMinutes, + int overlapMinutes, int initialLookbackMinutes); + + DaHuaVehicleAccessRecordSyncResultE sync(String startTime, String endTime, int pageSize); +} diff --git a/web-domain/src/main/java/com/zcloud/primeport/domain/model/DaHuaVehicleAccessRecordE.java b/web-domain/src/main/java/com/zcloud/primeport/domain/model/DaHuaVehicleAccessRecordE.java new file mode 100644 index 0000000..2c7112f --- /dev/null +++ b/web-domain/src/main/java/com/zcloud/primeport/domain/model/DaHuaVehicleAccessRecordE.java @@ -0,0 +1,45 @@ +package com.zcloud.primeport.domain.model; + +import lombok.Data; + +import java.time.LocalDateTime; + +@Data +public class DaHuaVehicleAccessRecordE { + + private Long id; + private String dahuaAccessId; + + private Long corpId; + private String corpName; + private Integer portArea; + + private Long dahuaDepartmentId; + private String dahuaDepartmentName; + private String dahuaOrgCode; + + private String parkingLotCode; + private String parkingLotName; + private Integer carType; + private String carTypeName; + + private String carNum; + private String exitCarNum; + + private String enterItcChannelCode; + private String enterSluiceChannelCode; + private String enterSluiceChannelName; + private LocalDateTime enterTime; + private String enterImageUrl; + + private String exitItcChannelCode; + private String exitSluiceChannelCode; + private String exitSluiceChannelName; + private LocalDateTime exitTime; + private String exitImageUrl; + + private String sourceType; + private String rawData; + private LocalDateTime lastSyncTime; + private String deleteEnum; +} diff --git a/web-domain/src/main/java/com/zcloud/primeport/domain/model/DaHuaVehicleAccessRecordSyncResultE.java b/web-domain/src/main/java/com/zcloud/primeport/domain/model/DaHuaVehicleAccessRecordSyncResultE.java new file mode 100644 index 0000000..96a32a1 --- /dev/null +++ b/web-domain/src/main/java/com/zcloud/primeport/domain/model/DaHuaVehicleAccessRecordSyncResultE.java @@ -0,0 +1,17 @@ +package com.zcloud.primeport.domain.model; + +import lombok.Data; + +import java.util.ArrayList; +import java.util.List; + +@Data +public class DaHuaVehicleAccessRecordSyncResultE { + + private int pageCount; + private int recordTotal; + private int recordCreated; + private int recordUpdated; + private int recordInvalid; + private List exceptionDetails = new ArrayList<>(); +} diff --git a/web-domain/src/main/java/com/zcloud/primeport/domain/model/VehicleApplySnapshotE.java b/web-domain/src/main/java/com/zcloud/primeport/domain/model/VehicleApplySnapshotE.java new file mode 100644 index 0000000..86c77ea --- /dev/null +++ b/web-domain/src/main/java/com/zcloud/primeport/domain/model/VehicleApplySnapshotE.java @@ -0,0 +1,16 @@ +package com.zcloud.primeport.domain.model; + +import lombok.Data; + +/** + * 车辆通行记录企业匹配所需的车辆档案快照。 + */ +@Data +public class VehicleApplySnapshotE { + + private String licenceNo; + private Long vehicleCorpId; + private String vehicleCorpName; + private Long vehicleDepartmentId; + private String vehicleDepartmentName; +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/api/DaHuaVehicleAccessRecordApi.java b/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/api/DaHuaVehicleAccessRecordApi.java new file mode 100644 index 0000000..83f33d1 --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/api/DaHuaVehicleAccessRecordApi.java @@ -0,0 +1,33 @@ +package com.zcloud.primeport.dahua.api; + +import com.dahuatech.icc.exception.ClientException; +import com.dahuatech.icc.oauth.model.v202010.GeneralResponse; +import com.zcloud.primeport.dahua.client.DaHuaApiClient; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; + +import java.util.HashMap; +import java.util.Map; + +@Component +public class DaHuaVehicleAccessRecordApi { + + private static final String PATH = "/evo-apigw/ipms/caraccess/find/his"; + + @Autowired(required = false) + private DaHuaApiClient daHuaApiClient; + + public boolean isConfigured() { + return daHuaApiClient != null; + } + + public GeneralResponse query(Map params) throws ClientException { + if (daHuaApiClient == null) { + return null; + } + String token = daHuaApiClient.getToken(); + Map headers = new HashMap<>(); + headers.put("Authorization", "bearer " + token); + return daHuaApiClient.getForm(PATH, params, headers); + } +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/config/DaHuaAutoConfig.java b/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/config/DaHuaAutoConfig.java index ada9518..b596f0e 100644 --- a/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/config/DaHuaAutoConfig.java +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/config/DaHuaAutoConfig.java @@ -9,7 +9,7 @@ import org.springframework.context.annotation.Configuration; @Configuration @EnableConfigurationProperties({DaHuaProperties.class, DaHuaResourceProperties.class, - DaHuaAccessRecordProperties.class}) + DaHuaAccessRecordProperties.class, DaHuaVehicleAccessRecordProperties.class}) public class DaHuaAutoConfig { @Bean("dahuaApiOauthConfig") diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/config/DaHuaVehicleAccessRecordProperties.java b/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/config/DaHuaVehicleAccessRecordProperties.java new file mode 100644 index 0000000..f172724 --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/dahua/config/DaHuaVehicleAccessRecordProperties.java @@ -0,0 +1,15 @@ +package com.zcloud.primeport.dahua.config; + +import lombok.Data; +import org.springframework.boot.context.properties.ConfigurationProperties; + +@Data +@ConfigurationProperties(prefix = "dahua.vehicle-access-record") +public class DaHuaVehicleAccessRecordProperties { + + private boolean syncEnabled = true; + private Integer syncPageSize = 1000; + private Integer syncDelayMinutes = 5; + private Integer syncOverlapMinutes = 10; + private Integer initialLookbackMinutes = 1440; +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaGatewayImpl.java b/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaGatewayImpl.java index b6f7a63..8dc4e4c 100644 --- a/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaGatewayImpl.java +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaGatewayImpl.java @@ -6,6 +6,7 @@ import cn.hutool.json.JSONUtil; import com.dahuatech.icc.exception.ClientException; import com.dahuatech.icc.oauth.model.v202010.GeneralResponse; import com.zcloud.primeport.dahua.api.DaHuaPersonApi; +import com.zcloud.primeport.dahua.api.DaHuaVehicleAccessRecordApi; import com.zcloud.primeport.domain.gateway.DaHuaGateway; import com.zcloud.gbscommon.dahua.cmd.DaHuaAccessRecordCmd; import com.zcloud.gbscommon.dahua.cmd.DaHuaCarAddCmd; @@ -33,6 +34,9 @@ public class DaHuaGatewayImpl implements DaHuaGateway { @Autowired(required = false) private DaHuaPersonApi daHuaPersonApi; + @Autowired(required = false) + private DaHuaVehicleAccessRecordApi daHuaVehicleAccessRecordApi; + @Autowired private DaHuaDepartmentCorpMappingRepository daHuaDeptMappingRepository; @@ -854,6 +858,39 @@ public class DaHuaGatewayImpl implements DaHuaGateway { } } + @Override + public Map queryVehicleAccessRecord(Map params) { + if (daHuaVehicleAccessRecordApi == null || !daHuaVehicleAccessRecordApi.isConfigured()) { + throw new RuntimeException("Dahua service is not configured"); + } + Map result = new HashMap<>(); + try { + GeneralResponse response = daHuaVehicleAccessRecordApi.query(params); + if (response == null) { + result.put("passed", false); + result.put("success", false); + result.put("code", "SYSTEM_ERROR"); + result.put("errMsg", "response is null"); + return result; + } + JSONObject json = JSONUtil.parseObj(JSONUtil.toJsonStr(response)); + result.put("success", json.getBool("success", false)); + result.put("code", json.getStr("code", response.getCode())); + result.put("errMsg", json.getStr("errMsg", "")); + result.put("data", json.get("data")); + result.put("passed", Boolean.TRUE.equals(result.get("success")) + || "0".equals(String.valueOf(result.get("code"))) + || "200".equals(String.valueOf(result.get("code")))); + return result; + } catch (ClientException e) { + result.put("passed", false); + result.put("success", false); + result.put("code", "CLIENT_EXCEPTION"); + result.put("errMsg", e.getMessage()); + return result; + } + } + @Override public String uploadPersonImg(String imageBase64) { if (daHuaPersonApi == null || !daHuaPersonApi.isConfigured()) { @@ -1202,4 +1239,4 @@ public class DaHuaGatewayImpl implements DaHuaGateway { } return mapping.getDahuaDepartmentId(); } -} \ No newline at end of file +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaResourceRepositoryGatewayImpl.java b/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaResourceRepositoryGatewayImpl.java index c9de675..fdfde20 100644 --- a/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaResourceRepositoryGatewayImpl.java +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaResourceRepositoryGatewayImpl.java @@ -8,20 +8,25 @@ import com.zcloud.primeport.domain.model.CorpInfoSnapshotE; import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE; import com.zcloud.primeport.domain.model.DaHuaDeviceE; import com.zcloud.primeport.domain.model.DaHuaDeviceChannelE; +import com.zcloud.primeport.domain.model.VehicleApplySnapshotE; import com.zcloud.primeport.persistence.dataobject.CorpInfoSnapshotDO; import com.zcloud.primeport.persistence.dataobject.DaHuaDepartmentCorpMappingDO; import com.zcloud.primeport.persistence.dataobject.DaHuaDeviceDO; import com.zcloud.primeport.persistence.dataobject.DaHuaDeviceChannelDO; +import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO; import com.zcloud.primeport.persistence.repository.CorpInfoSnapshotRepository; import com.zcloud.primeport.persistence.repository.DaHuaDepartmentCorpMappingRepository; import com.zcloud.primeport.persistence.repository.DaHuaDeviceRepository; import com.zcloud.primeport.persistence.repository.DaHuaDeviceChannelRepository; +import com.zcloud.primeport.persistence.repository.VehicleApplyRepository; import lombok.RequiredArgsConstructor; import org.springframework.beans.BeanUtils; import org.springframework.stereotype.Service; import java.time.LocalDateTime; import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; import java.util.List; @Service @@ -32,6 +37,7 @@ public class DaHuaResourceRepositoryGatewayImpl implements DaHuaResourceReposito private final DaHuaDepartmentCorpMappingRepository mappingRepository; private final DaHuaDeviceRepository deviceRepository; private final DaHuaDeviceChannelRepository deviceChannelRepository; + private final VehicleApplyRepository vehicleApplyRepository; private final DefaultIdentifierGenerator idGenerator = new DefaultIdentifierGenerator(); @Override @@ -89,6 +95,15 @@ public class DaHuaResourceRepositoryGatewayImpl implements DaHuaResourceReposito return convertList(deviceChannelRepository.list(new LambdaQueryWrapper<>()), DaHuaDeviceChannelE.class); } + @Override + public List listVehicleApplySnapshots(Collection licenceNos) { + if (licenceNos == null || licenceNos.isEmpty()) { + return Collections.emptyList(); + } + List snapshots = vehicleApplyRepository.listAccessMatchSnapshots(licenceNos); + return convertList(snapshots, VehicleApplySnapshotE.class); + } + @Override public void addDeviceChannel(DaHuaDeviceChannelE channel) { DaHuaDeviceChannelDO data = copy(channel, DaHuaDeviceChannelDO.class); diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaVehicleAccessRecordRepositoryGatewayImpl.java b/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaVehicleAccessRecordRepositoryGatewayImpl.java new file mode 100644 index 0000000..bf4a100 --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/gatewayimpl/DaHuaVehicleAccessRecordRepositoryGatewayImpl.java @@ -0,0 +1,90 @@ +package com.zcloud.primeport.gatewayimpl; + +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.baomidou.mybatisplus.core.incrementer.DefaultIdentifierGenerator; +import com.jjb.saas.framework.repository.basedo.BaseDO; +import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordRepositoryGateway; +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE; +import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO; +import com.zcloud.primeport.persistence.repository.DaHuaVehicleAccessRecordRepository; +import lombok.RequiredArgsConstructor; +import org.springframework.beans.BeanUtils; +import org.springframework.stereotype.Service; + +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.List; + +@Service +@RequiredArgsConstructor +public class DaHuaVehicleAccessRecordRepositoryGatewayImpl + implements DaHuaVehicleAccessRecordRepositoryGateway { + + private final DaHuaVehicleAccessRecordRepository repository; + private final DefaultIdentifierGenerator idGenerator = new DefaultIdentifierGenerator(); + + @Override + public List listByDahuaAccessIds(Collection accessIds) { + if (accessIds == null || accessIds.isEmpty()) { + return Collections.emptyList(); + } + LambdaQueryWrapper query = new LambdaQueryWrapper<>(); + query.in(DaHuaVehicleAccessRecordDO::getDahuaAccessId, accessIds); + return convertList(repository.list(query)); + } + + @Override + public LocalDateTime findSyncCheckpoint() { + return repository.findSyncCheckpoint(); + } + + @Override + public void saveSyncCheckpoint(LocalDateTime checkpointTime) { + repository.saveSyncCheckpoint(checkpointTime); + } + + @Override + public void add(DaHuaVehicleAccessRecordE record) { + DaHuaVehicleAccessRecordDO data = copy(record, DaHuaVehicleAccessRecordDO.class); + initializeBase(data); + repository.upsert(data); + record.setId(data.getId()); + } + + @Override + public void update(DaHuaVehicleAccessRecordE record) { + DaHuaVehicleAccessRecordDO data = copy(record, DaHuaVehicleAccessRecordDO.class); + data.setUpdateTime(LocalDateTime.now()); + repository.upsert(data); + } + + private void initializeBase(BaseDO target) { + LocalDateTime now = LocalDateTime.now(); + target.setId(idGenerator.nextId(target).longValue()); + target.setDeleteEnum("FALSE"); + target.setEnv("PROD"); + target.setVersion(0); + target.setCreateTime(now); + target.setUpdateTime(now); + } + + private List convertList(List source) { + List result = new ArrayList<>(); + for (DaHuaVehicleAccessRecordDO item : source) { + result.add(copy(item, DaHuaVehicleAccessRecordE.class)); + } + return result; + } + + private T copy(Object source, Class targetType) { + try { + T target = targetType.newInstance(); + BeanUtils.copyProperties(source, target); + return target; + } catch (InstantiationException | IllegalAccessException e) { + throw new IllegalStateException("Failed to convert vehicle access record", e); + } + } +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/dataobject/DaHuaVehicleAccessRecordDO.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/dataobject/DaHuaVehicleAccessRecordDO.java new file mode 100644 index 0000000..cd000b4 --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/dataobject/DaHuaVehicleAccessRecordDO.java @@ -0,0 +1,41 @@ +package com.zcloud.primeport.persistence.dataobject; + +import com.baomidou.mybatisplus.annotation.TableName; +import com.jjb.saas.framework.repository.basedo.BaseDO; +import lombok.Data; +import lombok.EqualsAndHashCode; + +import java.time.LocalDateTime; + +@Data +@TableName("dahua_vehicle_access_record") +@EqualsAndHashCode(callSuper = true) +public class DaHuaVehicleAccessRecordDO extends BaseDO { + + private String dahuaAccessId; + private Long corpId; + private String corpName; + private Integer portArea; + private Long dahuaDepartmentId; + private String dahuaDepartmentName; + private String dahuaOrgCode; + private String parkingLotCode; + private String parkingLotName; + private Integer carType; + private String carTypeName; + private String carNum; + private String exitCarNum; + private String enterItcChannelCode; + private String enterSluiceChannelCode; + private String enterSluiceChannelName; + private LocalDateTime enterTime; + private String enterImageUrl; + private String exitItcChannelCode; + private String exitSluiceChannelCode; + private String exitSluiceChannelName; + private LocalDateTime exitTime; + private String exitImageUrl; + private String sourceType; + private String rawData; + private LocalDateTime lastSyncTime; +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/dataobject/VehicleApplySnapshotDO.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/dataobject/VehicleApplySnapshotDO.java new file mode 100644 index 0000000..ac782ba --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/dataobject/VehicleApplySnapshotDO.java @@ -0,0 +1,15 @@ +package com.zcloud.primeport.persistence.dataobject; + +import com.baomidou.mybatisplus.annotation.TableName; +import lombok.Data; + +@Data +@TableName("vehicle_apply") +public class VehicleApplySnapshotDO { + + private String licenceNo; + private Long vehicleCorpId; + private String vehicleCorpName; + private Long vehicleDepartmentId; + private String vehicleDepartmentName; +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/mapper/DaHuaVehicleAccessRecordMapper.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/mapper/DaHuaVehicleAccessRecordMapper.java new file mode 100644 index 0000000..83574ec --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/mapper/DaHuaVehicleAccessRecordMapper.java @@ -0,0 +1,18 @@ +package com.zcloud.primeport.persistence.mapper; + +import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO; +import org.apache.ibatis.annotations.Mapper; +import org.apache.ibatis.annotations.Param; + +import java.time.LocalDateTime; + +@Mapper +public interface DaHuaVehicleAccessRecordMapper extends BaseMapper { + + LocalDateTime findSyncCheckpoint(); + + int upsertSyncCheckpoint(@Param("checkpointTime") LocalDateTime checkpointTime); + + int upsert(@Param("record") DaHuaVehicleAccessRecordDO record); +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/mapper/VehicleApplyMapper.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/mapper/VehicleApplyMapper.java index 2e792a7..72532ca 100644 --- a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/mapper/VehicleApplyMapper.java +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/mapper/VehicleApplyMapper.java @@ -9,6 +9,7 @@ import com.jjb.saas.framework.datascope.annotation.DataScopes; import com.zcloud.primeport.persistence.dataobject.AppCountDTO; import com.zcloud.primeport.persistence.dataobject.FgsVehicleCountDto; import com.zcloud.primeport.persistence.dataobject.VehicleApplyDO; +import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO; import com.zcloud.primeport.plan.mjDevice.dto.OneLevelCarSaveDto; import org.apache.ibatis.annotations.Mapper; import org.apache.ibatis.annotations.Param; @@ -16,6 +17,7 @@ import org.apache.ibatis.annotations.Param; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Collection; /** * web-infrastructure @@ -40,5 +42,7 @@ public interface VehicleApplyMapper extends BaseMapper { List getAppCount(@Param("params") Map parmas); List listTodayHasPort(String licenceNo); + + List listAccessMatchSnapshots(@Param("licenceNos") Collection licenceNos); } diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/DaHuaVehicleAccessRecordRepository.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/DaHuaVehicleAccessRecordRepository.java new file mode 100644 index 0000000..cc460ab --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/DaHuaVehicleAccessRecordRepository.java @@ -0,0 +1,15 @@ +package com.zcloud.primeport.persistence.repository; + +import com.jjb.saas.framework.repository.repo.BaseRepository; +import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO; + +import java.time.LocalDateTime; + +public interface DaHuaVehicleAccessRecordRepository extends BaseRepository { + + LocalDateTime findSyncCheckpoint(); + + void saveSyncCheckpoint(LocalDateTime checkpointTime); + + void upsert(DaHuaVehicleAccessRecordDO record); +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/VehicleApplyRepository.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/VehicleApplyRepository.java index 3e243e5..393c68e 100644 --- a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/VehicleApplyRepository.java +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/VehicleApplyRepository.java @@ -5,11 +5,13 @@ import com.jjb.saas.framework.repository.repo.BaseRepository; import com.zcloud.primeport.persistence.dataobject.AppCountDTO; import com.zcloud.primeport.persistence.dataobject.FgsVehicleCountDto; import com.zcloud.primeport.persistence.dataobject.VehicleApplyDO; +import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO; import com.zcloud.primeport.plan.mjDevice.dto.OneLevelCarSaveDto; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Collection; /** * web-infrastructure @@ -32,5 +34,7 @@ public interface VehicleApplyRepository extends BaseRepository { List getAppCount(Map parmas); List listTodayHasPort(String licenceNo); + + List listAccessMatchSnapshots(Collection licenceNos); } diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/impl/DaHuaVehicleAccessRecordRepositoryImpl.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/impl/DaHuaVehicleAccessRecordRepositoryImpl.java new file mode 100644 index 0000000..042809c --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/impl/DaHuaVehicleAccessRecordRepositoryImpl.java @@ -0,0 +1,30 @@ +package com.zcloud.primeport.persistence.repository.impl; + +import com.jjb.saas.framework.repository.repo.impl.BaseRepositoryImpl; +import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO; +import com.zcloud.primeport.persistence.mapper.DaHuaVehicleAccessRecordMapper; +import com.zcloud.primeport.persistence.repository.DaHuaVehicleAccessRecordRepository; +import org.springframework.stereotype.Repository; + +import java.time.LocalDateTime; + +@Repository +public class DaHuaVehicleAccessRecordRepositoryImpl + extends BaseRepositoryImpl + implements DaHuaVehicleAccessRecordRepository { + + @Override + public LocalDateTime findSyncCheckpoint() { + return baseMapper.findSyncCheckpoint(); + } + + @Override + public void saveSyncCheckpoint(LocalDateTime checkpointTime) { + baseMapper.upsertSyncCheckpoint(checkpointTime); + } + + @Override + public void upsert(DaHuaVehicleAccessRecordDO record) { + baseMapper.upsert(record); + } +} diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/impl/VehicleApplyRepositoryImpl.java b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/impl/VehicleApplyRepositoryImpl.java index ca16dba..df17373 100644 --- a/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/impl/VehicleApplyRepositoryImpl.java +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/persistence/repository/impl/VehicleApplyRepositoryImpl.java @@ -16,6 +16,7 @@ import com.zcloud.primeport.domain.enums.VehicleBelongTypeEnum; import com.zcloud.primeport.persistence.dataobject.AppCountDTO; import com.zcloud.primeport.persistence.dataobject.FgsVehicleCountDto; import com.zcloud.primeport.persistence.dataobject.VehicleApplyDO; +import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO; import com.zcloud.primeport.persistence.dataobject.VehicleMessageDO; import com.zcloud.primeport.persistence.mapper.VehicleApplyMapper; import com.zcloud.primeport.persistence.repository.VehicleApplyRepository; @@ -112,5 +113,13 @@ public class VehicleApplyRepositoryImpl extends BaseRepositoryImpl listTodayHasPort(String licenceNo) { return vehicleApplyMapper.listTodayHasPort(licenceNo); } + + @Override + public List listAccessMatchSnapshots(Collection licenceNos) { + if (licenceNos == null || licenceNos.isEmpty()) { + return Collections.emptyList(); + } + return vehicleApplyMapper.listAccessMatchSnapshots(licenceNos); + } } diff --git a/web-infrastructure/src/main/java/com/zcloud/primeport/plan/DaHuaVehicleAccessRecordSyncXxlJob.java b/web-infrastructure/src/main/java/com/zcloud/primeport/plan/DaHuaVehicleAccessRecordSyncXxlJob.java new file mode 100644 index 0000000..1db82ec --- /dev/null +++ b/web-infrastructure/src/main/java/com/zcloud/primeport/plan/DaHuaVehicleAccessRecordSyncXxlJob.java @@ -0,0 +1,43 @@ +package com.zcloud.primeport.plan; + +import com.jjb.saas.framework.job.Job; +import com.jjb.saas.framework.job.annotation.JobRegister; +import com.xxl.job.core.biz.model.ReturnT; +import com.xxl.job.core.handler.annotation.XxlJob; +import com.zcloud.primeport.dahua.config.DaHuaVehicleAccessRecordProperties; +import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordSyncGateway; +import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +@Component +@RequiredArgsConstructor +@Slf4j +public class DaHuaVehicleAccessRecordSyncXxlJob implements Job { + + private final DaHuaVehicleAccessRecordSyncGateway syncGateway; + private final DaHuaVehicleAccessRecordProperties properties; + + @Override + @JobRegister(cron = "0 */5 * * * ?", jobDesc = "大华车辆进出记录增量同步", triggerStatus = 1) + @XxlJob("com.zcloud.plan.DaHuaVehicleAccessRecordSyncXxlJob") + public ReturnT execute(String param) { + if (!properties.isSyncEnabled()) { + log.info("大华车辆进出记录增量同步已关闭"); + return ReturnT.SUCCESS; + } + try { + DaHuaVehicleAccessRecordSyncResultE result = syncGateway.syncIncremental( + properties.getSyncPageSize(), properties.getSyncDelayMinutes(), + properties.getSyncOverlapMinutes(), properties.getInitialLookbackMinutes()); + log.info("大华车辆进出记录增量同步完成: pages={}, total={}, created={}, updated={}, invalid={}", + result.getPageCount(), result.getRecordTotal(), result.getRecordCreated(), + result.getRecordUpdated(), result.getRecordInvalid()); + return ReturnT.SUCCESS; + } catch (Exception e) { + log.error("大华车辆进出记录增量同步失败", e); + return ReturnT.FAIL; + } + } +} diff --git a/web-infrastructure/src/main/resources/TableCreationDDL.sql b/web-infrastructure/src/main/resources/TableCreationDDL.sql index 69c7b98..690f3cd 100644 --- a/web-infrastructure/src/main/resources/TableCreationDDL.sql +++ b/web-infrastructure/src/main/resources/TableCreationDDL.sql @@ -498,3 +498,62 @@ CREATE TABLE `dahua_access_record` ( KEY `idx_dahua_access_record_person_time` (`person_code`, `swing_time`), KEY `idx_dahua_access_record_create_time` (`dahua_create_time`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华门禁通行记录表'; + +CREATE TABLE `dahua_vehicle_access_record` ( + `id` bigint NOT NULL COMMENT '主键ID', + `dahua_access_id` varchar(128) NOT NULL COMMENT '大华车辆进出记录accessId,幂等键', + `corp_id` bigint DEFAULT NULL COMMENT '平台企业ID快照', + `corp_name` varchar(255) DEFAULT NULL COMMENT '平台企业名称快照', + `port_area` int DEFAULT NULL COMMENT '企业港区快照', + `dahua_department_id` bigint DEFAULT NULL COMMENT '大华部门ID', + `dahua_department_name` varchar(255) DEFAULT NULL COMMENT '大华部门名称', + `dahua_org_code` varchar(64) DEFAULT NULL COMMENT '大华组织编码', + `parking_lot_code` varchar(128) DEFAULT NULL COMMENT '停车场编码', + `parking_lot_name` varchar(255) DEFAULT NULL COMMENT '停车场名称', + `car_type` int DEFAULT NULL COMMENT '车辆类型编码', + `car_type_name` varchar(128) DEFAULT NULL COMMENT '车辆类型名称', + `car_num` varchar(64) DEFAULT NULL COMMENT '入场车牌', + `exit_car_num` varchar(64) DEFAULT NULL COMMENT '出场车牌', + `enter_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '入场卡口相机通道编码', + `enter_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '入场道闸通道编码', + `enter_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '入场道闸通道名称', + `enter_time` datetime DEFAULT NULL COMMENT '入场时间', + `enter_image_url` varchar(1024) DEFAULT NULL COMMENT '入场图片', + `exit_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '出场卡口相机通道编码', + `exit_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '出场道闸通道编码', + `exit_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '出场道闸通道名称', + `exit_time` datetime DEFAULT NULL COMMENT '出场时间', + `exit_image_url` varchar(1024) DEFAULT NULL COMMENT '出场图片', + `source_type` varchar(32) NOT NULL DEFAULT 'PULL' COMMENT '数据来源', + `raw_data` json DEFAULT NULL COMMENT '大华原始记录JSON', + `last_sync_time` datetime DEFAULT NULL COMMENT '最近同步时间', + `delete_enum` varchar(32) NOT NULL DEFAULT 'FALSE' COMMENT '删除标识', + `create_time` datetime DEFAULT NULL, + `update_time` datetime DEFAULT NULL, + `create_id` bigint DEFAULT NULL, + `update_id` bigint DEFAULT NULL, + `env` varchar(32) NOT NULL DEFAULT 'PROD', + `create_name` varchar(255) DEFAULT NULL, + `update_name` varchar(255) DEFAULT NULL, + `tenant_id` bigint DEFAULT NULL, + `org_id` bigint DEFAULT NULL, + `version` int NOT NULL DEFAULT '0', + `remarks` varchar(500) DEFAULT NULL, + PRIMARY KEY (`id`), + UNIQUE KEY `uk_dahua_vehicle_access_record_access_id` (`dahua_access_id`), + KEY `idx_dahua_vehicle_access_record_corp_enter` (`corp_id`, `enter_time`), + KEY `idx_dahua_vehicle_access_record_corp_exit` (`corp_id`, `exit_time`), + KEY `idx_dahua_vehicle_access_record_car_enter` (`car_num`, `enter_time`), + KEY `idx_dahua_vehicle_access_record_car_exit` (`exit_car_num`, `exit_time`), + KEY `idx_dahua_vehicle_access_record_department` (`dahua_department_id`), + KEY `idx_dahua_vehicle_access_record_enter_channel` (`enter_sluice_channel_code`), + KEY `idx_dahua_vehicle_access_record_exit_channel` (`exit_sluice_channel_code`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华车辆进出记录表'; + +CREATE TABLE `dahua_sync_checkpoint` ( + `sync_type` varchar(64) NOT NULL COMMENT '同步任务类型', + `checkpoint_time` datetime NOT NULL COMMENT '最近完整成功同步窗口的结束时间', + `create_time` datetime DEFAULT NULL COMMENT '创建时间', + `update_time` datetime DEFAULT NULL COMMENT '更新时间', + PRIMARY KEY (`sync_type`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华同步检查点'; diff --git a/web-infrastructure/src/main/resources/db/migration/V20260806_01__create_dahua_vehicle_access_record.sql b/web-infrastructure/src/main/resources/db/migration/V20260806_01__create_dahua_vehicle_access_record.sql new file mode 100644 index 0000000..18ad08e --- /dev/null +++ b/web-infrastructure/src/main/resources/db/migration/V20260806_01__create_dahua_vehicle_access_record.sql @@ -0,0 +1,58 @@ +CREATE TABLE IF NOT EXISTS `dahua_vehicle_access_record` ( + `id` bigint NOT NULL COMMENT '主键ID', + `dahua_access_id` varchar(128) NOT NULL COMMENT '大华车辆进出记录accessId,幂等键', + `corp_id` bigint DEFAULT NULL COMMENT '平台企业ID快照', + `corp_name` varchar(255) DEFAULT NULL COMMENT '平台企业名称快照', + `port_area` int DEFAULT NULL COMMENT '企业港区快照', + `dahua_department_id` bigint DEFAULT NULL COMMENT '大华部门ID', + `dahua_department_name` varchar(255) DEFAULT NULL COMMENT '大华部门名称', + `dahua_org_code` varchar(64) DEFAULT NULL COMMENT '大华组织编码', + `parking_lot_code` varchar(128) DEFAULT NULL COMMENT '停车场编码', + `parking_lot_name` varchar(255) DEFAULT NULL COMMENT '停车场名称', + `car_type` int DEFAULT NULL COMMENT '车辆类型编码', + `car_type_name` varchar(128) DEFAULT NULL COMMENT '车辆类型名称', + `car_num` varchar(64) DEFAULT NULL COMMENT '入场车牌', + `exit_car_num` varchar(64) DEFAULT NULL COMMENT '出场车牌', + `enter_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '入场卡口相机通道编码', + `enter_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '入场道闸通道编码', + `enter_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '入场道闸通道名称', + `enter_time` datetime DEFAULT NULL COMMENT '入场时间', + `enter_image_url` varchar(1024) DEFAULT NULL COMMENT '入场图片', + `exit_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '出场卡口相机通道编码', + `exit_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '出场道闸通道编码', + `exit_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '出场道闸通道名称', + `exit_time` datetime DEFAULT NULL COMMENT '出场时间', + `exit_image_url` varchar(1024) DEFAULT NULL COMMENT '出场图片', + `source_type` varchar(32) NOT NULL DEFAULT 'PULL' COMMENT '数据来源', + `raw_data` json DEFAULT NULL COMMENT '大华原始记录JSON', + `last_sync_time` datetime DEFAULT NULL COMMENT '最近同步时间', + `delete_enum` varchar(32) NOT NULL DEFAULT 'FALSE' COMMENT '删除标识', + `create_time` datetime DEFAULT NULL, + `update_time` datetime DEFAULT NULL, + `create_id` bigint DEFAULT NULL, + `update_id` bigint DEFAULT NULL, + `env` varchar(32) NOT NULL DEFAULT 'PROD', + `create_name` varchar(255) DEFAULT NULL, + `update_name` varchar(255) DEFAULT NULL, + `tenant_id` bigint DEFAULT NULL, + `org_id` bigint DEFAULT NULL, + `version` int NOT NULL DEFAULT '0', + `remarks` varchar(500) DEFAULT NULL, + PRIMARY KEY (`id`), + UNIQUE KEY `uk_dahua_vehicle_access_record_access_id` (`dahua_access_id`), + KEY `idx_dahua_vehicle_access_record_corp_enter` (`corp_id`, `enter_time`), + KEY `idx_dahua_vehicle_access_record_corp_exit` (`corp_id`, `exit_time`), + KEY `idx_dahua_vehicle_access_record_car_enter` (`car_num`, `enter_time`), + KEY `idx_dahua_vehicle_access_record_car_exit` (`exit_car_num`, `exit_time`), + KEY `idx_dahua_vehicle_access_record_department` (`dahua_department_id`), + KEY `idx_dahua_vehicle_access_record_enter_channel` (`enter_sluice_channel_code`), + KEY `idx_dahua_vehicle_access_record_exit_channel` (`exit_sluice_channel_code`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华车辆进出记录表'; + +CREATE TABLE IF NOT EXISTS `dahua_sync_checkpoint` ( + `sync_type` varchar(64) NOT NULL COMMENT '同步任务类型', + `checkpoint_time` datetime NOT NULL COMMENT '最近完整成功同步窗口的结束时间', + `create_time` datetime DEFAULT NULL COMMENT '创建时间', + `update_time` datetime DEFAULT NULL COMMENT '更新时间', + PRIMARY KEY (`sync_type`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华同步检查点'; diff --git a/web-infrastructure/src/main/resources/mapper/DaHuaVehicleAccessRecordMapper.xml b/web-infrastructure/src/main/resources/mapper/DaHuaVehicleAccessRecordMapper.xml new file mode 100644 index 0000000..a46a1e3 --- /dev/null +++ b/web-infrastructure/src/main/resources/mapper/DaHuaVehicleAccessRecordMapper.xml @@ -0,0 +1,75 @@ + + + + + + + + + INSERT INTO dahua_sync_checkpoint (sync_type, checkpoint_time, create_time, update_time) + VALUES ('VEHICLE_ACCESS_RECORD', #{checkpointTime}, NOW(), NOW()) + ON DUPLICATE KEY UPDATE + checkpoint_time = GREATEST(checkpoint_time, VALUES(checkpoint_time)), + update_time = NOW() + + + + INSERT INTO dahua_vehicle_access_record ( + id, dahua_access_id, corp_id, corp_name, port_area, + dahua_department_id, dahua_department_name, dahua_org_code, + parking_lot_code, parking_lot_name, car_type, car_type_name, + car_num, exit_car_num, + enter_itc_channel_code, enter_sluice_channel_code, enter_sluice_channel_name, + enter_time, enter_image_url, + exit_itc_channel_code, exit_sluice_channel_code, exit_sluice_channel_name, + exit_time, exit_image_url, + source_type, raw_data, last_sync_time, + delete_enum, create_time, update_time, env, version + ) VALUES ( + #{record.id}, #{record.dahuaAccessId}, #{record.corpId}, #{record.corpName}, #{record.portArea}, + #{record.dahuaDepartmentId}, #{record.dahuaDepartmentName}, #{record.dahuaOrgCode}, + #{record.parkingLotCode}, #{record.parkingLotName}, #{record.carType}, #{record.carTypeName}, + #{record.carNum}, #{record.exitCarNum}, + #{record.enterItcChannelCode}, #{record.enterSluiceChannelCode}, #{record.enterSluiceChannelName}, + #{record.enterTime}, #{record.enterImageUrl}, + #{record.exitItcChannelCode}, #{record.exitSluiceChannelCode}, #{record.exitSluiceChannelName}, + #{record.exitTime}, #{record.exitImageUrl}, + #{record.sourceType}, #{record.rawData}, #{record.lastSyncTime}, + #{record.deleteEnum}, #{record.createTime}, #{record.updateTime}, #{record.env}, #{record.version} + ) + ON DUPLICATE KEY UPDATE + corp_id = COALESCE(VALUES(corp_id), corp_id), + corp_name = COALESCE(VALUES(corp_name), corp_name), + port_area = COALESCE(VALUES(port_area), port_area), + dahua_department_id = COALESCE(VALUES(dahua_department_id), dahua_department_id), + dahua_department_name = COALESCE(VALUES(dahua_department_name), dahua_department_name), + dahua_org_code = COALESCE(VALUES(dahua_org_code), dahua_org_code), + parking_lot_code = COALESCE(VALUES(parking_lot_code), parking_lot_code), + parking_lot_name = COALESCE(VALUES(parking_lot_name), parking_lot_name), + car_type = COALESCE(VALUES(car_type), car_type), + car_type_name = COALESCE(VALUES(car_type_name), car_type_name), + car_num = COALESCE(VALUES(car_num), car_num), + exit_car_num = COALESCE(VALUES(exit_car_num), exit_car_num), + enter_itc_channel_code = COALESCE(VALUES(enter_itc_channel_code), enter_itc_channel_code), + enter_sluice_channel_code = COALESCE(VALUES(enter_sluice_channel_code), enter_sluice_channel_code), + enter_sluice_channel_name = COALESCE(VALUES(enter_sluice_channel_name), enter_sluice_channel_name), + enter_time = COALESCE(VALUES(enter_time), enter_time), + enter_image_url = COALESCE(VALUES(enter_image_url), enter_image_url), + exit_itc_channel_code = COALESCE(VALUES(exit_itc_channel_code), exit_itc_channel_code), + exit_sluice_channel_code = COALESCE(VALUES(exit_sluice_channel_code), exit_sluice_channel_code), + exit_sluice_channel_name = COALESCE(VALUES(exit_sluice_channel_name), exit_sluice_channel_name), + exit_time = COALESCE(VALUES(exit_time), exit_time), + exit_image_url = COALESCE(VALUES(exit_image_url), exit_image_url), + source_type = VALUES(source_type), + raw_data = VALUES(raw_data), + last_sync_time = VALUES(last_sync_time), + delete_enum = VALUES(delete_enum), + update_time = VALUES(update_time) + + + diff --git a/web-infrastructure/src/main/resources/mapper/VehicleApplyDO.xml b/web-infrastructure/src/main/resources/mapper/VehicleApplyDO.xml index 7466948..11a0a5d 100644 --- a/web-infrastructure/src/main/resources/mapper/VehicleApplyDO.xml +++ b/web-infrastructure/src/main/resources/mapper/VehicleApplyDO.xml @@ -220,5 +220,20 @@ licence_no, licence_type +