feat(broadcast): finish multi-role insight experience

This commit is contained in:
2026-07-25 15:18:58 +08:00
parent 7990db6883
commit 230dfb4dfc
20 changed files with 552 additions and 81 deletions
@@ -30,6 +30,7 @@ import java.sql.Statement;
import java.util.List;
import java.util.Locale;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@@ -160,13 +161,13 @@ public class AihrBroadcastAttachmentService {
""");
jdbcTemplate.update("""
update aihr_broadcast_attachment
set insight_status = 'PENDING', update_time = now()
set insight_status = 'PENDING', insight_version = null, 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");
log.warn("broadcast attachment recovery failed type={}", error.getClass().getSimpleName());
}
}
@@ -178,7 +179,7 @@ public class AihrBroadcastAttachmentService {
.whenComplete((ignored, error) -> {
workerRunning.set(false);
if (error != null) {
log.warn("broadcast attachment worker failed");
log.warn("broadcast attachment worker failed type={}", error.getClass().getSimpleName());
}
if (hasQueued()) {
trigger();
@@ -234,16 +235,18 @@ public class AihrBroadcastAttachmentService {
limit 10
""", (rs, rowNum) -> new InsightWork(
rs.getLong("id"), rs.getString("tenant_id"), rs.getString("file_name"),
rs.getString("extracted_text")
rs.getString("extracted_text"), null
));
for (InsightWork row : rows) {
String claimToken = "claim-" + UUID.randomUUID().toString().replace("-", "").substring(0, 24);
int claimed = jdbcTemplate.update("""
update aihr_broadcast_attachment
set insight_status = 'PROCESSING', update_time = now()
set insight_status = 'PROCESSING', insight_version = ?, update_time = now()
where tenant_id = ? and id = ? and status = 'READY' and insight_status = 'PENDING'
""", row.tenantId(), row.id());
""", claimToken, row.tenantId(), row.id());
if (claimed == 1) {
return row;
return new InsightWork(
row.id(), row.tenantId(), row.fileName(), row.extractedText(), claimToken);
}
}
return null;
@@ -294,9 +297,10 @@ public class AihrBroadcastAttachmentService {
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'
where tenant_id = ? and id = ? and status = 'READY'
and insight_status = 'PROCESSING' and insight_version = ?
""", generation.status(), generation.json(), AihrBroadcastInsightService.VERSION,
work.tenantId(), work.id());
work.tenantId(), work.id(), work.claimToken());
if (updated != 1) {
log.info("broadcast attachment {} insights completed after its state changed", work.id());
}
@@ -306,8 +310,9 @@ public class AihrBroadcastAttachmentService {
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());
where tenant_id = ? and id = ? and status = 'READY'
and insight_status = 'PROCESSING' and insight_version = ?
""", AihrBroadcastInsightService.VERSION, work.tenantId(), work.id(), work.claimToken());
}
}
@@ -404,7 +409,8 @@ public class AihrBroadcastAttachmentService {
long id,
String tenantId,
String fileName,
String extractedText
String extractedText,
String claimToken
) {
}
}
@@ -3,6 +3,7 @@ package org.dromara.aihr.broadcast;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastEvidence;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastPerspective;
import org.dromara.aihr.domain.AihrModelDto.ChatRequest;
@@ -23,6 +24,7 @@ import static java.util.Map.entry;
@Service
@RequiredArgsConstructor
@Slf4j
public class AihrBroadcastInsightService {
public static final String VERSION = "role-perspective-v1";
@@ -49,7 +51,8 @@ public class AihrBroadcastInsightService {
private final AihrModelSeedService modelService;
public Generation generate(String fileName, String content) {
String source = content == null ? "" : content.replace("\r\n", "\n").replace('\r', '\n').trim();
String normalized = content == null ? "" : content.replace("\r\n", "\n").replace('\r', '\n').trim();
String source = AihrSensitiveText.forModel(normalized);
if (source.isEmpty()) {
return new Generation("FAILED", null, List.of());
}
@@ -61,12 +64,12 @@ public class AihrBroadcastInsightService {
for (SourceChunk chunk : chunks) {
try {
ChatResponse response = modelService.chat(new ChatRequest(
AihrSensitiveText.forModel(prompt(fileName, chunk)),
prompt(fileName, chunk),
null,
AihrSensitiveText.forModel("""
你负责解读物业公司内部文件。只返回 JSON,不要 markdown。
只允许使用给定原文;原文未提及的岗位、结论、行动或风险不要输出。
evidence.quote 必须逐字复制原文,不能改写或拼接。
每个结论的 evidenceQuote 必须逐字复制原文,不能改写或拼接。
""")
));
if (response == null || !"openai-compatible".equals(response.mode())
@@ -76,8 +79,9 @@ public class AihrBroadcastInsightService {
for (BroadcastPerspective perspective : parseModelPerspectives(response.answer(), chunk, source)) {
merged.computeIfAbsent(perspective.code(), PerspectiveAccumulator::new).add(perspective);
}
} catch (RuntimeException ignored) {
} catch (RuntimeException error) {
failedChunks++;
log.warn("broadcast insight chunk failed type={}", error.getClass().getSimpleName());
}
}
List<BroadcastPerspective> perspectives = merged.values().stream()
@@ -93,7 +97,8 @@ public class AihrBroadcastInsightService {
return new Generation(status,
objectMapper.writeValueAsString(new StoredInsights(VERSION, perspectives)),
perspectives);
} catch (Exception ignored) {
} catch (Exception error) {
log.warn("broadcast insight serialization failed type={}", error.getClass().getSimpleName());
return new Generation("FAILED", null, List.of());
}
}
@@ -118,7 +123,8 @@ public class AihrBroadcastInsightService {
}
}
return List.copyOf(perspectives);
} catch (Exception ignored) {
} catch (Exception error) {
log.warn("broadcast stored insights invalid type={}", error.getClass().getSimpleName());
return List.of();
}
}
@@ -193,37 +199,77 @@ public class AihrBroadcastInsightService {
return null;
}
List<BroadcastEvidence> evidence = new ArrayList<>();
JsonNode evidenceRows = row.path("evidence");
if (evidenceRows.isArray()) {
for (JsonNode item : evidenceRows) {
String quote = limited(item.path("quote").asText(), 200);
if (quote.length() < 20 || !chunk.text().contains(quote) || !source.contains(quote)) {
continue;
}
int paragraph = Math.max(chunk.paragraphIndex(), item.path("paragraphIndex").asInt(chunk.paragraphIndex()));
if (evidence.stream().noneMatch(existing -> existing.quote().equals(quote))) {
evidence.add(new BroadcastEvidence(paragraph, quote));
}
if (evidence.size() == MAX_EVIDENCE) {
break;
}
}
}
String summary = groundedText(row.path("summary"), 500, chunk, source, evidence);
List<String> concerns = groundedItems(row.path("concerns"), chunk, source, evidence);
List<String> impacts = groundedItems(row.path("impacts"), chunk, source, evidence);
List<String> actions = groundedItems(row.path("actions"), chunk, source, evidence);
List<String> risks = groundedItems(row.path("risks"), chunk, source, evidence);
if (evidence.isEmpty()) {
return null;
}
return new BroadcastPerspective(
code,
label,
limited(row.path("summary").asText(), 500),
items(row.path("concerns")),
items(row.path("impacts")),
items(row.path("actions")),
items(row.path("risks")),
summary,
concerns,
impacts,
actions,
risks,
evidence
);
}
private static List<String> groundedItems(
JsonNode node,
SourceChunk chunk,
String source,
List<BroadcastEvidence> evidence
) {
if (!node.isArray()) {
return List.of();
}
Set<String> values = new LinkedHashSet<>();
for (JsonNode item : node) {
String value = groundedText(item, 300, chunk, source, evidence);
if (!value.isEmpty()) {
values.add(value);
}
if (values.size() == MAX_ITEMS) {
break;
}
}
return List.copyOf(values);
}
private static String groundedText(
JsonNode item,
int maximum,
SourceChunk chunk,
String source,
List<BroadcastEvidence> evidence
) {
if (!item.isObject()) {
return "";
}
String text = limited(item.path("text").asText(), maximum);
String quote = limited(item.path("evidenceQuote").asText(), 200);
if (text.isEmpty() || quote.length() < 20) {
return "";
}
int offset = source.indexOf(quote, chunk.startOffset());
if (offset < chunk.startOffset() || offset + quote.length() > chunk.endOffset()) {
return "";
}
boolean existing = evidence.stream().anyMatch(value -> value.quote().equals(quote));
if (!existing && evidence.size() == MAX_EVIDENCE) {
return "";
}
if (!existing) {
evidence.add(new BroadcastEvidence(paragraphIndex(source, offset), quote));
}
return text;
}
private static BroadcastPerspective persistedPerspective(JsonNode row) {
String code = clean(row.path("code").asText());
String label = ROLE_LABELS.get(code);
@@ -274,13 +320,18 @@ public class AihrBroadcastInsightService {
finance=财务;hr=人力;operations=业务运营;audit_risk=审计/风控;management=管理层。
返回格式:
{"perspectives":[{"code":"finance","summary":"...","concerns":["..."],"impacts":["..."],
"actions":["..."],"risks":["..."],"evidence":[{"paragraphIndex":1,"quote":"原文逐字短句"}]}]}
{"perspectives":[{"code":"finance",
"summary":{"text":"...","evidenceQuote":"原文逐字短句"},
"concerns":[{"text":"...","evidenceQuote":"原文逐字短句"}],
"impacts":[{"text":"...","evidenceQuote":"原文逐字短句"}],
"actions":[{"text":"...","evidenceQuote":"原文逐字短句"}],
"risks":[{"text":"...","evidenceQuote":"原文逐字短句"}]}]}
规则:
1. 只输出本切片明确涉及的岗位;没有直接内容就不要输出该岗位。
2. 每类最多 5 条,evidence 最多 3 条且每条 20 到 200 字。
3. 不得使用常识补充原文没有写的责任、时限、处罚或流程。
2. 每个 summary 或列表项都必须携带自己的 evidenceQuote;没有原文证据就删除该项。
3. 每类最多 5 条,证据短句 20 到 200 字。
4. 不得使用常识补充原文没有写的责任、时限、处罚或流程。
文件名:%s
切片:%d
@@ -366,6 +417,14 @@ public class AihrBroadcastInsightService {
}
private void add(BroadcastPerspective perspective) {
long newEvidence = perspective.evidence().stream()
.map(BroadcastEvidence::quote)
.filter(quote -> !evidence.containsKey(quote))
.distinct()
.count();
if (evidence.size() + newEvidence > MAX_EVIDENCE) {
return;
}
if (summary.isEmpty()) summary = perspective.summary();
add(concerns, perspective.concerns(), MAX_ITEMS);
add(impacts, perspective.impacts(), MAX_ITEMS);
@@ -645,7 +645,7 @@ public class AihrBroadcastService {
throw employeeOnly();
}
Long matched = jdbcTemplate.queryForObject("""
select count(*)
select count(distinct o.ext_party_id)
from sys_user u
join aihr_org_snapshot o
on binary o.tenant_id = binary u.tenant_id
@@ -655,9 +655,23 @@ public class AihrBroadcastService {
and u.user_type = ?
and u.status = '0'
and u.del_flag = '0'
and char_length(u.phonenumber) = 11
and u.phonenumber regexp '^[0-9]{11}$'
and o.employment_status = 'active'
and o.ext_party_id is not null
and o.ext_party_id <> ''
and not exists (
select 1
from aihr_org_snapshot o2
where binary o2.tenant_id = binary o.tenant_id
and binary o2.ext_party_id = binary o.ext_party_id
and o2.employment_status = 'active'
and o2.person_phone is not null
and o2.person_phone <> ''
and binary o2.person_phone <> binary u.phonenumber
)
""", Long.class, principal.tenantId(), principal.userId(), UserType.APP_USER.getUserType());
if (matched == null || matched == 0) {
if (matched == null || matched != 1) {
throw employeeOnly();
}
}
@@ -674,7 +688,32 @@ public class AihrBroadcastService {
and u.user_type = ?
and u.status = '0'
and u.del_flag = '0'
and char_length(u.phonenumber) = 11
and u.phonenumber regexp '^[0-9]{11}$'
and o.employment_status = 'active'
and o.ext_party_id is not null
and o.ext_party_id <> ''
and o.ext_party_id = (
select min(o2.ext_party_id)
from aihr_org_snapshot o2
where binary o2.tenant_id = binary u.tenant_id
and binary o2.person_phone = binary u.phonenumber
and o2.employment_status = 'active'
and o2.ext_party_id is not null
and o2.ext_party_id <> ''
group by o2.person_phone
having count(distinct o2.ext_party_id) = 1
)
and not exists (
select 1
from aihr_org_snapshot o3
where binary o3.tenant_id = binary o.tenant_id
and binary o3.ext_party_id = binary o.ext_party_id
and o3.employment_status = 'active'
and o3.person_phone is not null
and o3.person_phone <> ''
and binary o3.person_phone <> binary u.phonenumber
)
and o.position_name is not null
and o.position_name <> ''
order by o.position_name
@@ -2,17 +2,33 @@ package org.dromara.aihr.broadcast;
import org.dromara.aihr.broadcast.AihrBroadcastDto.BroadcastAttachmentResponse;
import org.dromara.aihr.broadcast.AihrBroadcastDto.PublishRequest;
import org.dromara.aihr.knowledge.parse.KnowledgeDocumentParser;
import org.dromara.aihr.knowledge.service.AihrKnowledgePrincipalResolver;
import org.dromara.aihr.service.AihrSopSeedService;
import org.dromara.system.service.ISysOssService;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.RowMapper;
import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
import java.nio.file.Files;
import java.nio.file.Path;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ScheduledExecutorService;
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@Tag("dev")
class AihrBroadcastAttachmentContractTest {
@@ -58,5 +74,90 @@ class AihrBroadcastAttachmentContractTest {
assertTrue(source.contains("insight_status = 'PENDING'"));
assertTrue(source.contains("insight_status = 'FAILED'"));
assertTrue(source.contains("insight_status = 'PROCESSING'"));
assertTrue(source.contains("set insight_status = 'PENDING', insight_version = null"));
}
@Test
void insightWorkerUsesItsClaimTokenWhenCompletingWithoutChangingFileReadiness() throws Exception {
InsightJdbcTemplate jdbcTemplate = new InsightJdbcTemplate();
AihrBroadcastInsightService insightService = mock(AihrBroadcastInsightService.class);
when(insightService.generate(any(), any())).thenReturn(
new AihrBroadcastInsightService.Generation("FAILED", null, List.of()));
AihrBroadcastAttachmentService service = new AihrBroadcastAttachmentService(
jdbcTemplate,
mock(ISysOssService.class),
mock(KnowledgeDocumentParser.class),
mock(AihrSopSeedService.class),
insightService,
mock(AihrKnowledgePrincipalResolver.class),
mock(AihrBroadcastService.class),
mock(ScheduledExecutorService.class)
);
Method claim = AihrBroadcastAttachmentService.class.getDeclaredMethod("claimInsight");
claim.setAccessible(true);
Object work = claim.invoke(service);
Method claimTokenAccessor = work.getClass().getDeclaredMethod("claimToken");
claimTokenAccessor.setAccessible(true);
String claimToken = (String) claimTokenAccessor.invoke(work);
Method process = AihrBroadcastAttachmentService.class.getDeclaredMethod("processInsight", work.getClass());
process.setAccessible(true);
process.invoke(service, work);
assertTrue(claimToken.startsWith("claim-"));
assertArrayEquals(new Object[]{claimToken, "tenant-a", 9L}, jdbcTemplate.updateArgs.get(0));
assertTrue(jdbcTemplate.updateSql.get(1).contains("and insight_version = ?"));
assertFalse(jdbcTemplate.updateSql.get(1).contains("set status ="));
assertEquals(claimToken, jdbcTemplate.updateArgs.get(1)[5]);
}
private static final class InsightJdbcTemplate extends JdbcTemplate {
private final List<String> updateSql = new ArrayList<>();
private final List<Object[]> updateArgs = new ArrayList<>();
@Override
public <T> List<T> query(String sql, RowMapper<T> rowMapper, Object... args) {
return rows(rowMapper);
}
@Override
public <T> List<T> query(String sql, RowMapper<T> rowMapper) {
return rows(rowMapper);
}
private static <T> List<T> rows(RowMapper<T> rowMapper) {
try {
return List.of(rowMapper.mapRow(resultSet(Map.of(
"id", 9L,
"tenant_id", "tenant-a",
"file_name", "制度.txt",
"extracted_text", "制度原文"
)), 0));
} catch (SQLException error) {
throw new AssertionError(error);
}
}
@Override
public int update(String sql, Object... args) {
updateSql.add(sql);
updateArgs.add(args);
return 1;
}
private static ResultSet resultSet(Map<String, Object> values) {
return (ResultSet) Proxy.newProxyInstance(
InsightJdbcTemplate.class.getClassLoader(), new Class<?>[]{ResultSet.class}, (proxy, method, args) -> {
Object value = values.get(args[0]);
if ("getLong".equals(method.getName())) {
return ((Number) value).longValue();
}
if ("getString".equals(method.getName())) {
return value == null ? null : value.toString();
}
throw new UnsupportedOperationException(method.getName());
});
}
}
}
@@ -39,22 +39,21 @@ class AihrBroadcastInsightServiceTest {
}
@Test
void fabricatedEvidenceIsRemovedBeforePersistence() {
String source = "财务部应在每月五日前完成费用复核并向项目负责人提交差异说明。";
void everyConclusionRequiresItsOwnExactEvidenceAndUsesTheActualParagraph() {
String source = "本通知自发布之日起执行。\n\n财务部应在每月五日前完成费用复核并向项目负责人提交差异说明。";
AihrModelSeedService model = mock(AihrModelSeedService.class);
when(model.chat(any())).thenReturn(success("""
{
"perspectives": [{
"code": "finance",
"summary": "财务需要按月完成费用复核。",
"concerns": ["费用复核时限"],
"impacts": ["需形成差异说明"],
"actions": ["每月五日前完成复核"],
"risks": ["逾期影响项目核算"],
"evidence": [
{"paragraphIndex": 1, "quote": "财务部应在每月五日前完成费用复核并向项目负责人提交差异说明"},
{"paragraphIndex": 1, "quote": "模型编造但原文不存在的财务处罚规则"}
]
"summary": {"text": "财务需要按月完成费用复核。",
"evidenceQuote": "财务部应在每月五日前完成费用复核并向项目负责人提交差异说明"},
"concerns": [{"text": "逾期将被处罚",
"evidenceQuote": "模型编造但原文不存在的财务处罚规则"}],
"impacts": [{"text": "需形成差异说明",
"evidenceQuote": "财务部应在每月五日前完成费用复核并向项目负责人提交差异说明"}],
"actions": [],
"risks": []
}]
}
"""));
@@ -65,19 +64,41 @@ class AihrBroadcastInsightServiceTest {
assertEquals("READY", result.status());
assertEquals(1, result.perspectives().size());
assertEquals(1, result.perspectives().get(0).evidence().size());
assertEquals(2, result.perspectives().get(0).evidence().get(0).paragraphIndex());
assertTrue(result.perspectives().get(0).concerns().isEmpty());
assertEquals(List.of("需形成差异说明"), result.perspectives().get(0).impacts());
assertFalse(result.json().contains("模型编造"));
assertTrue(source.contains(result.perspectives().get(0).evidence().get(0).quote()));
}
@Test
void evidenceIsCheckedAgainstTheSameRedactedTextSentToTheModel() {
String source = "财务复核遇到疑问时请联系13912345678,并在每月五日前提交书面差异说明。";
AihrModelSeedService model = mock(AihrModelSeedService.class);
when(model.chat(any())).thenReturn(success("""
{"perspectives":[{"code":"finance",
"summary":{"text":"按时提交差异说明",
"evidenceQuote":"财务复核遇到疑问时请联系1391****5678,并在每月五日前提交书面差异说明。"},
"concerns":[],"impacts":[],"actions":[],"risks":[]}]}
"""));
AihrBroadcastInsightService service = new AihrBroadcastInsightService(new ObjectMapper(), model);
AihrBroadcastInsightService.Generation result = service.generate("费用通知.txt", source);
assertEquals("READY", result.status());
assertTrue(result.json().contains("1391****5678"));
assertFalse(result.json().contains("13912345678"));
}
@Test
void oneFailedChunkProducesPartialResultAndAllFailuresProduceFailed() {
String source = "财务部需要复核本月项目费用并形成书面差异说明。" + "甲".repeat(4_100);
AihrModelSeedService partialModel = mock(AihrModelSeedService.class);
when(partialModel.chat(any()))
.thenReturn(success("""
{"perspectives":[{"code":"finance","summary":"复核费用","concerns":[],"impacts":[],
"actions":[],"risks":[],"evidence":[{"paragraphIndex":1,
"quote":"财务部需要复核本月项目费用并形成书面差异说明"}]}]}
{"perspectives":[{"code":"finance",
"summary":{"text":"复核费用","evidenceQuote":"财务部需要复核本月项目费用并形成书面差异说明"},
"concerns":[],"impacts":[],"actions":[],"risks":[]}]}
"""))
.thenThrow(new IllegalStateException("provider timeout"));
AihrBroadcastInsightService partial = new AihrBroadcastInsightService(new ObjectMapper(), partialModel);
@@ -70,6 +70,22 @@ class AihrBroadcastServiceTest {
assertArrayEquals(new Object[]{"tenant-a", 7L, "app_user"}, jdbcTemplate.queryForObjectArgs);
}
@Test
void employeeBroadcastEndpointsFailClosedWhenOnePhoneMatchesMultiplePeople() {
RecordingJdbcTemplate jdbcTemplate = new RecordingJdbcTemplate();
jdbcTemplate.audienceCount = 2L;
AihrBroadcastService service = new AihrBroadcastService(jdbcTemplate, RESOLVER);
ServiceException error = assertThrows(ServiceException.class, service::unreadCount);
assertEquals(HttpStatus.FORBIDDEN, error.getCode());
assertTrue(jdbcTemplate.queryForObjectSql.contains("count(distinct o.ext_party_id)"));
assertTrue(jdbcTemplate.queryForObjectSql.contains("not exists"));
assertTrue(jdbcTemplate.queryForObjectSql.contains("binary o2.person_phone <> binary u.phonenumber"));
assertTrue(jdbcTemplate.queryForObjectSql.contains("char_length(u.phonenumber) = 11"));
assertTrue(jdbcTemplate.queryForObjectSql.contains("u.phonenumber regexp '^[0-9]{11}$'"));
}
@Test
void unavailableMessagesDoNotLeakAcrossTenantsOrAfterWithdrawal() {
RecordingJdbcTemplate jdbcTemplate = new RecordingJdbcTemplate();