feat(broadcast): expose employee role insights

This commit is contained in:
2026-07-25 14:33:27 +08:00
parent 91ebe41259
commit 7990db6883
5 changed files with 184 additions and 19 deletions
@@ -49,6 +49,7 @@ public class AihrBroadcastAttachmentService {
private final ISysOssService ossService;
private final KnowledgeDocumentParser documentParser;
private final AihrSopSeedService sopService;
private final AihrBroadcastInsightService insightService;
private final AihrKnowledgePrincipalResolver principalResolver;
private final AihrBroadcastService broadcastService;
private final ScheduledExecutorService scheduledExecutorService;
@@ -109,13 +110,14 @@ public class AihrBroadcastAttachmentService {
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
select id, file_name, file_size, content_type, status, summary, error_message,
insight_status, insights_json
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")
rs.getString("error_message"), rs.getString("insight_status"), rs.getString("insights_json")
), tenantId, id);
if (rows.isEmpty()) {
throw new ServiceException("公司文件不存在或无权查看", HttpStatus.NOT_FOUND);
@@ -156,6 +158,12 @@ public class AihrBroadcastAttachmentService {
set status = 'QUEUED', error_message = null, update_time = now()
where status = 'PROCESSING' and update_time < date_sub(now(), interval 30 minute)
""");
jdbcTemplate.update("""
update aihr_broadcast_attachment
set insight_status = 'PENDING', update_time = now()
where status = 'READY' and insight_status = 'PROCESSING'
and update_time < date_sub(now(), interval 2 hour)
""");
trigger();
} catch (RuntimeException error) {
log.warn("broadcast attachment recovery failed");
@@ -179,9 +187,18 @@ public class AihrBroadcastAttachmentService {
}
private void processQueue() {
AttachmentWork work;
while ((work = claim()) != null) {
process(work);
while (true) {
AttachmentWork attachment = claim();
if (attachment != null) {
process(attachment);
continue;
}
InsightWork insight = claimInsight();
if (insight != null) {
processInsight(insight);
continue;
}
return;
}
}
@@ -208,6 +225,30 @@ public class AihrBroadcastAttachmentService {
return null;
}
private InsightWork claimInsight() {
List<InsightWork> rows = jdbcTemplate.query("""
select id, tenant_id, file_name, extracted_text
from aihr_broadcast_attachment
where status = 'READY' and insight_status = 'PENDING'
order by id
limit 10
""", (rs, rowNum) -> new InsightWork(
rs.getLong("id"), rs.getString("tenant_id"), rs.getString("file_name"),
rs.getString("extracted_text")
));
for (InsightWork row : rows) {
int claimed = jdbcTemplate.update("""
update aihr_broadcast_attachment
set insight_status = 'PROCESSING', update_time = now()
where tenant_id = ? and id = ? and status = 'READY' and insight_status = 'PENDING'
""", 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()));
@@ -228,7 +269,9 @@ public class AihrBroadcastAttachmentService {
() -> 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()
set status = 'READY', extracted_text = ?, summary = ?,
insight_status = 'PENDING', insights_json = null, insight_version = null,
error_message = null, update_time = now()
where tenant_id = ? and id = ? and status = 'PROCESSING'
""", text, summary, work.tenantId(), work.id());
if (updated != 1) {
@@ -244,9 +287,34 @@ public class AihrBroadcastAttachmentService {
}
}
private void processInsight(InsightWork work) {
try {
AihrBroadcastInsightService.Generation generation = TenantHelper.dynamic(work.tenantId(),
() -> insightService.generate(work.fileName(), work.extractedText()));
int updated = jdbcTemplate.update("""
update aihr_broadcast_attachment
set insight_status = ?, insights_json = ?, insight_version = ?, update_time = now()
where tenant_id = ? and id = ? and status = 'READY' and insight_status = 'PROCESSING'
""", generation.status(), generation.json(), AihrBroadcastInsightService.VERSION,
work.tenantId(), work.id());
if (updated != 1) {
log.info("broadcast attachment {} insights completed after its state changed", work.id());
}
} catch (RuntimeException error) {
log.warn("broadcast attachment {} insight generation failed", work.id());
jdbcTemplate.update("""
update aihr_broadcast_attachment
set insight_status = 'FAILED', insights_json = null, insight_version = ?,
update_time = now()
where tenant_id = ? and id = ? and status = 'READY' and insight_status = 'PROCESSING'
""", AihrBroadcastInsightService.VERSION, work.tenantId(), work.id());
}
}
private boolean hasQueued() {
Long count = jdbcTemplate.queryForObject("""
select count(*) from aihr_broadcast_attachment where status = 'QUEUED'
select count(*) from aihr_broadcast_attachment
where status = 'QUEUED' or (status = 'READY' and insight_status = 'PENDING')
""", Long.class);
return count != null && count > 0;
}
@@ -260,17 +328,25 @@ public class AihrBroadcastAttachmentService {
}
}
private static BroadcastAttachmentResponse response(
private BroadcastAttachmentResponse response(
long id,
String fileName,
long fileSize,
String contentType,
String status,
String summary,
String errorMessage
String errorMessage,
String rawInsightStatus,
String insightsJson
) {
String insightStatus = "PROCESSING".equals(rawInsightStatus) ? "PENDING"
: rawInsightStatus == null ? "NOT_PROCESSED" : rawInsightStatus;
List<String> labels = AihrBroadcastInsightService.readPerspectives(insightsJson).stream()
.map(AihrBroadcastDto.BroadcastPerspective::label)
.toList();
return new BroadcastAttachmentResponse(id, fileName, fileSize, contentType, status, summary,
errorMessage, "READY".equals(status) ? "/api/aihr/broadcast/attachments/" + id + "/content" : null);
errorMessage, "READY".equals(status) ? "/api/aihr/broadcast/attachments/" + id + "/content" : null,
insightStatus, null, labels, List.of());
}
private static void requireBroadcastAdmin(AihrKnowledgePrincipal principal) {
@@ -323,4 +399,12 @@ public class AihrBroadcastAttachmentService {
long uploadedBy
) {
}
private record InsightWork(
long id,
String tenantId,
String fileName,
String extractedText
) {
}
}
@@ -32,6 +32,7 @@ public class AihrBroadcastInsightService {
private static final int MAX_PERSPECTIVES = 9;
private static final int MAX_ITEMS = 5;
private static final int MAX_EVIDENCE = 3;
private static final ObjectMapper STORED_JSON_MAPPER = new ObjectMapper();
private static final Map<String, String> ROLE_LABELS = Map.ofEntries(
entry("living_advisor", "生活顾问/客服"),
entry("cleaning", "保洁"),
@@ -97,12 +98,12 @@ public class AihrBroadcastInsightService {
}
}
public List<BroadcastPerspective> readPerspectives(String storedJson) {
public static List<BroadcastPerspective> readPerspectives(String storedJson) {
if (storedJson == null || storedJson.isBlank()) {
return List.of();
}
try {
JsonNode rows = objectMapper.readTree(storedJson).path("perspectives");
JsonNode rows = STORED_JSON_MAPPER.readTree(storedJson).path("perspectives");
if (!rows.isArray()) {
return List.of();
}
@@ -203,6 +203,8 @@ public class AihrBroadcastService {
AihrKnowledgePrincipal principal = principalResolver.current();
requireEmployeeAudience(principal);
long messageId = messageId(rawMessageId);
String defaultPerspectiveCode = AihrBroadcastInsightService.defaultPerspectiveCode(
currentPositionNames(principal));
List<BroadcastDetailResponse> rows = jdbcTemplate.query("""
select m.id, m.title, m.content, m.published_time,
case when r.id is null then 0 else 1 end as read_flag,
@@ -212,7 +214,9 @@ public class AihrBroadcastService {
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
a.error_message as attachment_error_message,
a.insight_status as attachment_insight_status,
a.insights_json as attachment_insights_json
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 = ?
@@ -231,7 +235,8 @@ public class AihrBroadcastService {
""", (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"), attachment(rs)
rs.getBoolean("targeted_flag"), rs.getString("target_reason"),
attachment(rs, defaultPerspectiveCode)
), principal.userId(), principal.tenantId(), principal.userId(), principal.tenantId(), messageId);
if (rows.isEmpty()) {
throw unavailable();
@@ -549,7 +554,7 @@ public class AihrBroadcastService {
from aihr_org_snapshot o
left join sys_user u
on binary o.tenant_id = binary u.tenant_id
and (binary o.person_phone = binary u.phonenumber or binary o.ext_party_id = binary u.user_name)
and binary o.person_phone = binary u.phonenumber
and u.user_type = 'app_user'
and u.status = '0'
and u.del_flag = '0'
@@ -644,7 +649,7 @@ public class AihrBroadcastService {
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 or binary o.ext_party_id = binary u.user_name)
and binary o.person_phone = binary u.phonenumber
where u.tenant_id = ?
and u.user_id = ?
and u.user_type = ?
@@ -657,6 +662,26 @@ public class AihrBroadcastService {
}
}
private List<String> currentPositionNames(AihrKnowledgePrincipal principal) {
return jdbcTemplate.query("""
select distinct o.position_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'
and o.position_name is not null
and o.position_name <> ''
order by o.position_name
""", (rs, rowNum) -> rs.getString("position_name"),
principal.tenantId(), principal.userId(), UserType.APP_USER.getUserType());
}
/**
* Broadcast records are read through explicit JDBC predicates, so they do not receive
* the ORM tenant interceptor automatically. Dynamic-tenant mode must therefore use the
@@ -710,14 +735,29 @@ public class AihrBroadcastService {
}
}
private static BroadcastAttachmentResponse attachment(java.sql.ResultSet rs) throws java.sql.SQLException {
private static BroadcastAttachmentResponse attachment(
java.sql.ResultSet rs,
String requestedDefaultPerspectiveCode
) throws java.sql.SQLException {
Long id = rs.getObject("attachment_id", Long.class);
if (id == null) return null;
List<AihrBroadcastDto.BroadcastPerspective> perspectives =
AihrBroadcastInsightService.readPerspectives(rs.getString("attachment_insights_json"));
String defaultPerspectiveCode = perspectives.stream()
.anyMatch(perspective -> perspective.code().equals(requestedDefaultPerspectiveCode))
? requestedDefaultPerspectiveCode : null;
String rawInsightStatus = rs.getString("attachment_insight_status");
String insightStatus = "PROCESSING".equals(rawInsightStatus) ? "PENDING"
: rawInsightStatus == null ? "NOT_PROCESSED" : rawInsightStatus;
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"
"/api/aihr/broadcast/attachments/" + id + "/content",
insightStatus,
defaultPerspectiveCode,
perspectives.stream().map(AihrBroadcastDto.BroadcastPerspective::label).toList(),
perspectives
);
}
@@ -54,5 +54,9 @@ class AihrBroadcastAttachmentContractTest {
assertTrue(source.contains(
"TenantHelper.dynamic(work.tenantId(), () -> ossService.getById(work.ossId()))"));
assertTrue(source.contains("set status = 'READY', extracted_text = ?, summary = ?,"));
assertTrue(source.contains("insight_status = 'PENDING'"));
assertTrue(source.contains("insight_status = 'FAILED'"));
assertTrue(source.contains("insight_status = 'PROCESSING'"));
}
}
@@ -32,6 +32,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import static java.util.Map.entry;
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -65,7 +66,7 @@ class AihrBroadcastServiceTest {
assertTrue(jdbcTemplate.queryForObjectSql.contains("from sys_user u"));
assertTrue(jdbcTemplate.queryForObjectSql.contains("binary o.tenant_id = binary u.tenant_id"));
assertTrue(jdbcTemplate.queryForObjectSql.contains("binary o.person_phone = binary u.phonenumber"));
assertTrue(jdbcTemplate.queryForObjectSql.contains("binary o.ext_party_id = binary u.user_name"));
assertFalse(jdbcTemplate.queryForObjectSql.contains(" or binary o.ext_party_id"));
assertArrayEquals(new Object[]{"tenant-a", 7L, "app_user"}, jdbcTemplate.queryForObjectArgs);
}
@@ -113,6 +114,38 @@ class AihrBroadcastServiceTest {
assertEquals(null, jdbcTemplate.querySql);
}
@Test
void employeeAttachmentKeepsAllPerspectivesAndMarksOnlyTheMatchingDefault() throws Exception {
Method attachment = AihrBroadcastService.class.getDeclaredMethod(
"attachment", ResultSet.class, String.class);
attachment.setAccessible(true);
ResultSet row = RecordingJdbcTemplate.resultSet(Map.ofEntries(
entry("attachment_id", 9L),
entry("attachment_file_name", "制度.txt"),
entry("attachment_file_size", 1024L),
entry("attachment_content_type", "text/plain"),
entry("attachment_status", "READY"),
entry("attachment_summary", "文件摘要"),
entry("attachment_error_message", ""),
entry("attachment_insight_status", "READY"),
entry("attachment_insights_json", """
{"version":"role-perspective-v1","perspectives":[
{"code":"finance","label":"财务","summary":"财务视角","concerns":[],"impacts":[],
"actions":[],"risks":[],"evidence":[{"paragraphIndex":1,"quote":"财务部应在每月五日前完成费用复核并提交差异说明"}]},
{"code":"hr","label":"人力","summary":"人力视角","concerns":[],"impacts":[],
"actions":[],"risks":[],"evidence":[{"paragraphIndex":2,"quote":"人力资源部负责组织全体员工完成制度培训和考试"}]}
]}
""")
));
AihrBroadcastDto.BroadcastAttachmentResponse response =
(AihrBroadcastDto.BroadcastAttachmentResponse) attachment.invoke(null, row, "finance");
assertEquals("finance", response.defaultPerspectiveCode());
assertEquals(List.of("财务", "人力"), response.perspectiveLabels());
assertEquals(2, response.perspectives().size());
}
@Test
void repeatedReadUsesTheDatabaseUniqueRecordInsteadOfCreatingDuplicates() {
RecordingJdbcTemplate jdbcTemplate = new RecordingJdbcTemplate();
@@ -462,6 +495,9 @@ class AihrBroadcastServiceTest {
Object value = values.get(args[0]);
return value == null ? 0L : ((Number) value).longValue();
}
if ("getObject".equals(name)) {
return values.get(args[0]);
}
if ("getTimestamp".equals(name)) {
return values.get(args[0]);
}