package jnpf.limsService.impl; import jnpf.base.UserInfo; import jnpf.limsEntity.LimsMessageDeliveryEntity; import jnpf.limsEntity.LimsMessageDeliveryStatus; import jnpf.limsEntity.LimsMessageEventEntity; import jnpf.limsEntity.LimsMessageEventStatus; import jnpf.limsEntity.LimsMessageLedger; import jnpf.limsEntity.LimsMessageRequest; import jnpf.limsEntity.LimsMessageSendResult; import jnpf.limsMapper.LimsMessageDeliveryMapper; import jnpf.limsMapper.LimsMessageEventMapper; import jnpf.limsService.LimsMessageService; import jnpf.limsService.LimsRecipientResolver; import jnpf.limsService.client.LimsMessageClient; import jnpf.message.model.SentMessageForm; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; import java.util.Collections; import java.util.Date; import java.util.HashMap; import java.util.List; import java.util.Map; @Slf4j @Service @RequiredArgsConstructor public class LimsMessageServiceImpl implements LimsMessageService { private static final int MAX_ERROR_LENGTH = 1000; private final LimsMessageEventMapper eventMapper; private final LimsMessageDeliveryMapper deliveryMapper; private final LimsRecipientResolver recipientResolver; private final LimsMessageClient messageClient; private final LimsMessageLedgerService ledgerService; private final LimsMessageConfigValidator configValidator; @Override public LimsMessageSendResult send(LimsMessageRequest request) { validate(request); request.setTenantId(normalizeTenant(request.getTenantId())); List recipients = recipientResolver.resolve(request.getTenantId(), request.getRecipients()); LimsMessageConfigValidator.ValidationResult config = validateConfig(request); if (request.isDryRun()) { log.info("[lims-message][dry] event={}, bizId={}, occurrence={}, configValid={}, configError={}, " + "recipients={}, parameters={}", request.getEvent().getCode(), request.getBizId(), request.getOccurrenceKey(), config.isValid(), config.getError(), recipients, request.getParameters()); int sendable = config.isValid() ? recipients.size() : 0; return new LimsMessageSendResult(true, false, recipients.size(), sendable, 0); } boolean ready = config.isValid() && !recipients.isEmpty(); LimsMessageLedger ledger = ledgerService.initialize(request, ready ? recipients : Collections.emptyList()); if (!ledger.isAcquired()) { return LimsMessageSendResult.duplicate(); } LimsMessageEventEntity event = ledger.getEvent(); if (!config.isValid()) { event.setRecipientCount(recipients.size()); failEvent(event, config.getError()); log.error("[lims-message] configuration rejected event={}, bizId={}, occurrence={}, reason={}", request.getEvent().getCode(), request.getBizId(), request.getOccurrenceKey(), config.getError()); return new LimsMessageSendResult(true, false, recipients.size(), 0, 0); } if (recipients.isEmpty()) { failEvent(event, "RECIPIENT_NOT_FOUND: no active user resolved"); log.error("[lims-message] no recipient event={}, bizId={}, occurrence={}, refs={}, tenant={}", request.getEvent().getCode(), request.getBizId(), request.getOccurrenceKey(), request.getRecipients(), request.getTenantId()); return new LimsMessageSendResult(true, false, 0, 0, 0); } List deliveries = ledger.getDeliveries(); event.setStatus(LimsMessageEventStatus.PROCESSING); event.setLastModifyTime(new Date()); eventMapper.updateById(event); int sent = 0; int failed = 0; for (LimsMessageDeliveryEntity delivery : deliveries) { if (deliveryMapper.claim(delivery.getId()) == 0) { failed++; log.warn("[lims-message] delivery claim lost, event={}, delivery={}", request.getEvent().getCode(), delivery.getId()); continue; } try { List errors = messageClient.sendScheduleMessage( buildMessageForm(request, delivery)); if (errors != null && !errors.isEmpty()) { markFailed(delivery, String.join("; ", errors)); failed++; } else { markSent(delivery); sent++; } } catch (Exception e) { markFailed(delivery, e.getMessage()); failed++; log.error("[lims-message] send failed event={}, bizId={}, recipient={}", request.getEvent().getCode(), request.getBizId(), delivery.getRecipientUserId(), e); } } aggregate(event, recipients.size(), sent, failed); return new LimsMessageSendResult(true, false, recipients.size(), sent, failed); } private SentMessageForm buildMessageForm(LimsMessageRequest request, LimsMessageDeliveryEntity delivery) { UserInfo userInfo = new UserInfo(); userInfo.setUserId("system"); userInfo.setUserName("LIMS"); userInfo.setTenantId(normalizeTenant(request.getTenantId())); Map parameters = request.getParameters() == null ? new HashMap<>() : new HashMap<>(request.getParameters()); parameters.put("deliveryKey", delivery.getDeliveryKey()); SentMessageForm form = new SentMessageForm(); form.setToUserIds(Collections.singletonList(delivery.getRecipientUserId())); form.setTemplateId(request.getEvent().getCode()); form.setParameterMap(parameters); form.setContentMsg(new HashMap<>()); form.setUserInfo(userInfo); return form; } private void markSent(LimsMessageDeliveryEntity delivery) { Date now = new Date(); delivery.setStatus(LimsMessageDeliveryStatus.SENT); delivery.setSentTime(now); delivery.setLastModifyTime(now); delivery.setErrorCode(null); delivery.setErrorMessage(null); deliveryMapper.updateById(delivery); } private void markFailed(LimsMessageDeliveryEntity delivery, String error) { delivery.setStatus(LimsMessageDeliveryStatus.FAILED); delivery.setRetryCount(delivery.getRetryCount() == null ? 1 : delivery.getRetryCount() + 1); delivery.setErrorCode("SEND_FAILED"); delivery.setErrorMessage(limit(error)); delivery.setNextRetryTime(new Date(System.currentTimeMillis() + 2 * 60_000L)); delivery.setLastModifyTime(new Date()); deliveryMapper.updateById(delivery); } private void aggregate(LimsMessageEventEntity event, int recipients, int sent, int failed) { Date now = new Date(); event.setRecipientCount(recipients); event.setSentCount(sent); event.setFailedCount(failed); event.setStatus(failed == 0 ? LimsMessageEventStatus.COMPLETED : sent == 0 ? LimsMessageEventStatus.FAILED : LimsMessageEventStatus.PARTIAL_FAILED); event.setLastModifyTime(now); if (failed == 0) { event.setCompletedTime(now); } eventMapper.updateById(event); } private void failEvent(LimsMessageEventEntity event, String error) { event.setStatus(LimsMessageEventStatus.FAILED); event.setErrorMessage(limit(error)); event.setLastModifyTime(new Date()); eventMapper.updateById(event); } private void validate(LimsMessageRequest request) { if (request == null || request.getEvent() == null) { throw new IllegalArgumentException("message event is required"); } if (!StringUtils.hasText(request.getBizId()) || !StringUtils.hasText(request.getOccurrenceKey())) { throw new IllegalArgumentException("bizId and occurrenceKey are required"); } } private LimsMessageConfigValidator.ValidationResult validateConfig(LimsMessageRequest request) { try { return configValidator.validate(request.getEvent().getCode()); } catch (Exception e) { log.error("[lims-message] configuration check failed event={}, bizId={}", request.getEvent().getCode(), request.getBizId(), e); return new LimsMessageConfigValidator.ValidationResult(false, "CONFIG_INVALID: configuration check failed: " + limit(e.getMessage())); } } private String normalizeTenant(String tenantId) { return StringUtils.hasText(tenantId) ? tenantId : "0"; } private String limit(String value) { if (!StringUtils.hasText(value)) { return "Unknown send error"; } return value.length() <= MAX_ERROR_LENGTH ? value : value.substring(0, MAX_ERROR_LENGTH); } }