feat(aihr): complete company message R1-R3

This commit is contained in:
2026-07-24 23:25:15 +08:00
parent 5d864a1444
commit 342f84a9de
37 changed files with 2479 additions and 42 deletions
@@ -0,0 +1,326 @@
package org.dromara.aihr.broadcast;
import jakarta.annotation.PostConstruct;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastAttachmentResponse;
import org.dromara.aihr.knowledge.domain.AihrKnowledgePrincipal;
import org.dromara.aihr.knowledge.parse.KnowledgeDocumentParser;
import org.dromara.aihr.knowledge.service.AihrKnowledgePrincipalResolver;
import org.dromara.aihr.service.AihrMultipartFiles;
import org.dromara.aihr.service.AihrSopSeedService;
import org.dromara.common.core.constant.HttpStatus;
import org.dromara.common.core.constant.TenantConstants;
import org.dromara.common.core.enums.UserType;
import org.dromara.common.core.exception.ServiceException;
import org.dromara.common.oss.core.OssClient;
import org.dromara.common.oss.factory.OssFactory;
import org.dromara.common.tenant.helper.TenantHelper;
import org.dromara.system.domain.vo.SysOssVo;
import org.dromara.system.service.ISysOssService;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.support.GeneratedKeyHolder;
import org.springframework.jdbc.support.KeyHolder;
import org.springframework.stereotype.Service;
import org.springframework.web.multipart.MultipartFile;
import java.io.InputStream;
import java.sql.PreparedStatement;
import java.sql.Statement;
import java.util.List;
import java.util.Locale;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@Service
@RequiredArgsConstructor
@Slf4j
public class AihrBroadcastAttachmentService {
private static final long MAX_FILE_BYTES = 100L * 1024 * 1024;
private static final Set<String> SUPPORTED_SUFFIXES = Set.of(
"txt", "md", "markdown", "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx"
);
private final JdbcTemplate jdbcTemplate;
private final ISysOssService ossService;
private final KnowledgeDocumentParser documentParser;
private final AihrSopSeedService sopService;
private final AihrKnowledgePrincipalResolver principalResolver;
private final AihrBroadcastService broadcastService;
private final ScheduledExecutorService scheduledExecutorService;
private final AtomicBoolean workerRunning = new AtomicBoolean();
@PostConstruct
void scheduleRecovery() {
scheduledExecutorService.scheduleWithFixedDelay(this::recoverAndTrigger, 20, 30, TimeUnit.SECONDS);
}
public BroadcastAttachmentResponse upload(MultipartFile file, String rawExpectedTenantId) {
AihrKnowledgePrincipal principal = principalResolver.current();
String tenantId = adminTenant(principal);
requireExpectedTenant(rawExpectedTenantId, tenantId);
if (file == null || file.isEmpty()) {
throw new ServiceException("请选择公司文件", HttpStatus.BAD_REQUEST);
}
if (file.getSize() <= 0 || file.getSize() > MAX_FILE_BYTES) {
throw new ServiceException("公司文件不能超过 100MB", HttpStatus.BAD_REQUEST);
}
String fileName = AihrMultipartFiles.sanitizeFileName(file.getOriginalFilename(), "company-file.txt");
if (!SUPPORTED_SUFFIXES.contains(suffix(fileName))) {
throw new ServiceException("仅支持 TXT、Markdown、PDF、Word、Excel 和 PPT 文件", HttpStatus.BAD_REQUEST);
}
MultipartFile sanitized = AihrMultipartFiles.withOriginalFilename(file, fileName);
SysOssVo oss = ossService.upload(sanitized);
KeyHolder keyHolder = new GeneratedKeyHolder();
try {
jdbcTemplate.update(connection -> {
PreparedStatement statement = connection.prepareStatement("""
insert into aihr_broadcast_attachment
(tenant_id, oss_id, file_name, file_size, content_type, status, uploaded_by, create_time, update_time)
values (?, ?, ?, ?, ?, 'QUEUED', ?, now(), now())
""", Statement.RETURN_GENERATED_KEYS);
statement.setString(1, tenantId);
statement.setLong(2, oss.getOssId());
statement.setString(3, fileName);
statement.setLong(4, file.getSize());
statement.setString(5, contentType(file.getContentType()));
statement.setLong(6, principal.userId());
return statement;
}, keyHolder);
} catch (RuntimeException error) {
deleteOssQuietly(oss.getOssId());
throw error;
}
Number key = keyHolder.getKey();
if (key == null) {
deleteOssQuietly(oss.getOssId());
throw new ServiceException("公司文件入队失败", 500);
}
trigger();
return adminAttachment(key.longValue());
}
public BroadcastAttachmentResponse adminAttachment(Long rawId) {
AihrKnowledgePrincipal principal = principalResolver.current();
String tenantId = adminTenant(principal);
long id = attachmentId(rawId);
List<BroadcastAttachmentResponse> rows = jdbcTemplate.query("""
select id, file_name, file_size, content_type, status, summary, error_message
from aihr_broadcast_attachment
where tenant_id = ? and id = ?
""", (rs, rowNum) -> response(
rs.getLong("id"), rs.getString("file_name"), rs.getLong("file_size"),
rs.getString("content_type"), rs.getString("status"), rs.getString("summary"),
rs.getString("error_message")
), tenantId, id);
if (rows.isEmpty()) {
throw new ServiceException("公司文件不存在或无权查看", HttpStatus.NOT_FOUND);
}
return rows.get(0);
}
public Long authorizedDownloadOssId(Long rawAttachmentId) {
long attachmentId = attachmentId(rawAttachmentId);
AihrKnowledgePrincipal principal = principalResolver.current();
if (UserType.APP_USER.getUserType().equals(principal.userType())) {
List<Long> messageIds = jdbcTemplate.query("""
select message_id from aihr_broadcast_attachment
where tenant_id = ? and id = ? and status = 'READY' and message_id is not null
""", (rs, rowNum) -> rs.getLong("message_id"), principal.tenantId(), attachmentId);
if (messageIds.isEmpty()) {
return null;
}
broadcastService.message(messageIds.get(0));
return ossId(principal.tenantId(), attachmentId);
}
requireBroadcastAdmin(principal);
return ossId(adminTenant(principal), attachmentId);
}
private Long ossId(String tenantId, long attachmentId) {
List<Long> rows = jdbcTemplate.query("""
select oss_id from aihr_broadcast_attachment
where tenant_id = ? and id = ? and status = 'READY'
""", (rs, rowNum) -> rs.getLong("oss_id"), tenantId, attachmentId);
return rows.isEmpty() ? null : rows.get(0);
}
private void recoverAndTrigger() {
try {
jdbcTemplate.update("""
update aihr_broadcast_attachment
set status = 'QUEUED', error_message = null, update_time = now()
where status = 'PROCESSING' and update_time < date_sub(now(), interval 30 minute)
""");
trigger();
} catch (RuntimeException error) {
log.warn("broadcast attachment recovery failed");
}
}
private void trigger() {
if (!workerRunning.compareAndSet(false, true)) {
return;
}
CompletableFuture.runAsync(this::processQueue, scheduledExecutorService)
.whenComplete((ignored, error) -> {
workerRunning.set(false);
if (error != null) {
log.warn("broadcast attachment worker failed");
}
if (hasQueued()) {
trigger();
}
});
}
private void processQueue() {
AttachmentWork work;
while ((work = claim()) != null) {
process(work);
}
}
private AttachmentWork claim() {
List<AttachmentWork> rows = jdbcTemplate.query("""
select id, tenant_id, oss_id, file_name, content_type, uploaded_by
from aihr_broadcast_attachment
where status = 'QUEUED'
order by id
limit 10
""", (rs, rowNum) -> new AttachmentWork(
rs.getLong("id"), rs.getString("tenant_id"), rs.getLong("oss_id"),
rs.getString("file_name"), rs.getString("content_type"), rs.getLong("uploaded_by")
));
for (AttachmentWork row : rows) {
int claimed = jdbcTemplate.update("""
update aihr_broadcast_attachment set status = 'PROCESSING', update_time = now()
where tenant_id = ? and id = ? and status = 'QUEUED'
""", row.tenantId(), row.id());
if (claimed == 1) {
return row;
}
}
return null;
}
private void process(AttachmentWork work) {
try {
SysOssVo object = TenantHelper.dynamic(work.tenantId(), () -> ossService.getById(work.ossId()));
Integer owned = jdbcTemplate.queryForObject("""
select count(*) from sys_oss
where binary tenant_id = binary ? and oss_id = ? and create_by = ?
""", Integer.class, work.tenantId(), work.ossId(), work.uploadedBy());
if (object == null || object.getFileName() == null || object.getFileName().isBlank()
|| owned == null || owned != 1) {
throw new IllegalStateException("stored company file is unavailable");
}
OssClient storage = OssFactory.instance(object.getService());
String text;
try (InputStream input = storage.getObjectContent(object.getFileName())) {
text = documentParser.parse(work.fileName(), work.contentType(), input).text();
}
String summary = TenantHelper.dynamic(work.tenantId(),
() -> sopService.summarizeBroadcastAttachment(work.fileName(), text));
int updated = jdbcTemplate.update("""
update aihr_broadcast_attachment
set status = 'READY', extracted_text = ?, summary = ?, error_message = null, update_time = now()
where tenant_id = ? and id = ? and status = 'PROCESSING'
""", text, summary, work.tenantId(), work.id());
if (updated != 1) {
log.info("broadcast attachment {} completed after its state changed", work.id());
}
} catch (Exception error) {
log.warn("broadcast attachment {} processing failed", work.id());
jdbcTemplate.update("""
update aihr_broadcast_attachment
set status = 'FAILED', error_message = '文件提炼失败,请重新上传', update_time = now()
where tenant_id = ? and id = ? and status = 'PROCESSING'
""", work.tenantId(), work.id());
}
}
private boolean hasQueued() {
Long count = jdbcTemplate.queryForObject("""
select count(*) from aihr_broadcast_attachment where status = 'QUEUED'
""", Long.class);
return count != null && count > 0;
}
private void deleteOssQuietly(Long ossId) {
if (ossId == null) return;
try {
ossService.deleteWithValidByIds(List.of(ossId), false);
} catch (RuntimeException cleanupError) {
log.warn("broadcast attachment OSS cleanup failed ossId={}", ossId);
}
}
private static BroadcastAttachmentResponse response(
long id,
String fileName,
long fileSize,
String contentType,
String status,
String summary,
String errorMessage
) {
return new BroadcastAttachmentResponse(id, fileName, fileSize, contentType, status, summary,
errorMessage, "READY".equals(status) ? "/api/aihr/broadcast/attachments/" + id + "/content" : null);
}
private static void requireBroadcastAdmin(AihrKnowledgePrincipal principal) {
if (UserType.APP_USER.getUserType().equals(principal.userType())
|| (!principal.roles().contains(TenantConstants.SUPER_ADMIN_ROLE_KEY)
&& !principal.roles().contains("hr_operator"))) {
throw new ServiceException("无权管理公司文件", HttpStatus.FORBIDDEN);
}
}
private static String adminTenant(AihrKnowledgePrincipal principal) {
requireBroadcastAdmin(principal);
String tenantId = TenantHelper.getTenantId();
return tenantId == null || tenantId.isBlank() ? principal.tenantId() : tenantId.trim();
}
private static void requireExpectedTenant(String rawExpectedTenantId, String tenantId) {
String expected = rawExpectedTenantId == null ? "" : rawExpectedTenantId.trim();
if (expected.isEmpty() || expected.length() > 64) {
throw new ServiceException("当前租户不能为空", HttpStatus.BAD_REQUEST);
}
if (!expected.equals(tenantId)) {
throw new ServiceException("当前租户已切换,请刷新页面后重试", HttpStatus.CONFLICT);
}
}
private static long attachmentId(Long value) {
if (value == null || value < 1) {
throw new ServiceException("公司文件编号无效", HttpStatus.BAD_REQUEST);
}
return value;
}
private static String suffix(String fileName) {
int dot = fileName.lastIndexOf('.');
return dot < 0 ? "" : fileName.substring(dot + 1).toLowerCase(Locale.ROOT);
}
private static String contentType(String value) {
String normalized = value == null ? "" : value.trim();
return normalized.isEmpty() ? "application/octet-stream" : normalized.substring(0, Math.min(100, normalized.length()));
}
private record AttachmentWork(
long id,
String tenantId,
long ossId,
String fileName,
String contentType,
long uploadedBy
) {
}
}
@@ -4,11 +4,13 @@ import cn.dev33.satoken.annotation.SaCheckLogin;
import cn.dev33.satoken.annotation.SaCheckRole;
import cn.dev33.satoken.annotation.SaMode;
import jakarta.validation.Valid;
import jakarta.servlet.http.HttpServletResponse;
import lombok.RequiredArgsConstructor;
import org.dromara.aihr.broadcast.AihrBroadcastDto.AdminBroadcastListResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastDetailResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastListResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastTargetOptions;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastAttachmentResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.PublishRequest;
import org.dromara.aihr.broadcast.AihrBroadcastDto.PublishResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.UnreadCountResponse;
@@ -21,7 +23,13 @@ import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RequestPart;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.multipart.MultipartFile;
import org.springframework.http.MediaType;
import org.dromara.system.service.ISysOssService;
import java.io.IOException;
@RestController
@RequiredArgsConstructor
@@ -32,6 +40,8 @@ public class AihrBroadcastController {
private static final String HR_OPERATOR_ROLE = "hr_operator";
private final AihrBroadcastService broadcastService;
private final AihrBroadcastAttachmentService attachmentService;
private final ISysOssService ossService;
@GetMapping("/unread-count")
public R<UnreadCountResponse> unreadCount() {
@@ -63,6 +73,21 @@ public class AihrBroadcastController {
return R.ok(broadcastService.targetOptions());
}
@SaCheckRole(value = {TenantConstants.SUPER_ADMIN_ROLE_KEY, HR_OPERATOR_ROLE}, mode = SaMode.OR)
@PostMapping(value = "/admin/attachments", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
public R<BroadcastAttachmentResponse> uploadAttachment(
@RequestPart("file") MultipartFile file,
@RequestParam String expectedTenantId
) {
return R.ok(attachmentService.upload(file, expectedTenantId));
}
@SaCheckRole(value = {TenantConstants.SUPER_ADMIN_ROLE_KEY, HR_OPERATOR_ROLE}, mode = SaMode.OR)
@GetMapping("/admin/attachments/{id}")
public R<BroadcastAttachmentResponse> attachmentStatus(@PathVariable Long id) {
return R.ok(attachmentService.adminAttachment(id));
}
@GetMapping("/messages/{id}")
public R<BroadcastDetailResponse> message(@PathVariable Long id) {
return R.ok(broadcastService.message(id));
@@ -74,6 +99,16 @@ public class AihrBroadcastController {
return R.ok();
}
@GetMapping("/attachments/{id}/content")
public void attachmentContent(@PathVariable Long id, HttpServletResponse response) throws IOException {
Long ossId = attachmentService.authorizedDownloadOssId(id);
if (ossId == null) {
response.sendError(HttpServletResponse.SC_NOT_FOUND, "公司文件不存在或无权访问");
return;
}
ossService.download(ossId, response);
}
@SaCheckRole(value = {TenantConstants.SUPER_ADMIN_ROLE_KEY, HR_OPERATOR_ROLE}, mode = SaMode.OR)
@PostMapping("/messages")
public R<PublishResponse> publish(@Valid @RequestBody PublishRequest request) {
@@ -45,13 +45,26 @@ public final class AihrBroadcastDto {
String publishedAt,
boolean requiredRead,
boolean targeted,
String targetReason
String targetReason,
BroadcastAttachmentResponse attachment
) {
public BroadcastDetailResponse(Long id, String title, String content, boolean read, String publishedAt) {
this(id, title, content, read, publishedAt, false, false, null);
this(id, title, content, read, publishedAt, false, false, null, null);
}
}
public record BroadcastAttachmentResponse(
Long id,
String fileName,
long fileSize,
String contentType,
String status,
String summary,
String errorMessage,
String downloadUrl
) {
}
public record AdminBroadcastListResponse(
String tenantId,
long total,
@@ -138,10 +151,22 @@ public final class AihrBroadcastDto {
@Size(max = 64, message = "当前租户不能超过 64 个字符")
String expectedTenantId,
Boolean requiredRead,
@Valid BroadcastTargetRequest targets
@Valid BroadcastTargetRequest targets,
Long attachmentId
) {
public PublishRequest(String requestId, String title, String content, String expectedTenantId) {
this(requestId, title, content, expectedTenantId, false, BroadcastTargetRequest.empty());
this(requestId, title, content, expectedTenantId, false, BroadcastTargetRequest.empty(), null);
}
public PublishRequest(
String requestId,
String title,
String content,
String expectedTenantId,
Boolean requiredRead,
BroadcastTargetRequest targets
) {
this(requestId, title, content, expectedTenantId, requiredRead, targets, null);
}
}
@@ -4,6 +4,7 @@ import lombok.RequiredArgsConstructor;
import org.dromara.aihr.broadcast.AihrBroadcastDto.AdminBroadcastListItem;
import org.dromara.aihr.broadcast.AihrBroadcastDto.AdminBroadcastListResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastDetailResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastAttachmentResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastListItem;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastListResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastTargetOption;
@@ -207,7 +208,11 @@ public class AihrBroadcastService {
case when r.id is null then 0 else 1 end as read_flag,
case when m.required_read = 1 then 1 else 0 end as required_read,
case when tr.user_id is null then 0 else 1 end as targeted_flag,
tr.target_reason
tr.target_reason,
a.id as attachment_id, a.file_name as attachment_file_name,
a.file_size as attachment_file_size, a.content_type as attachment_content_type,
a.status as attachment_status, a.summary as attachment_summary,
a.error_message as attachment_error_message
from aihr_broadcast_message m
left join aihr_broadcast_read r
on r.tenant_id = m.tenant_id and r.message_id = m.id and r.user_id = ?
@@ -219,12 +224,14 @@ public class AihrBroadcastService {
and recipient.target_status = 'MATCHED'
group by recipient.tenant_id, recipient.message_id, recipient.user_id
) tr on tr.tenant_id = m.tenant_id and tr.message_id = m.id
left join aihr_broadcast_attachment a
on a.tenant_id = m.tenant_id and a.id = m.attachment_id and a.status = 'READY'
where m.tenant_id = ? and m.id = ? and m.status = 'PUBLISHED'
limit 1
""", (rs, rowNum) -> new BroadcastDetailResponse(
rs.getLong("id"), rs.getString("title"), rs.getString("content"), rs.getBoolean("read_flag"),
format(rs.getTimestamp("published_time")), rs.getBoolean("required_read"),
rs.getBoolean("targeted_flag"), rs.getString("target_reason")
rs.getBoolean("targeted_flag"), rs.getString("target_reason"), attachment(rs)
), principal.userId(), principal.tenantId(), principal.userId(), principal.tenantId(), messageId);
if (rows.isEmpty()) {
throw unavailable();
@@ -245,12 +252,16 @@ public class AihrBroadcastService {
requireEmployeeAudience(principal);
long messageId = messageId(rawMessageId);
List<BroadcastQuestionContext> rows = jdbcTemplate.query("""
select m.id, m.title, m.content, m.published_time
select m.id, m.title, m.content, m.published_time, a.extracted_text
from aihr_broadcast_message m
left join aihr_broadcast_attachment a
on a.tenant_id = m.tenant_id and a.id = m.attachment_id and a.status = 'READY'
where m.tenant_id = ? and m.id = ? and m.status = 'PUBLISHED'
limit 1
""", (rs, rowNum) -> new BroadcastQuestionContext(
rs.getLong("id"), rs.getString("title"), rs.getString("content"), format(rs.getTimestamp("published_time"))
rs.getLong("id"), rs.getString("title"), questionContent(
rs.getString("content"), rs.getString("extracted_text")),
format(rs.getTimestamp("published_time"))
), principal.tenantId(), messageId);
if (rows.isEmpty()) {
throw unavailable();
@@ -289,29 +300,33 @@ public class AihrBroadcastService {
String content = required(request == null ? null : request.content(), "消息正文", 10_000);
boolean requiredRead = request != null && Boolean.TRUE.equals(request.requiredRead());
TargetSelection targetSelection = targetSelection(request == null ? null : request.targets());
Long attachmentId = optionalAttachmentId(request == null ? null : request.attachmentId());
String requestHash = requestHash(title, content);
String targetPayloadHash = targetPayloadHash(requiredRead, targetSelection);
String targetPayloadHash = targetPayloadHash(requiredRead, targetSelection, attachmentId);
PublishReplay existing = findPublishReplay(tenantId, principal.userId(), requestKey);
if (existing != null) {
return replayOrConflict(existing, title, content, requiredRead, targetSelection, targetPayloadHash);
}
requirePublishableAttachment(tenantId, principal.userId(), attachmentId);
KeyHolder keyHolder = new GeneratedKeyHolder();
try {
jdbcTemplate.update(connection -> {
PreparedStatement statement = connection.prepareStatement("""
insert into aihr_broadcast_message
(tenant_id, title, content, required_read, status, published_by, publish_request_key, publish_request_hash,
(tenant_id, title, content, attachment_id, required_read, status, published_by, publish_request_key, publish_request_hash,
target_payload_hash, published_time, create_time, update_time)
values (?, ?, ?, ?, 'PUBLISHED', ?, ?, ?, ?, now(), now(), now())
values (?, ?, ?, ?, ?, 'PUBLISHED', ?, ?, ?, ?, now(), now(), now())
""", Statement.RETURN_GENERATED_KEYS);
statement.setString(1, tenantId);
statement.setString(2, title);
statement.setString(3, content);
statement.setBoolean(4, requiredRead);
statement.setLong(5, principal.userId());
statement.setString(6, requestKey);
statement.setString(7, requestHash);
statement.setString(8, targetPayloadHash);
if (attachmentId == null) statement.setNull(4, java.sql.Types.BIGINT);
else statement.setLong(4, attachmentId);
statement.setBoolean(5, requiredRead);
statement.setLong(6, principal.userId());
statement.setString(7, requestKey);
statement.setString(8, requestHash);
statement.setString(9, targetPayloadHash);
return statement;
}, keyHolder);
} catch (DuplicateKeyException duplicate) {
@@ -338,6 +353,7 @@ public class AihrBroadcastService {
if (snapshots != 1) {
throw new ServiceException("保存消息版本失败");
}
bindAttachment(tenantId, principal.userId(), id.longValue(), attachmentId);
TargetAudit targetAudit = persistTargetSnapshot(tenantId, id.longValue(), targetSelection);
return new PublishResponse(id.longValue(), 1, targetAudit.matchedRecipientCount(), targetAudit.unmatchedTargetCount());
}
@@ -599,7 +615,7 @@ public class AihrBroadcastService {
return List.copyOf(sorted);
}
private static String targetPayloadHash(boolean requiredRead, TargetSelection selection) {
private static String targetPayloadHash(boolean requiredRead, TargetSelection selection, Long attachmentId) {
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
updateDigestField(digest, requiredRead ? "1" : "0");
@@ -609,6 +625,10 @@ public class AihrBroadcastService {
selection.positionNames().forEach(value -> updateDigestField(digest, value));
updateDigestField(digest, "LEVEL");
selection.positionLevels().forEach(value -> updateDigestField(digest, value));
if (attachmentId != null) {
updateDigestField(digest, "ATTACHMENT");
updateDigestField(digest, attachmentId.toString());
}
return HexFormat.of().formatHex(digest.digest());
} catch (NoSuchAlgorithmException impossible) {
throw new IllegalStateException(impossible);
@@ -658,6 +678,56 @@ public class AihrBroadcastService {
return value;
}
private static Long optionalAttachmentId(Long value) {
if (value == null) return null;
if (value < 1) {
throw new ServiceException("公司文件编号无效", HttpStatus.BAD_REQUEST);
}
return value;
}
private void requirePublishableAttachment(String tenantId, long userId, Long attachmentId) {
if (attachmentId == null) return;
Long count = jdbcTemplate.queryForObject("""
select count(*) from aihr_broadcast_attachment
where tenant_id = ? and id = ? and uploaded_by = ?
and status = 'READY' and message_id is null
""", Long.class, tenantId, attachmentId, userId);
if (count == null || count != 1) {
throw new ServiceException("公司文件尚未提炼完成、已被使用或不属于当前账号", HttpStatus.CONFLICT);
}
}
private void bindAttachment(String tenantId, long userId, long messageId, Long attachmentId) {
if (attachmentId == null) return;
int updated = jdbcTemplate.update("""
update aihr_broadcast_attachment set message_id = ?, update_time = now()
where tenant_id = ? and id = ? and uploaded_by = ?
and status = 'READY' and message_id is null
""", messageId, tenantId, attachmentId, userId);
if (updated != 1) {
throw new ServiceException("公司文件绑定失败,请刷新后重试", HttpStatus.CONFLICT);
}
}
private static BroadcastAttachmentResponse attachment(java.sql.ResultSet rs) throws java.sql.SQLException {
Long id = rs.getObject("attachment_id", Long.class);
if (id == null) return null;
return new BroadcastAttachmentResponse(
id, rs.getString("attachment_file_name"), rs.getLong("attachment_file_size"),
rs.getString("attachment_content_type"), rs.getString("attachment_status"),
rs.getString("attachment_summary"), rs.getString("attachment_error_message"),
"/api/aihr/broadcast/attachments/" + id + "/content"
);
}
private static String questionContent(String messageContent, String attachmentText) {
if (attachmentText == null || attachmentText.isBlank()) {
return messageContent;
}
return messageContent + "\n\n公司文件提炼内容:\n" + attachmentText;
}
private static String adminStatus(String rawStatus) {
String status = rawStatus == null ? "" : rawStatus.trim().toUpperCase(Locale.ROOT);
if (status.isEmpty() || "ALL".equals(status)) {
@@ -0,0 +1,69 @@
package org.dromara.aihr.direct;
import cn.dev33.satoken.annotation.SaCheckLogin;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import org.dromara.aihr.direct.AihrDirectDto.AdminFeedbackItem;
import org.dromara.aihr.direct.AihrDirectDto.AdminFeedbackPage;
import org.dromara.aihr.direct.AihrDirectDto.ChannelItem;
import org.dromara.aihr.direct.AihrDirectDto.FeedbackItem;
import org.dromara.aihr.direct.AihrDirectDto.FeedbackPage;
import org.dromara.aihr.direct.AihrDirectDto.ReplyRequest;
import org.dromara.aihr.direct.AihrDirectDto.SubmitRequest;
import org.dromara.common.core.domain.R;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
@RestController
@RequiredArgsConstructor
@SaCheckLogin
@RequestMapping("/api/aihr/direct")
public class AihrDirectController {
private final AihrDirectService directService;
@GetMapping("/channels")
public R<List<ChannelItem>> channels() {
return R.ok(directService.channels());
}
@PostMapping("/feedback")
public R<FeedbackItem> submit(@Valid @RequestBody SubmitRequest request) {
return R.ok(directService.submit(request));
}
@GetMapping("/mine")
public R<FeedbackPage> mine(
@RequestParam(required = false) Integer pageNum,
@RequestParam(required = false) Integer pageSize
) {
return R.ok(directService.mine(pageNum, pageSize));
}
@GetMapping("/mine/{id}")
public R<FeedbackItem> mineDetail(@PathVariable Long id) {
return R.ok(directService.mineDetail(id));
}
@GetMapping("/admin/feedback")
public R<AdminFeedbackPage> adminFeedback(
@RequestParam(required = false) Integer pageNum,
@RequestParam(required = false) Integer pageSize,
@RequestParam(required = false) String status,
@RequestParam(required = false) String channelCode
) {
return R.ok(directService.adminFeedback(pageNum, pageSize, status, channelCode));
}
@PostMapping("/admin/feedback/{id}/reply")
public R<AdminFeedbackItem> reply(@PathVariable Long id, @Valid @RequestBody ReplyRequest request) {
return R.ok(directService.reply(id, request));
}
}
@@ -0,0 +1,85 @@
package org.dromara.aihr.direct;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.Pattern;
import jakarta.validation.constraints.Size;
import java.util.List;
public final class AihrDirectDto {
private AihrDirectDto() {
}
public record ChannelItem(String code, String name, boolean defaultAnonymous) {
}
public record SubmitRequest(
@NotBlank(message = "请求ID不能为空")
@Size(max = 100, message = "请求ID不能超过 100 个字符")
@Pattern(regexp = "[A-Za-z0-9._:-]+", message = "请求ID格式不正确")
String requestId,
@NotBlank(message = "请选择直达对象")
@Size(max = 20, message = "直达对象格式不正确")
String channelCode,
Boolean anonymous,
@NotBlank(message = "反馈内容不能为空")
@Size(max = 5000, message = "反馈内容不能超过 5000 个字符")
String content
) {
}
public record ReplyRequest(
@NotBlank(message = "回复内容不能为空")
@Size(max = 5000, message = "回复内容不能超过 5000 个字符")
String content,
@NotBlank(message = "当前租户不能为空")
@Size(max = 64, message = "当前租户不能超过 64 个字符")
String expectedTenantId
) {
}
public record FeedbackItem(
Long id,
String channelCode,
String channelName,
boolean anonymous,
String content,
String status,
String replyContent,
String repliedAt,
String createdAt
) {
}
public record FeedbackPage(
long total,
int pageNum,
int pageSize,
List<FeedbackItem> rows
) {
}
public record AdminFeedbackItem(
Long id,
String channelCode,
String channelName,
String senderDisplayName,
boolean anonymous,
String content,
String status,
String replyContent,
String repliedAt,
String createdAt
) {
}
public record AdminFeedbackPage(
String tenantId,
long total,
int pageNum,
int pageSize,
List<AdminFeedbackItem> rows
) {
}
}
@@ -0,0 +1,444 @@
package org.dromara.aihr.direct;
import lombok.RequiredArgsConstructor;
import org.dromara.aihr.direct.AihrDirectDto.AdminFeedbackItem;
import org.dromara.aihr.direct.AihrDirectDto.AdminFeedbackPage;
import org.dromara.aihr.direct.AihrDirectDto.ChannelItem;
import org.dromara.aihr.direct.AihrDirectDto.FeedbackItem;
import org.dromara.aihr.direct.AihrDirectDto.FeedbackPage;
import org.dromara.aihr.direct.AihrDirectDto.ReplyRequest;
import org.dromara.aihr.direct.AihrDirectDto.SubmitRequest;
import org.dromara.aihr.knowledge.domain.AihrKnowledgePrincipal;
import org.dromara.aihr.knowledge.service.AihrKnowledgePrincipalResolver;
import org.dromara.common.core.constant.HttpStatus;
import org.dromara.common.core.constant.TenantConstants;
import org.dromara.common.core.enums.UserType;
import org.dromara.common.core.exception.ServiceException;
import org.dromara.common.tenant.helper.TenantHelper;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.support.GeneratedKeyHolder;
import org.springframework.jdbc.support.KeyHolder;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.sql.PreparedStatement;
import java.sql.Statement;
import java.sql.Timestamp;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HexFormat;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Locale;
import java.util.Set;
import java.util.regex.Pattern;
@Service
@RequiredArgsConstructor
public class AihrDirectService {
private static final int DEFAULT_PAGE_SIZE = 20;
private static final int MAX_PAGE_SIZE = 100;
private static final Pattern REQUEST_ID_PATTERN = Pattern.compile("[A-Za-z0-9._:-]+");
private static final List<DirectChannel> CHANNELS = List.of(
new DirectChannel("PRESIDENT", "总裁直达", "direct_president", true),
new DirectChannel("FINANCE", "财务直达", "direct_finance", false),
new DirectChannel("HR", "人力直达", "direct_hr", false),
new DirectChannel("AUDIT", "审计直达", "direct_audit", true),
new DirectChannel("OPERATIONS", "运营直达", "direct_operations", false)
);
private final JdbcTemplate jdbcTemplate;
private final AihrKnowledgePrincipalResolver principalResolver;
public List<ChannelItem> channels() {
requireEmployee(principalResolver.current());
return CHANNELS.stream()
.map(channel -> new ChannelItem(channel.code(), channel.name(), channel.defaultAnonymous()))
.toList();
}
@Transactional
public FeedbackItem submit(SubmitRequest request) {
AihrKnowledgePrincipal principal = principalResolver.current();
String senderName = requireEmployee(principal);
DirectChannel channel = channel(request.channelCode());
String requestKey = requestKey(request.requestId());
String content = required(request.content(), "反馈内容", 5000);
boolean anonymous = request.anonymous() == null ? channel.defaultAnonymous() : request.anonymous();
String requestHash = requestHash(channel.code(), anonymous, content);
FeedbackItem replay = replay(principal, requestKey, requestHash);
if (replay != null) {
return replay;
}
KeyHolder keyHolder = new GeneratedKeyHolder();
try {
jdbcTemplate.update(connection -> {
PreparedStatement statement = connection.prepareStatement("""
insert into aihr_direct_feedback
(tenant_id, channel_code, sender_user_id, sender_name, anonymous_flag, content,
status, submit_request_key, submit_request_hash, create_time, update_time)
values (?, ?, ?, ?, ?, ?, 'SUBMITTED', ?, ?, now(), now())
""", Statement.RETURN_GENERATED_KEYS);
statement.setString(1, principal.tenantId());
statement.setString(2, channel.code());
statement.setLong(3, principal.userId());
statement.setString(4, senderName);
statement.setBoolean(5, anonymous);
statement.setString(6, content);
statement.setString(7, requestKey);
statement.setString(8, requestHash);
return statement;
}, keyHolder);
} catch (DuplicateKeyException concurrentReplay) {
FeedbackItem completed = replay(principal, requestKey, requestHash);
if (completed != null) {
return completed;
}
throw concurrentReplay;
}
Number key = keyHolder.getKey();
if (key == null) {
throw new ServiceException("提交反馈失败", 500);
}
return mineDetail(principal, key.longValue());
}
public FeedbackPage mine(Integer rawPageNum, Integer rawPageSize) {
AihrKnowledgePrincipal principal = principalResolver.current();
requireEmployee(principal);
Page page = page(rawPageNum, rawPageSize);
Long total = jdbcTemplate.queryForObject("""
select count(*) from aihr_direct_feedback
where tenant_id = ? and sender_user_id = ?
""", Long.class, principal.tenantId(), principal.userId());
List<FeedbackItem> rows = jdbcTemplate.query("""
select id, channel_code, anonymous_flag, content, status, reply_content, replied_time, create_time
from aihr_direct_feedback
where tenant_id = ? and sender_user_id = ?
order by create_time desc, id desc
limit ? offset ?
""", (rs, rowNum) -> feedbackItem(
rs.getLong("id"), rs.getString("channel_code"), rs.getBoolean("anonymous_flag"),
rs.getString("content"), rs.getString("status"), rs.getString("reply_content"),
rs.getTimestamp("replied_time"), rs.getTimestamp("create_time")
), principal.tenantId(), principal.userId(), page.pageSize(), page.offset());
return new FeedbackPage(total == null ? 0 : total, page.pageNum(), page.pageSize(), rows);
}
public FeedbackItem mineDetail(Long rawId) {
AihrKnowledgePrincipal principal = principalResolver.current();
requireEmployee(principal);
return mineDetail(principal, id(rawId));
}
public AdminFeedbackPage adminFeedback(
Integer rawPageNum,
Integer rawPageSize,
String rawStatus,
String rawChannelCode
) {
AihrKnowledgePrincipal principal = principalResolver.current();
String tenantId = adminTenant(principal);
List<DirectChannel> allowedChannels = allowedChannels(principal);
DirectChannel requestedChannel = optionalChannel(rawChannelCode);
if (requestedChannel != null && !allowedChannels.contains(requestedChannel)) {
throw forbidden();
}
List<String> channelCodes = requestedChannel == null
? allowedChannels.stream().map(DirectChannel::code).toList()
: List.of(requestedChannel.code());
String status = status(rawStatus);
Page page = page(rawPageNum, rawPageSize);
Query query = adminQuery(tenantId, channelCodes, status);
Long total = jdbcTemplate.queryForObject("select count(*) " + query.fromWhere(),
Long.class, query.parameters().toArray());
List<Object> parameters = new ArrayList<>(query.parameters());
parameters.add(page.pageSize());
parameters.add(page.offset());
List<AdminFeedbackItem> rows = jdbcTemplate.query("""
select id, channel_code, sender_name, anonymous_flag, content, status,
reply_content, replied_time, create_time
""" + query.fromWhere() + """
order by create_time desc, id desc
limit ? offset ?
""", (rs, rowNum) -> adminItem(
rs.getLong("id"), rs.getString("channel_code"), rs.getString("sender_name"),
rs.getBoolean("anonymous_flag"), rs.getString("content"), rs.getString("status"),
rs.getString("reply_content"), rs.getTimestamp("replied_time"), rs.getTimestamp("create_time")
), parameters.toArray());
return new AdminFeedbackPage(tenantId, total == null ? 0 : total, page.pageNum(), page.pageSize(), rows);
}
@Transactional
public AdminFeedbackItem reply(Long rawId, ReplyRequest request) {
AihrKnowledgePrincipal principal = principalResolver.current();
String tenantId = adminTenant(principal);
requireExpectedTenant(request.expectedTenantId(), tenantId);
List<String> channelCodes = allowedChannels(principal).stream().map(DirectChannel::code).toList();
long feedbackId = id(rawId);
String replyContent = required(request.content(), "回复内容", 5000);
AdminFeedbackItem existing = adminDetail(tenantId, feedbackId, channelCodes);
if ("REPLIED".equals(existing.status())) {
if (replyContent.equals(existing.replyContent())) {
return existing;
}
throw new ServiceException("该反馈已经回复,不能覆盖原回复", HttpStatus.CONFLICT);
}
int updated = jdbcTemplate.update("""
update aihr_direct_feedback
set status = 'REPLIED', reply_content = ?, replied_by = ?, replied_time = now(), update_time = now()
where tenant_id = ? and id = ? and status = 'SUBMITTED'
""", replyContent, principal.userId(), tenantId, feedbackId);
if (updated != 1) {
AdminFeedbackItem concurrent = adminDetail(tenantId, feedbackId, channelCodes);
if ("REPLIED".equals(concurrent.status()) && replyContent.equals(concurrent.replyContent())) {
return concurrent;
}
throw new ServiceException("该反馈已经回复,不能覆盖原回复", HttpStatus.CONFLICT);
}
return adminDetail(tenantId, feedbackId, channelCodes);
}
private FeedbackItem replay(AihrKnowledgePrincipal principal, String requestKey, String requestHash) {
List<Replay> rows = jdbcTemplate.query("""
select id, submit_request_hash
from aihr_direct_feedback
where tenant_id = ? and sender_user_id = ? and submit_request_key = ?
""", (rs, rowNum) -> new Replay(rs.getLong("id"), rs.getString("submit_request_hash")),
principal.tenantId(), principal.userId(), requestKey);
if (rows.isEmpty()) {
return null;
}
Replay replay = rows.get(0);
if (!requestHash.equals(replay.requestHash())) {
throw new ServiceException("该请求ID已用于不同的反馈", HttpStatus.CONFLICT);
}
return mineDetail(principal, replay.id());
}
private FeedbackItem mineDetail(AihrKnowledgePrincipal principal, long feedbackId) {
List<FeedbackItem> rows = jdbcTemplate.query("""
select id, channel_code, anonymous_flag, content, status, reply_content, replied_time, create_time
from aihr_direct_feedback
where tenant_id = ? and sender_user_id = ? and id = ?
""", (rs, rowNum) -> feedbackItem(
rs.getLong("id"), rs.getString("channel_code"), rs.getBoolean("anonymous_flag"),
rs.getString("content"), rs.getString("status"), rs.getString("reply_content"),
rs.getTimestamp("replied_time"), rs.getTimestamp("create_time")
), principal.tenantId(), principal.userId(), feedbackId);
if (rows.isEmpty()) {
throw new ServiceException("反馈不存在或无权查看", HttpStatus.NOT_FOUND);
}
return rows.get(0);
}
private AdminFeedbackItem adminDetail(String tenantId, long feedbackId, List<String> channelCodes) {
String placeholders = String.join(",", Collections.nCopies(channelCodes.size(), "?"));
List<Object> parameters = new ArrayList<>();
parameters.add(tenantId);
parameters.add(feedbackId);
parameters.addAll(channelCodes);
List<AdminFeedbackItem> rows = jdbcTemplate.query("""
select id, channel_code, sender_name, anonymous_flag, content, status,
reply_content, replied_time, create_time
from aihr_direct_feedback
where tenant_id = ? and id = ? and channel_code in (""" + placeholders + ")",
(rs, rowNum) -> adminItem(
rs.getLong("id"), rs.getString("channel_code"), rs.getString("sender_name"),
rs.getBoolean("anonymous_flag"), rs.getString("content"), rs.getString("status"),
rs.getString("reply_content"), rs.getTimestamp("replied_time"), rs.getTimestamp("create_time")
), parameters.toArray());
if (rows.isEmpty()) {
throw new ServiceException("反馈不存在或无权处理", HttpStatus.NOT_FOUND);
}
return rows.get(0);
}
private String requireEmployee(AihrKnowledgePrincipal principal) {
if (!UserType.APP_USER.getUserType().equals(principal.userType())) {
throw new ServiceException("仅在职员工可使用直达反馈", HttpStatus.FORBIDDEN);
}
List<String> names = jdbcTemplate.query("""
select distinct coalesce(nullif(trim(o.person_name), ''), '员工') as person_name
from sys_user u
join aihr_org_snapshot o
on binary o.tenant_id = binary u.tenant_id
and binary o.person_phone = binary u.phonenumber
where u.tenant_id = ? and u.user_id = ? and u.user_type = ?
and u.status = '0' and u.del_flag = '0' and o.employment_status = 'active'
order by person_name limit 1
""", (rs, rowNum) -> rs.getString("person_name"),
principal.tenantId(), principal.userId(), UserType.APP_USER.getUserType());
if (names.isEmpty()) {
throw new ServiceException("仅在职员工可使用直达反馈", HttpStatus.FORBIDDEN);
}
return names.get(0);
}
private static List<DirectChannel> allowedChannels(AihrKnowledgePrincipal principal) {
if (UserType.APP_USER.getUserType().equals(principal.userType())) {
throw forbidden();
}
Set<String> roles = new LinkedHashSet<>(principal.roles());
if (roles.contains(TenantConstants.SUPER_ADMIN_ROLE_KEY)) {
return CHANNELS;
}
List<DirectChannel> allowed = CHANNELS.stream().filter(channel -> roles.contains(channel.role())).toList();
if (allowed.isEmpty()) {
throw forbidden();
}
return allowed;
}
private static Query adminQuery(String tenantId, List<String> channelCodes, String status) {
String placeholders = String.join(",", Collections.nCopies(channelCodes.size(), "?"));
String statusFilter = status == null ? "" : " and status = ?";
List<Object> parameters = new ArrayList<>();
parameters.add(tenantId);
parameters.addAll(channelCodes);
if (status != null) {
parameters.add(status);
}
return new Query(" from aihr_direct_feedback where tenant_id = ? and channel_code in ("
+ placeholders + ")" + statusFilter, parameters);
}
private static FeedbackItem feedbackItem(
long id,
String channelCode,
boolean anonymous,
String content,
String status,
String replyContent,
Timestamp repliedAt,
Timestamp createdAt
) {
return new FeedbackItem(id, channelCode, channel(channelCode).name(), anonymous, content, status,
replyContent, format(repliedAt), format(createdAt));
}
private static AdminFeedbackItem adminItem(
long id,
String channelCode,
String senderName,
boolean anonymous,
String content,
String status,
String replyContent,
Timestamp repliedAt,
Timestamp createdAt
) {
return new AdminFeedbackItem(id, channelCode, channel(channelCode).name(),
anonymous ? "匿名员工" : required(senderName, "提交人姓名", 100), anonymous,
content, status, replyContent, format(repliedAt), format(createdAt));
}
private static DirectChannel channel(String rawCode) {
String code = rawCode == null ? "" : rawCode.trim().toUpperCase(Locale.ROOT);
return CHANNELS.stream().filter(channel -> channel.code().equals(code)).findFirst()
.orElseThrow(() -> new ServiceException("直达对象无效", HttpStatus.BAD_REQUEST));
}
private static DirectChannel optionalChannel(String rawCode) {
return rawCode == null || rawCode.isBlank() ? null : channel(rawCode);
}
private static String status(String rawStatus) {
String normalized = rawStatus == null ? "" : rawStatus.trim().toUpperCase(Locale.ROOT);
if (normalized.isEmpty() || "ALL".equals(normalized)) {
return null;
}
if ("SUBMITTED".equals(normalized) || "REPLIED".equals(normalized)) {
return normalized;
}
throw new ServiceException("status 只支持 SUBMITTED、REPLIED 或 all", HttpStatus.BAD_REQUEST);
}
private static Page page(Integer rawPageNum, Integer rawPageSize) {
int pageNum = rawPageNum == null ? 1 : rawPageNum;
int pageSize = rawPageSize == null ? DEFAULT_PAGE_SIZE : rawPageSize;
if (pageNum < 1 || pageSize < 1 || pageSize > MAX_PAGE_SIZE) {
throw new ServiceException("分页参数无效", HttpStatus.BAD_REQUEST);
}
return new Page(pageNum, pageSize, (long) (pageNum - 1) * pageSize);
}
private static long id(Long value) {
if (value == null || value < 1) {
throw new ServiceException("反馈编号无效", HttpStatus.BAD_REQUEST);
}
return value;
}
private static String requestKey(String value) {
String normalized = required(value, "请求ID", 100);
if (!REQUEST_ID_PATTERN.matcher(normalized).matches()) {
throw new ServiceException("请求ID格式不正确", HttpStatus.BAD_REQUEST);
}
return normalized;
}
private static String requestHash(String channelCode, boolean anonymous, String content) {
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
digest.update(channelCode.getBytes(StandardCharsets.UTF_8));
digest.update((byte) 0);
digest.update(anonymous ? (byte) 1 : (byte) 0);
digest.update((byte) 0);
digest.update(content.getBytes(StandardCharsets.UTF_8));
return HexFormat.of().formatHex(digest.digest());
} catch (NoSuchAlgorithmException impossible) {
throw new IllegalStateException(impossible);
}
}
private static String required(String value, String label, int maximumLength) {
String normalized = value == null ? "" : value.trim();
if (normalized.isEmpty()) {
throw new ServiceException(label + "不能为空", HttpStatus.BAD_REQUEST);
}
if (normalized.length() > maximumLength) {
throw new ServiceException(label + "不能超过 " + maximumLength + " 个字符", HttpStatus.BAD_REQUEST);
}
return normalized;
}
private static String adminTenant(AihrKnowledgePrincipal principal) {
String tenantId = TenantHelper.getTenantId();
return tenantId == null || tenantId.isBlank() ? principal.tenantId() : tenantId.trim();
}
private static void requireExpectedTenant(String rawExpectedTenantId, String tenantId) {
if (!required(rawExpectedTenantId, "当前租户", 64).equals(tenantId)) {
throw new ServiceException("当前租户已切换,请刷新页面后重试", HttpStatus.CONFLICT);
}
}
private static String format(Timestamp value) {
return value == null ? null
: value.toLocalDateTime().withNano(0).format(DateTimeFormatter.ISO_LOCAL_DATE_TIME);
}
private static ServiceException forbidden() {
return new ServiceException("无权处理该直达反馈", HttpStatus.FORBIDDEN);
}
private record DirectChannel(String code, String name, String role, boolean defaultAnonymous) {
}
private record Replay(long id, String requestHash) {
}
private record Query(String fromWhere, List<Object> parameters) {
}
private record Page(int pageNum, int pageSize, long offset) {
}
}
@@ -107,16 +107,26 @@ public class AihrKnowledgeQueryService {
boolean finalizeAudit = false;
QueryResponse response = null;
if (context.broadcastMessageId() != null) {
// Project selection has already been verified above. Knowledge-space authorization
// deliberately precedes message resolution, so a caller never learns message
// state while attempting to bypass the selected project or authorized SOP scope.
spaceIds = accessService.resolveInternalSpaceIds(
principal, app, request.spaceCodes(), "READ");
// A published company message is an independently authorized source. Extra SOP
// evidence is optional unless the caller explicitly requested a knowledge space.
try {
spaceIds = accessService.resolveInternalSpaceIds(
principal, app, request.spaceCodes(), "READ");
} catch (ServiceException ex) {
if (ex.getCode() != null && ex.getCode() == HttpStatus.FORBIDDEN
&& request.spaceCodes().isEmpty()) {
spaceIds = Set.of();
} else {
throw ex;
}
}
BroadcastQuestionContext broadcast = resolveBroadcastContext(principal, context.broadcastMessageId());
response = answerBroadcastQuestion(
broadcast,
request.queryText(),
queryDocuments(principal, app, spaceIds, routed, request.queryText())
spaceIds.isEmpty()
? emptyProjectResponse(request.queryText())
: queryDocuments(principal, app, spaceIds, routed, request.queryText())
);
// Revalidate immediately before the response can be persisted into the short
// conversation. A withdrawal (or an employee becoming ineligible) during the
@@ -1432,6 +1432,14 @@ public class AihrSopSeedService {
}
}
/**
* Reuses the configured document-insight model for a broadcast attachment without
* importing that attachment into the enterprise knowledge base.
*/
public String summarizeBroadcastAttachment(String fileName, String content) {
return documentInsight(fileName, content, "公司文件").summary();
}
private String callInsightModel(ChatRuntime runtime, String fileName, String content, boolean autoCategory, String manualCategory) throws Exception {
ObjectNode body = objectMapper.createObjectNode();
body.put("model", runtime.modelName());
@@ -0,0 +1,54 @@
package org.dromara.aihr.broadcast;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastAttachmentResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.PublishRequest;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Method;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
@Tag("dev")
class AihrBroadcastAttachmentContractTest {
@Test
void questionContextAddsServerSideExtractedTextWithoutAddingItToTheApiDto() throws Exception {
Method method = AihrBroadcastService.class.getDeclaredMethod("questionContent", String.class, String.class);
method.setAccessible(true);
String context = (String) method.invoke(null, "消息正文", "原文件提取内容");
List<String> responseFields = List.of(BroadcastAttachmentResponse.class.getRecordComponents()).stream()
.map(component -> component.getName())
.toList();
assertTrue(context.contains("消息正文"));
assertTrue(context.contains("原文件提取内容"));
assertFalse(responseFields.contains("extractedText"));
}
@Test
void publishRequestKeepsLegacyConstructorAndAcceptsOneAttachment() {
PublishRequest legacy = new PublishRequest("request-1", "标题", "正文", "000000");
PublishRequest attached = new PublishRequest(
"request-2", "标题", "正文", "000000", false,
AihrBroadcastDto.BroadcastTargetRequest.empty(), 42L
);
assertEquals(null, legacy.attachmentId());
assertEquals(42L, attached.attachmentId());
}
@Test
void backgroundWorkerReadsOssInsideTheOwningTenant() throws Exception {
String source = Files.readString(Path.of(
"src/main/java/org/dromara/aihr/broadcast/AihrBroadcastAttachmentService.java"));
assertTrue(source.contains(
"TenantHelper.dynamic(work.tenantId(), () -> ossService.getById(work.ossId()))"));
}
}
@@ -0,0 +1,85 @@
package org.dromara.aihr.direct;
import cn.dev33.satoken.annotation.SaCheckLogin;
import org.dromara.aihr.direct.AihrDirectDto.AdminFeedbackItem;
import org.dromara.aihr.knowledge.domain.AihrKnowledgePrincipal;
import org.dromara.common.core.constant.HttpStatus;
import org.dromara.common.core.constant.TenantConstants;
import org.dromara.common.core.enums.UserType;
import org.dromara.common.core.exception.ServiceException;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import org.springframework.transaction.annotation.Transactional;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Set;
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;
@Tag("dev")
class AihrDirectSecurityContractTest {
@Test
void departmentRoleCanOnlyOpenItsOwnInboxWhileSuperadminCanOpenAll() throws Exception {
List<?> auditChannels = allowedChannels(principal("sys_user", Set.of("direct_audit")));
List<?> allChannels = allowedChannels(principal("sys_user", Set.of(TenantConstants.SUPER_ADMIN_ROLE_KEY)));
assertEquals(1, auditChannels.size());
assertTrue(auditChannels.get(0).toString().contains("AUDIT"));
assertEquals(5, allChannels.size());
}
@Test
void employeeCannotCallTheHandlerScope() {
InvocationTargetException wrapped = assertThrows(InvocationTargetException.class,
() -> allowedChannels(principal(UserType.APP_USER.getUserType(), Set.of("employee"))));
ServiceException error = (ServiceException) wrapped.getCause();
assertEquals(HttpStatus.FORBIDDEN, error.getCode());
}
@Test
void businessInboxDtoNeverExposesInternalIdentityKeys() {
List<String> fields = List.of(AdminFeedbackItem.class.getRecordComponents()).stream()
.map(component -> component.getName().toLowerCase())
.toList();
assertFalse(fields.contains("senderuserid"));
assertFalse(fields.contains("extpartyid"));
assertFalse(fields.contains("phone"));
assertTrue(fields.contains("senderdisplayname"));
}
@Test
void endpointsRequireLoginAndRepliesAreTransactional() throws Exception {
assertTrue(AihrDirectController.class.isAnnotationPresent(SaCheckLogin.class));
assertTrue(AihrDirectService.class.getMethod("reply", Long.class, AihrDirectDto.ReplyRequest.class)
.isAnnotationPresent(Transactional.class));
}
@Test
void employeeIdentityUsesOnlyTheVerifiedMobilePhone() throws Exception {
String source = Files.readString(Path.of(
"src/main/java/org/dromara/aihr/direct/AihrDirectService.java"));
assertTrue(source.contains("binary o.person_phone = binary u.phonenumber"));
assertFalse(source.contains("o.ext_party_id = binary u.user_name"));
}
private static List<?> allowedChannels(AihrKnowledgePrincipal principal) throws Exception {
Method method = AihrDirectService.class.getDeclaredMethod("allowedChannels", AihrKnowledgePrincipal.class);
method.setAccessible(true);
return (List<?>) method.invoke(null, principal);
}
private static AihrKnowledgePrincipal principal(String userType, Set<String> roles) {
return new AihrKnowledgePrincipal("tenant-a", 7L, userType, "", roles, Set.of(), "client");
}
}
@@ -607,6 +607,40 @@ class AihrKnowledgeQueryServiceTest {
eq(List.of("BCAST")), eq("SUCCESS"), anyLong(), eq("broadcast-context-v1"));
}
@Test
void broadcastQuestionWorksWithoutAnyKnowledgeSpace() {
var resolver = mock(AihrKnowledgePrincipalResolver.class);
var appService = mock(AihrKnowledgeAppService.class);
var access = mock(AihrKnowledgeAccessService.class);
var sop = mock(AihrSopSeedService.class);
var broadcast = mock(AihrBroadcastService.class);
var model = mock(AihrModelSeedService.class);
var principal = new AihrKnowledgePrincipal("000000", 7L, "app_user", "employee-7",
Set.of("employee"), Set.of(), "app");
var app = new AuthenticatedApp(3L, "000000", "yc_mobile", "员工端", "SESSION", 60, null);
when(resolver.current()).thenReturn(principal);
when(appService.requireSessionApp("000000", "app")).thenReturn(app);
when(access.resolveInternalSpaceIds(principal, app, List.of(), "READ"))
.thenThrow(new ServiceException("当前应用和账号没有共同可访问的知识空间", 403));
when(broadcast.resolveQuestionContext(principal, 42L)).thenReturn(
new AihrBroadcastService.BroadcastQuestionContext(
42L, "客户投诉首问负责制", "30分钟内登记,2小时内反馈首次处理进展。", "2026-07-24T20:10:00"));
when(model.tryChat(anyString(), anyString(), eq(0.0))).thenReturn(
Optional.of("应在30分钟内登记,并在2小时内反馈首次处理进展。"));
var service = new AihrKnowledgeQueryService(resolver, appService, access, sop,
mock(AihrKnowledgeQueryAuditService.class), mock(JdbcTemplate.class),
mock(AihrKnowledgeDataToolService.class), mock(AihrKnowledgeConversationService.class),
mock(AihrMemoryService.class), broadcast, model);
var result = service.queryInternal(new QueryRequest(
"多久登记和反馈?", List.of(), "sop", null, "mobile", 5, null,
null, null, null, 42L));
assertTrue(result.answer().contains("30分钟"));
assertEquals(42L, result.broadcastContext().messageId());
verify(sop, never()).searchAuthorized(any(), any(), any());
}
@Test
void broadcastContextIsRejectedForMediaAndExternalCalls() {
var service = service(mock(AihrKnowledgePrincipalResolver.class), mock(AihrKnowledgeAppService.class),
@@ -0,0 +1,58 @@
-- 公司消息 R2:单文件异步提炼、原文件受控下载及消息追问上下文。
-- 可重复执行,需在 aihr_20260723_broadcast_targeting_mysql8.sql 之后执行。
CREATE TABLE IF NOT EXISTS `aihr_broadcast_attachment` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '公司文件ID',
`tenant_id` varchar(20) NOT NULL COMMENT '租户编号',
`message_id` bigint DEFAULT NULL COMMENT '绑定的公司消息ID',
`oss_id` bigint NOT NULL COMMENT '原文件OSS编号',
`file_name` varchar(255) NOT NULL COMMENT '原文件名',
`file_size` bigint NOT NULL COMMENT '文件字节数',
`content_type` varchar(100) NOT NULL COMMENT '内容类型',
`status` varchar(20) NOT NULL DEFAULT 'QUEUED' COMMENT 'QUEUED/PROCESSING/READY/FAILED',
`extracted_text` longtext DEFAULT NULL COMMENT '提取文本,仅供服务端上下文使用',
`summary` varchar(1000) DEFAULT NULL COMMENT 'AI提炼摘要',
`error_message` varchar(500) DEFAULT NULL COMMENT '公开失败说明',
`uploaded_by` bigint NOT NULL COMMENT '上传账号ID',
`create_time` datetime NOT NULL COMMENT '上传时间',
`update_time` datetime NOT NULL COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_aihr_broadcast_attachment_oss` (`tenant_id`, `oss_id`),
UNIQUE KEY `uk_aihr_broadcast_attachment_message` (`tenant_id`, `message_id`),
KEY `idx_aihr_broadcast_attachment_queue` (`status`, `update_time`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='银城大喇叭公司文件提炼';
SET @has_broadcast_attachment_id := (
SELECT COUNT(*) FROM information_schema.COLUMNS
WHERE table_schema = DATABASE() AND table_name = 'aihr_broadcast_message'
AND column_name = 'attachment_id'
);
SET @broadcast_attachment_ddl := IF(
@has_broadcast_attachment_id = 0,
'ALTER TABLE aihr_broadcast_message ADD COLUMN attachment_id bigint DEFAULT NULL COMMENT ''公司文件ID'' AFTER content, ADD KEY idx_aihr_broadcast_message_attachment (tenant_id, attachment_id)',
'SELECT 1'
);
PREPARE aihr_broadcast_attachment_stmt FROM @broadcast_attachment_ddl; EXECUTE aihr_broadcast_attachment_stmt; DEALLOCATE PREPARE aihr_broadcast_attachment_stmt;
SET @broadcast_attachment_target_collation := (
SELECT collation_name FROM information_schema.columns
WHERE table_schema = DATABASE() AND table_name = 'aihr_knowledge_info' AND column_name = 'tenant_id'
LIMIT 1
);
SET @broadcast_attachment_target_collation := COALESCE(@broadcast_attachment_target_collation, @@collation_database);
SET @broadcast_attachment_table_collation := (
SELECT table_collation FROM information_schema.tables
WHERE table_schema = DATABASE() AND table_name = 'aihr_broadcast_attachment'
LIMIT 1
);
SET @broadcast_attachment_ddl := IF(
@broadcast_attachment_table_collation = @broadcast_attachment_target_collation,
'SELECT 1',
CONCAT('ALTER TABLE aihr_broadcast_attachment CONVERT TO CHARACTER SET utf8mb4 COLLATE ', @broadcast_attachment_target_collation)
);
PREPARE aihr_broadcast_attachment_stmt FROM @broadcast_attachment_ddl; EXECUTE aihr_broadcast_attachment_stmt; DEALLOCATE PREPARE aihr_broadcast_attachment_stmt;
SET @has_broadcast_attachment_id := NULL;
SET @broadcast_attachment_target_collation := NULL;
SET @broadcast_attachment_table_collation := NULL;
SET @broadcast_attachment_ddl := NULL;
@@ -0,0 +1,71 @@
-- 公司消息 R1:员工点对点直达反馈。业务匿名不显示提交人,但保留内部账号用于本人查询和审计。
CREATE TABLE IF NOT EXISTS `aihr_direct_feedback` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '反馈ID',
`tenant_id` varchar(20) NOT NULL COMMENT '租户编号',
`channel_code` varchar(20) NOT NULL COMMENT 'PRESIDENT/FINANCE/HR/AUDIT/OPERATIONS',
`sender_user_id` bigint NOT NULL COMMENT '内部提交账号ID,不在业务处理接口返回',
`sender_name` varchar(100) NOT NULL COMMENT '提交时员工姓名快照',
`anonymous_flag` tinyint(1) NOT NULL DEFAULT 0 COMMENT '业务展示是否匿名',
`content` text NOT NULL COMMENT '反馈内容',
`status` varchar(20) NOT NULL DEFAULT 'SUBMITTED' COMMENT 'SUBMITTED/REPLIED',
`submit_request_key` varchar(100) NOT NULL COMMENT '员工提交幂等键',
`submit_request_hash` char(64) NOT NULL COMMENT '提交载荷SHA-256',
`reply_content` text DEFAULT NULL COMMENT '正式回复',
`replied_by` bigint DEFAULT NULL COMMENT '回复账号ID',
`replied_time` datetime DEFAULT NULL COMMENT '回复时间',
`create_time` datetime NOT NULL COMMENT '提交时间',
`update_time` datetime NOT NULL COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_aihr_direct_feedback_submit` (`tenant_id`, `sender_user_id`, `submit_request_key`),
KEY `idx_aihr_direct_feedback_inbox` (`tenant_id`, `channel_code`, `status`, `create_time`),
KEY `idx_aihr_direct_feedback_sender` (`tenant_id`, `sender_user_id`, `create_time`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='员工点对点直达反馈';
-- 与租户事实表保持相同排序规则,兼容历史库。
SET @direct_target_collation := (
SELECT collation_name FROM information_schema.columns
WHERE table_schema = DATABASE() AND table_name = 'aihr_knowledge_info' AND column_name = 'tenant_id'
LIMIT 1
);
SET @direct_target_collation := COALESCE(@direct_target_collation, @@collation_database);
SET @direct_table_collation := (
SELECT table_collation FROM information_schema.tables
WHERE table_schema = DATABASE() AND table_name = 'aihr_direct_feedback'
LIMIT 1
);
SET @direct_ddl := IF(
@direct_table_collation = @direct_target_collation,
'SELECT 1',
CONCAT('ALTER TABLE aihr_direct_feedback CONVERT TO CHARACTER SET utf8mb4 COLLATE ', @direct_target_collation)
);
PREPARE aihr_direct_stmt FROM @direct_ddl; EXECUTE aihr_direct_stmt; DEALLOCATE PREPARE aihr_direct_stmt;
SET @direct_target_collation := NULL;
SET @direct_table_collation := NULL;
SET @direct_ddl := NULL;
-- 当前正式租户预置五个处理角色;具体人员仍由管理员按职责分配。
INSERT INTO `sys_role`
(`role_id`, `tenant_id`, `role_name`, `role_key`, `role_sort`, `data_scope`,
`menu_check_strictly`, `dept_check_strictly`, `status`, `del_flag`,
`create_by`, `create_time`, `remark`)
SELECT seed.role_id, '000000', seed.role_name, seed.role_key, seed.role_sort, '1',
1, 1, '0', '0', 1, NOW(), '公司消息直通车固定处理角色'
FROM (
SELECT 202607240000000101 AS role_id, '总裁直达处理' AS role_name, 'direct_president' AS role_key, 21 AS role_sort
UNION ALL SELECT 202607240000000102, '财务直达处理', 'direct_finance', 22
UNION ALL SELECT 202607240000000103, '人力直达处理', 'direct_hr', 23
UNION ALL SELECT 202607240000000104, '审计直达处理', 'direct_audit', 24
UNION ALL SELECT 202607240000000105, '运营直达处理', 'direct_operations', 25
) seed
WHERE NOT EXISTS (
SELECT 1 FROM `sys_role` existing
WHERE existing.tenant_id = '000000'
AND existing.role_key = seed.role_key
AND existing.del_flag = '0'
)
AND NOT EXISTS (
SELECT 1 FROM `sys_role` existing
WHERE existing.role_id = seed.role_id
);