fix: harden knowledge lifecycle rollout

This commit is contained in:
key
2026-08-03 16:00:45 +08:00
parent df77998e0c
commit 99e8c509d5
18 changed files with 709 additions and 156 deletions
@@ -9,6 +9,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.springframework.transaction.support.TransactionTemplate;
import java.util.List;
@@ -19,6 +20,7 @@ public class AihrKnowledgeIndexOutboxService {
private final JdbcTemplate jdbcTemplate;
private final AihrSopSeedService sopSeedService;
private final TransactionTemplate transactionTemplate;
@Value("${aihr.qdrant.governed-collection:${AIHR_QDRANT_GOVERNED_COLLECTION:aihr_knowledge_governed_v1}}")
private String productionCollection = AihrKnowledgeRolloutService.DEFAULT_GOVERNED_COLLECTION;
@@ -102,15 +104,18 @@ public class AihrKnowledgeIndexOutboxService {
return;
}
try {
Integer vectors = TenantHelper.dynamic(row.tenantId(), () -> {
AihrSopSeedService.VectorizationResult vectorization = TenantHelper.dynamic(row.tenantId(), () -> {
if ("DELETE".equals(row.operation())) {
sopSeedService.deleteVectorDocument(row.knowledgeId(), row.docId(),
row.collection(), row.generation());
return null;
}
return sopSeedService.vectorizePublishedDocument(row.knowledgeId(), row.docId(),
return sopSeedService.vectorizePublishedDocumentWithManifest(row.knowledgeId(), row.docId(),
row.versionId(), row.collection(), row.generation());
});
if ("UPSERT".equals(row.operation()) && vectorization != null && vectorization.count() > 0) {
recordGenerationVectorization(row, vectorization);
}
jdbcTemplate.update("""
update aihr_index_outbox
set status = 'SUCCEEDED', last_error = null, update_time = now()
@@ -121,51 +126,9 @@ public class AihrKnowledgeIndexOutboxService {
set index_status = ?, update_time = now()
where tenant_id = ? and id = ?
and coalesce(published_version_id, current_version_id) = ?
""", "DELETE".equals(row.operation()) || vectors == null || vectors == 0 ? "NOT_REQUIRED" : "READY",
""", "DELETE".equals(row.operation()) || vectorization == null || vectorization.count() == 0
? "NOT_REQUIRED" : "READY",
row.tenantId(), row.assetId(), row.versionId());
if ("UPSERT".equals(row.operation()) && vectors != null && vectors > 0) {
jdbcTemplate.update("""
insert into aihr_knowledge_generation
(tenant_id, knowledge_id, generation, collection_name, status,
version_count, point_count, create_time, ready_time, activated_time)
values (?, ?, ?, ?, ?, 0, 0, now(), ?, ?)
on duplicate key update collection_name = values(collection_name),
status = case when status = 'FAILED' then 'BUILDING' else status end,
ready_time = coalesce(ready_time, values(ready_time)),
activated_time = coalesce(activated_time, values(activated_time)), last_error = null
""", row.tenantId(), row.knowledgeId(), row.generation(), row.collection(),
row.generation() == row.activeGeneration() ? "ACTIVE" : "BUILDING",
row.generation() == row.activeGeneration() ? java.time.LocalDateTime.now() : null,
row.generation() == row.activeGeneration() ? java.time.LocalDateTime.now() : null);
jdbcTemplate.update("""
insert into aihr_generation_version
(tenant_id, knowledge_id, generation, asset_id, version_id, content_sha256,
point_count, status, create_time)
select a.tenant_id, a.knowledge_id, ?, a.id, ?, v.content_sha256, ?, 'INCLUDED', now()
from aihr_data_asset a
join aihr_data_version v on v.tenant_id = a.tenant_id and v.id = ?
where a.tenant_id = ? and a.id = ?
on duplicate key update version_id = values(version_id),
content_sha256 = values(content_sha256), point_count = values(point_count), status = 'INCLUDED'
""", row.generation(), row.versionId(), vectors, row.versionId(),
row.tenantId(), row.assetId());
jdbcTemplate.update("""
update aihr_knowledge_generation generation_row
set version_count = (select count(*) from aihr_generation_version member
where member.tenant_id = generation_row.tenant_id
and member.knowledge_id = generation_row.knowledge_id
and member.generation = generation_row.generation
and member.status = 'INCLUDED'),
point_count = (select coalesce(sum(member.point_count), 0)
from aihr_generation_version member
where member.tenant_id = generation_row.tenant_id
and member.knowledge_id = generation_row.knowledge_id
and member.generation = generation_row.generation
and member.status = 'INCLUDED')
where generation_row.tenant_id = ? and generation_row.knowledge_id = ?
and generation_row.generation = ?
""", row.tenantId(), row.knowledgeId(), row.generation());
}
} catch (RuntimeException ex) {
String message = ex.getMessage() == null ? ex.getClass().getSimpleName() : ex.getMessage();
int attempted = row.attemptCount() + 1;
@@ -200,6 +163,84 @@ public class AihrKnowledgeIndexOutboxService {
}
}
private void recordGenerationVectorization(OutboxRow row,
AihrSopSeedService.VectorizationResult vectorization) {
String mismatch = transactionTemplate.execute(status -> {
List<GenerationManifest> rows = jdbcTemplate.query("""
select embedding_model, embedding_dimension, status, last_error
from aihr_knowledge_generation
where tenant_id = ? and knowledge_id = ? and generation = ? for update
""", (rs, rowNum) -> new GenerationManifest(rs.getString("embedding_model"),
rs.getObject("embedding_dimension", Integer.class), rs.getString("status"),
rs.getString("last_error")), row.tenantId(), row.knowledgeId(), row.generation());
GenerationManifest existing = rows.isEmpty() ? null : rows.get(0);
if (existing != null && existing.lastError() != null
&& existing.lastError().startsWith("EMBEDDING_MANIFEST_MISMATCH")) {
return existing.lastError();
}
if (existing != null && !manifestCompatible(existing.embeddingModel(), existing.embeddingDimension(),
vectorization.embeddingModel(), vectorization.embeddingDimension())) {
String reason = "EMBEDDING_MANIFEST_MISMATCH expected=" + existing.embeddingModel() + "/"
+ existing.embeddingDimension() + " actual=" + vectorization.embeddingModel() + "/"
+ vectorization.embeddingDimension();
jdbcTemplate.update("""
update aihr_knowledge_generation set status = 'FAILED', last_error = ?, ready_time = null
where tenant_id = ? and knowledge_id = ? and generation = ?
""", reason, row.tenantId(), row.knowledgeId(), row.generation());
return reason;
}
jdbcTemplate.update("""
insert into aihr_knowledge_generation
(tenant_id, knowledge_id, generation, collection_name, status,
version_count, point_count, embedding_model, embedding_dimension,
create_time, ready_time, activated_time)
values (?, ?, ?, ?, ?, 0, 0, ?, ?, now(), ?, ?)
on duplicate key update collection_name = values(collection_name),
embedding_model = coalesce(embedding_model, values(embedding_model)),
embedding_dimension = coalesce(embedding_dimension, values(embedding_dimension)),
status = case when status = 'FAILED' then 'BUILDING' else status end,
ready_time = coalesce(ready_time, values(ready_time)),
activated_time = coalesce(activated_time, values(activated_time)), last_error = null
""", row.tenantId(), row.knowledgeId(), row.generation(), row.collection(),
row.generation() == row.activeGeneration() ? "ACTIVE" : "BUILDING",
vectorization.embeddingModel(), vectorization.embeddingDimension(),
row.generation() == row.activeGeneration() ? java.time.LocalDateTime.now() : null,
row.generation() == row.activeGeneration() ? java.time.LocalDateTime.now() : null);
jdbcTemplate.update("""
insert into aihr_generation_version
(tenant_id, knowledge_id, generation, asset_id, version_id, content_sha256,
point_count, status, create_time)
select a.tenant_id, a.knowledge_id, ?, a.id, ?, v.content_sha256, ?, 'INCLUDED', now()
from aihr_data_asset a
join aihr_data_version v on v.tenant_id = a.tenant_id and v.id = ?
where a.tenant_id = ? and a.id = ?
on duplicate key update version_id = values(version_id),
content_sha256 = values(content_sha256), point_count = values(point_count), status = 'INCLUDED'
""", row.generation(), row.versionId(), vectorization.count(), row.versionId(),
row.tenantId(), row.assetId());
jdbcTemplate.update("""
update aihr_knowledge_generation generation_row
set version_count = (select count(*) from aihr_generation_version member
where member.tenant_id = generation_row.tenant_id
and member.knowledge_id = generation_row.knowledge_id
and member.generation = generation_row.generation
and member.status = 'INCLUDED'),
point_count = (select coalesce(sum(member.point_count), 0)
from aihr_generation_version member
where member.tenant_id = generation_row.tenant_id
and member.knowledge_id = generation_row.knowledge_id
and member.generation = generation_row.generation
and member.status = 'INCLUDED')
where generation_row.tenant_id = ? and generation_row.knowledge_id = ?
and generation_row.generation = ?
""", row.tenantId(), row.knowledgeId(), row.generation());
return null;
});
if (mismatch != null) {
throw new IllegalStateException(mismatch);
}
}
public OutboxHealth health(String tenantId) {
OutboxCounts counts = jdbcTemplate.queryForObject("""
select sum(case when status = 'PENDING' then 1 else 0 end) pending_count,
@@ -240,6 +281,15 @@ public class AihrKnowledgeIndexOutboxService {
|| ("DELETE".equals(operation) && !"DEPRECATED".equals(lifecycleStatus));
}
static boolean manifestCompatible(String expectedModel, Integer expectedDimension,
String actualModel, int actualDimension) {
if (actualModel == null || actualModel.isBlank() || actualDimension <= 0) {
return false;
}
return (expectedModel == null || expectedModel.isBlank() || expectedModel.equals(actualModel))
&& (expectedDimension == null || expectedDimension <= 0 || expectedDimension == actualDimension);
}
public record OutboxHealth(int pending, int processing, int failed, int deadLetter,
long oldestPendingSeconds, int staleAssets, boolean vectorMatched,
Integer mysqlFragments, Integer ungovernedFragments, Long qdrantPoints,
@@ -255,4 +305,8 @@ public class AihrKnowledgeIndexOutboxService {
String docId, String lifecycleStatus,
long currentVersionId, long activeGeneration) {
}
private record GenerationManifest(String embeddingModel, Integer embeddingDimension,
String status, String lastError) {
}
}
@@ -15,10 +15,12 @@ import org.springframework.jdbc.support.GeneratedKeyHolder;
import org.springframework.jdbc.support.KeyHolder;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.support.TransactionTemplate;
import java.sql.PreparedStatement;
import java.sql.Statement;
import java.time.LocalDateTime;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Locale;
import java.util.Map;
@@ -33,9 +35,9 @@ public class AihrKnowledgeRetentionService {
private final JdbcTemplate jdbcTemplate;
private final ObjectMapper objectMapper;
private final AihrKnowledgeLifecycleService lifecycleService;
private final AihrSopSeedService sopSeedService;
private final ISysOssService ossService;
private final TransactionTemplate transactionTemplate;
public ImpactAnalysis impact(String tenantId, long assetId) {
AssetRetention asset = requireAsset(tenantId, assetId, false);
@@ -62,17 +64,28 @@ public class AihrKnowledgeRetentionService {
int generations = count("""
select count(*) from aihr_generation_version where tenant_id = ? and asset_id = ?
""", tenantId, assetId);
int activeGenerationMemberships = count("""
select count(*)
from aihr_generation_version member
where member.tenant_id = ? and member.asset_id = ? and member.status = 'INCLUDED'
and member.generation = coalesce((
select scope.active_generation from aihr_knowledge_rollout_scope scope
where scope.tenant_id = member.tenant_id and scope.knowledge_id = member.knowledge_id
limit 1), 1)
""", tenantId, assetId);
boolean canRequestPurge = canRequestPurge(asset.lifecycleStatus(), asset.retentionClass(),
asset.legalHold(), learningTasks, broadcasts);
asset.legalHold(), learningTasks, broadcasts, activeGenerationMemberships);
List<String> blockers = new java.util.ArrayList<>();
if (asset.legalHold()) blockers.add("LEGAL_HOLD");
if ("PERMANENT".equals(asset.retentionClass())) blockers.add("PERMANENT_RETENTION");
if ("PUBLISHED".equals(asset.lifecycleStatus())) blockers.add("WITHDRAW_FIRST");
if (learningTasks > 0) blockers.add("LEARNING_TASK_REFERENCE");
if (broadcasts > 0) blockers.add("BROADCAST_REFERENCE");
if (activeGenerationMemberships > 0) blockers.add("ACTIVE_GENERATION_REFERENCE");
return new ImpactAnalysis(assetId, asset.lifecycleStatus(), asset.retentionClass(), asset.legalHold(),
asset.attachmentId(), asset.rawOssId(), fragments, publishedFragments, queryEvidence,
learningTasks, broadcasts, generations, canRequestPurge, List.copyOf(blockers));
learningTasks, broadcasts, generations, activeGenerationMemberships,
canRequestPurge, List.copyOf(blockers));
}
@Transactional
@@ -145,13 +158,52 @@ public class AihrKnowledgeRetentionService {
return requireRequest(tenantId, requestId, false);
}
@Transactional
public DeletionRequest execute(String tenantId, long requestId, long executorId) {
if (executorId <= 0) throw new ServiceException("A signed-in executor is required", HttpStatus.BAD_REQUEST);
ExecutionPlan plan = transactionTemplate.execute(status -> prepareExecution(tenantId, requestId, executorId));
if (plan == null) {
throw new ServiceException("Failed to prepare physical purge", HttpStatus.CONFLICT);
}
if (plan.completed()) {
return plan.request();
}
try {
TenantHelper.dynamic(tenantId, () -> {
sopSeedService.deleteLegacyVectorDocument(plan.asset().knowledgeId(), plan.asset().docId());
for (VectorDeletionTarget target : plan.vectorTargets()) {
sopSeedService.deleteVectorDocument(plan.asset().knowledgeId(), plan.asset().docId(),
target.collection(), target.generation());
}
return null;
});
if (plan.asset().rawOssId() != null && strictCount(
"select count(*) from aihr_knowledge_attach where tenant_id = ? and oss_id = ?",
tenantId, plan.asset().rawOssId()) <= 1 && strictCount(
"select count(*) from sys_oss where tenant_id = ? and oss_id = ?",
tenantId, plan.asset().rawOssId()) > 0) {
ossService.deleteWithValidByIds(List.of(plan.asset().rawOssId()), true);
}
transactionTemplate.executeWithoutResult(status -> completePurge(tenantId, requestId));
} catch (RuntimeException ex) {
transactionTemplate.executeWithoutResult(status -> markPurgeFailed(tenantId, requestId, ex));
throw new ServiceException("Physical purge failed; the request remains retryable", 503)
.setDetailMessage(ex.getMessage());
}
return requireRequest(tenantId, requestId, false);
}
private ExecutionPlan prepareExecution(String tenantId, long requestId, long executorId) {
DeletionRequest request = requireRequest(tenantId, requestId, true);
if ("EXECUTED".equals(request.status())) return request;
if (!Set.of("APPROVED", "FAILED").contains(request.status())) {
throw new ServiceException("Purge request is not approved", HttpStatus.CONFLICT);
if ("EXECUTED".equals(request.status())) {
return new ExecutionPlan(request, null, List.of(), true);
}
boolean staleExecution = "EXECUTING".equals(request.status()) && strictCount("""
select count(*) from aihr_deletion_request
where tenant_id = ? and id = ? and status = 'EXECUTING'
and update_time < date_sub(now(), interval 10 minute)
""", tenantId, requestId) > 0;
if (!Set.of("APPROVED", "FAILED").contains(request.status()) && !staleExecution) {
throw new ServiceException("Purge request is not approved or is already executing", HttpStatus.CONFLICT);
}
if (request.requesterId() == executorId) {
throw new ServiceException("Requester cannot execute the irreversible purge", HttpStatus.CONFLICT);
@@ -163,36 +215,69 @@ public class AihrKnowledgeRetentionService {
if (asset.legalHold() || "PUBLISHED".equals(asset.lifecycleStatus())) {
throw new ServiceException("Legal hold or active publication blocks purge", HttpStatus.CONFLICT);
}
jdbcTemplate.update("""
update aihr_deletion_request set status = 'EXECUTING', executed_by = ?, last_error = null, update_time = now()
where tenant_id = ? and id = ? and status in ('APPROVED','FAILED')
""", executorId, tenantId, requestId);
try {
AihrKnowledgeRolloutService.ModeGeneration target = rolloutTarget(tenantId, asset.knowledgeId());
TenantHelper.dynamic(tenantId, () -> {
sopSeedService.deleteVectorDocument(asset.knowledgeId(), asset.docId(),
target.governedCollection(), target.activeGeneration());
return null;
});
if (asset.rawOssId() != null && count(
"select count(*) from aihr_knowledge_attach where tenant_id = ? and oss_id = ?",
tenantId, asset.rawOssId()) <= 1) {
ossService.deleteWithValidByIds(List.of(asset.rawOssId()), true);
}
purgeDatabase(tenantId, asset);
jdbcTemplate.update("""
update aihr_deletion_request
set status = 'EXECUTED', executed_time = now(), last_error = null, update_time = now()
where tenant_id = ? and id = ?
""", tenantId, requestId);
} catch (RuntimeException ex) {
jdbcTemplate.update("""
update aihr_deletion_request set status = 'FAILED', last_error = left(?, 1000), update_time = now()
where tenant_id = ? and id = ?
""", ex.getMessage(), tenantId, requestId);
return requireRequest(tenantId, requestId, false);
if (strictCount("""
select count(*)
from aihr_generation_version member
where member.tenant_id = ? and member.asset_id = ? and member.status = 'INCLUDED'
and member.generation = coalesce((
select scope.active_generation from aihr_knowledge_rollout_scope scope
where scope.tenant_id = member.tenant_id and scope.knowledge_id = member.knowledge_id
limit 1), 1)
""", tenantId, asset.assetId()) > 0) {
throw new ServiceException(
"Active vector generation still contains this asset; activate a replacement generation first",
HttpStatus.CONFLICT);
}
return requireRequest(tenantId, requestId, false);
jdbcTemplate.update("""
update aihr_deletion_request
set status = 'EXECUTING', executed_by = ?, last_error = null, update_time = now()
where tenant_id = ? and id = ?
""", executorId, tenantId, requestId);
LinkedHashSet<VectorDeletionTarget> targets = new LinkedHashSet<>(jdbcTemplate.query("""
select distinct vector_target.collection_name, vector_target.generation
from (
select g.collection_name, member.generation
from aihr_generation_version member
join aihr_knowledge_generation g on g.tenant_id = member.tenant_id
and g.knowledge_id = member.knowledge_id and g.generation = member.generation
where member.tenant_id = ? and member.asset_id = ?
union all
select target_collection, index_generation
from aihr_index_outbox
where tenant_id = ? and asset_id = ?
) vector_target
order by vector_target.generation
""", (rs, rowNum) -> new VectorDeletionTarget(rs.getString("collection_name"),
rs.getLong("generation")), tenantId, asset.assetId(), tenantId, asset.assetId()));
AihrKnowledgeRolloutService.ModeGeneration active = rolloutTarget(tenantId, asset.knowledgeId());
targets.add(new VectorDeletionTarget(active.governedCollection(), active.activeGeneration()));
return new ExecutionPlan(request, asset, List.copyOf(targets), false);
}
private void completePurge(String tenantId, long requestId) {
DeletionRequest request = requireRequest(tenantId, requestId, true);
if ("EXECUTED".equals(request.status())) {
return;
}
if (!"EXECUTING".equals(request.status())) {
throw new ServiceException("Purge request execution state changed", HttpStatus.CONFLICT);
}
AssetRetention asset = requireAsset(tenantId, request.assetId(), true);
purgeDatabase(tenantId, asset);
jdbcTemplate.update("""
update aihr_deletion_request
set status = 'EXECUTED', executed_time = now(), last_error = null, update_time = now()
where tenant_id = ? and id = ? and status = 'EXECUTING'
""", tenantId, requestId);
}
private void markPurgeFailed(String tenantId, long requestId, RuntimeException error) {
String message = error.getMessage() == null ? error.getClass().getSimpleName() : error.getMessage();
jdbcTemplate.update("""
update aihr_deletion_request
set status = 'FAILED', last_error = left(?, 1000), update_time = now()
where tenant_id = ? and id = ? and status = 'EXECUTING'
""", message, tenantId, requestId);
}
public List<DeletionRequest> requests(String tenantId, String status, int limit) {
@@ -269,40 +354,74 @@ public class AihrKnowledgeRetentionService {
}
private void purgeDatabase(String tenantId, AssetRetention asset) {
safeUpdate("delete from aihr_query_evidence where tenant_id = ? and asset_id = ?",
List<GenerationRef> affectedGenerations = jdbcTemplate.query("""
select distinct knowledge_id, generation
from aihr_generation_version
where tenant_id = ? and asset_id = ?
order by knowledge_id, generation
""", (rs, rowNum) -> new GenerationRef(rs.getLong("knowledge_id"), rs.getLong("generation")),
tenantId, asset.assetId());
safeUpdate("delete from aihr_knowledge_fragment_locator where tenant_id = ? and doc_id = ?",
jdbcTemplate.update("delete from aihr_query_evidence where tenant_id = ? and asset_id = ?",
tenantId, asset.assetId());
jdbcTemplate.update("delete from aihr_knowledge_fragment_locator where tenant_id = ? and doc_id = ?",
tenantId, asset.docId());
safeUpdate("delete from aihr_knowledge_fragment where tenant_id = ? and knowledge_id = ? and doc_id = ?",
jdbcTemplate.update("delete from aihr_knowledge_fragment where tenant_id = ? and knowledge_id = ? and doc_id = ?",
tenantId, asset.knowledgeId(), asset.docId());
safeUpdate("""
jdbcTemplate.update("""
delete from aihr_claim_conflict where tenant_id = ? and (
left_claim_id in (select id from aihr_knowledge_claim where tenant_id = ? and asset_id = ?)
or right_claim_id in (select id from aihr_knowledge_claim where tenant_id = ? and asset_id = ?))
""", tenantId, tenantId, asset.assetId(), tenantId, asset.assetId());
safeUpdate("""
jdbcTemplate.update("""
delete from aihr_duplicate_relation where tenant_id = ?
and (left_asset_id = ? or right_asset_id = ?)
""", tenantId, asset.assetId(), asset.assetId());
safeUpdate("delete from aihr_review_batch_item where tenant_id = ? and asset_id = ?",
jdbcTemplate.update("delete from aihr_review_batch_item where tenant_id = ? and asset_id = ?",
tenantId, asset.assetId());
safeUpdate("""
jdbcTemplate.update("""
delete from aihr_dataset_membership where tenant_id = ? and version_id in
(select id from aihr_data_version where tenant_id = ? and asset_id = ?)
""", tenantId, tenantId, asset.assetId());
safeUpdate("delete from aihr_quality_assessment where tenant_id = ? and asset_id = ?",
jdbcTemplate.update("delete from aihr_quality_assessment where tenant_id = ? and asset_id = ?",
tenantId, asset.assetId());
for (String table : List.of("aihr_generation_version", "aihr_index_outbox", "aihr_quality_issue",
"aihr_review_assistance", "aihr_version_diff", "aihr_version_rollback")) {
safeUpdate("delete from " + table + " where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
jdbcTemplate.update("delete from " + table + " where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
}
safeUpdate("delete from aihr_knowledge_claim where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
safeUpdate("delete from aihr_review_decision where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
safeUpdate("delete from aihr_chunk_revision where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
safeUpdate("delete from aihr_data_version where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
safeUpdate("delete from aihr_data_asset where tenant_id = ? and id = ?", tenantId, asset.assetId());
for (GenerationRef generation : affectedGenerations) {
jdbcTemplate.update("""
update aihr_knowledge_generation generation_row
set version_count = (select count(*) from aihr_generation_version member
where member.tenant_id = generation_row.tenant_id
and member.knowledge_id = generation_row.knowledge_id
and member.generation = generation_row.generation
and member.status = 'INCLUDED'),
point_count = (select coalesce(sum(member.point_count), 0)
from aihr_generation_version member
where member.tenant_id = generation_row.tenant_id
and member.knowledge_id = generation_row.knowledge_id
and member.generation = generation_row.generation
and member.status = 'INCLUDED'),
status = 'FAILED', manifest_sha256 = null,
last_error = concat('INVALIDATED_BY_ASSET_PURGE:', ?)
where generation_row.tenant_id = ? and generation_row.knowledge_id = ?
and generation_row.generation = ? and generation_row.status <> 'ACTIVE'
""", asset.assetId(), tenantId, generation.knowledgeId(), generation.generation());
jdbcTemplate.update("""
update aihr_knowledge_rollout_scope
set candidate_generation = null,
reason = concat('Candidate invalidated by physical purge of asset ', ?),
update_time = now()
where tenant_id = ? and knowledge_id = ? and candidate_generation = ?
""", asset.assetId(), tenantId, generation.knowledgeId(), generation.generation());
}
jdbcTemplate.update("delete from aihr_knowledge_claim where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
jdbcTemplate.update("delete from aihr_review_decision where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
jdbcTemplate.update("delete from aihr_chunk_revision where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
jdbcTemplate.update("delete from aihr_data_version where tenant_id = ? and asset_id = ?", tenantId, asset.assetId());
jdbcTemplate.update("delete from aihr_data_asset where tenant_id = ? and id = ?", tenantId, asset.assetId());
if (asset.attachmentId() != null) {
safeUpdate("delete from aihr_knowledge_attach where tenant_id = ? and id = ?",
jdbcTemplate.update("delete from aihr_knowledge_attach where tenant_id = ? and id = ?",
tenantId, asset.attachmentId());
}
}
@@ -371,12 +490,9 @@ public class AihrKnowledgeRetentionService {
}
}
private void safeUpdate(String sql, Object... args) {
try {
jdbcTemplate.update(sql, args);
} catch (DataAccessException ignored) {
// Optional consumers can be absent in an incremental deployment; reconciliation records residue.
}
private int strictCount(String sql, Object... args) {
Integer value = jdbcTemplate.queryForObject(sql, Integer.class, args);
return value == null ? 0 : value;
}
private String json(Object value) {
@@ -398,9 +514,10 @@ public class AihrKnowledgeRetentionService {
}
static boolean canRequestPurge(String lifecycleStatus, String retentionClass, boolean legalHold,
int learningTasks, int broadcasts) {
int learningTasks, int broadcasts, int activeGenerationMemberships) {
return !legalHold && !"PERMANENT".equals(normalize(retentionClass))
&& !"PUBLISHED".equals(normalize(lifecycleStatus)) && learningTasks == 0 && broadcasts == 0;
&& !"PUBLISHED".equals(normalize(lifecycleStatus)) && learningTasks == 0 && broadcasts == 0
&& activeGenerationMemberships == 0;
}
private static String normalize(String value) {
@@ -420,7 +537,8 @@ public class AihrKnowledgeRetentionService {
public record ImpactAnalysis(long assetId, String lifecycleStatus, String retentionClass,
boolean legalHold, Long attachmentId, Long rawOssId, int chunks,
int publishedFragments, int queryEvidence, int learningTasks,
int broadcasts, int generations, boolean canRequestPurge,
int broadcasts, int generations, int activeGenerationMemberships,
boolean canRequestPurge,
List<String> blockers) {
}
@@ -447,4 +565,14 @@ public class AihrKnowledgeRetentionService {
Long rawOssId, String lifecycleStatus, Long publishedVersionId,
String retentionClass, boolean legalHold) {
}
private record VectorDeletionTarget(String collection, long generation) {
}
private record GenerationRef(long knowledgeId, long generation) {
}
private record ExecutionPlan(DeletionRequest request, AssetRetention asset,
List<VectorDeletionTarget> vectorTargets, boolean completed) {
}
}
@@ -3,8 +3,10 @@ package org.dromara.aihr.knowledge.quality;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor;
import org.dromara.aihr.service.AihrSopSeedService;
import org.dromara.common.core.constant.HttpStatus;
import org.dromara.common.core.exception.ServiceException;
import org.dromara.common.tenant.helper.TenantHelper;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.support.GeneratedKeyHolder;
import org.springframework.jdbc.support.KeyHolder;
@@ -29,6 +31,7 @@ public class AihrKnowledgeShadowRolloutService {
private final JdbcTemplate jdbcTemplate;
private final ObjectMapper objectMapper;
private final AihrSopSeedService sopSeedService;
@Transactional
public GenerationBuild startBuild(String tenantId, long knowledgeId, long reviewerId, BuildCommand command) {
@@ -73,7 +76,8 @@ public class AihrKnowledgeShadowRolloutService {
set candidate_generation = ?, updated_by = ?, reason = ?, update_time = now()
where tenant_id = ? and knowledge_id = ?
""", generation, reviewerId, command.reason().trim(), tenantId, knowledgeId);
return new GenerationBuild(knowledgeId, generation, collection, "BUILDING", jobs, 0, 0);
return new GenerationBuild(knowledgeId, generation, collection, "BUILDING", jobs, 0, 0,
null, null, null);
}
@Transactional
@@ -104,16 +108,28 @@ public class AihrKnowledgeShadowRolloutService {
where o.tenant_id = ? and a.knowledge_id = ? and o.index_generation = ?
and o.status in ('PENDING','PROCESSING')
""", tenantId, knowledgeId, generation);
String status = failed > 0 ? "FAILED" : (pending == 0 && included == expected ? "READY" : "BUILDING");
GenerationManifest manifest = generationManifest(tenantId, knowledgeId, generation);
boolean manifestReady = manifest != null && manifest.embeddingModel() != null
&& !manifest.embeddingModel().isBlank() && manifest.embeddingDimension() != null
&& manifest.embeddingDimension() > 0;
String failureReason = failed > 0 ? "INDEX_OUTBOX_FAILURE"
: (pending == 0 && expected == 0 ? "EMPTY_GENERATION"
: (pending == 0 && included == expected && !manifestReady ? "EMBEDDING_MANIFEST_MISSING" : null));
String status = failureReason != null ? "FAILED"
: (pending == 0 && included == expected ? "READY" : "BUILDING");
String manifestSha256 = "READY".equals(status)
? generationManifestSha256(tenantId, knowledgeId, generation, manifest) : null;
jdbcTemplate.update("""
update aihr_knowledge_generation
set status = ?, ready_time = case when ? = 'READY' then coalesce(ready_time, now()) else ready_time end,
last_error = case when ? = 'FAILED' then 'INDEX_OUTBOX_FAILURE' else null end
manifest_sha256 = case when ? = 'READY' then ? else manifest_sha256 end,
last_error = ?
where tenant_id = ? and knowledge_id = ? and generation = ?
""", status, status, status, tenantId, knowledgeId, generation);
""", status, status, status, manifestSha256, failureReason, tenantId, knowledgeId, generation);
GenerationBuild result = buildStatus(tenantId, knowledgeId, generation);
return new GenerationBuild(result.knowledgeId(), result.generation(), result.collection(), status,
expected, included, failed + pending);
expected, included, failed + pending, result.embeddingModel(), result.embeddingDimension(),
result.manifestSha256());
}
@Transactional
@@ -125,26 +141,24 @@ public class AihrKnowledgeShadowRolloutService {
if (!"SHADOW".equals(scope.mode()) || scope.candidateGeneration() == null) {
throw new ServiceException("Shadow comparison requires a candidate generation", HttpStatus.CONFLICT);
}
GenerationManifest manifest = generationManifest(tenantId, knowledgeId, scope.candidateGeneration());
if (manifest == null || !"READY".equals(manifest.status())) {
throw new ServiceException("Shadow comparison requires a READY candidate generation",
HttpStatus.CONFLICT);
}
if (manifest.embeddingModel() == null || manifest.embeddingModel().isBlank()
|| manifest.embeddingDimension() == null || manifest.embeddingDimension() <= 0) {
throw new ServiceException("Candidate generation embedding manifest is incomplete",
HttpStatus.CONFLICT);
}
String query = command.queryText().trim().substring(0, Math.min(1000, command.queryText().trim().length()));
long started = System.nanoTime();
List<Long> legacy = jdbcTemplate.queryForList("""
select f.id from aihr_knowledge_fragment f
join aihr_knowledge_attach a on a.tenant_id = f.tenant_id
and a.knowledge_id = f.knowledge_id and a.doc_id = f.doc_id and a.status = 2
where f.tenant_id = ? and f.knowledge_id = ? and f.content like concat('%', ?, '%')
order by f.id limit 5
""", Long.class, tenantId, knowledgeId, query);
List<Long> governed = jdbcTemplate.queryForList("""
select c.published_fragment_id
from aihr_chunk_revision c
join aihr_data_asset a on a.tenant_id = c.tenant_id and a.id = c.asset_id
and a.published_version_id = c.version_id
join aihr_dataset_membership d on d.tenant_id = c.tenant_id and d.version_id = c.version_id
and d.dataset_code = 'production' and d.status = 'ACTIVE'
where c.tenant_id = ? and a.knowledge_id = ? and a.lifecycle_status = 'PUBLISHED'
and c.published_fragment_id is not null and c.content like concat('%', ?, '%')
order by c.chunk_index, c.id limit 5
""", Long.class, tenantId, knowledgeId, query);
List<Long> legacy = TenantHelper.dynamic(tenantId,
() -> sopSeedService.shadowServingFragmentIds(knowledgeId, query, 5));
List<Long> governed = TenantHelper.dynamic(tenantId,
() -> sopSeedService.shadowCandidateFragmentIds(knowledgeId, query, 5,
manifest.collection(), scope.candidateGeneration(), manifest.embeddingModel(),
manifest.embeddingDimension()));
long latencyMs = Math.max(0, (System.nanoTime() - started) / 1_000_000L);
boolean top1 = !legacy.isEmpty() && !governed.isEmpty() && legacy.get(0).equals(governed.get(0));
Set<Long> intersection = new LinkedHashSet<>(legacy);
@@ -305,10 +319,13 @@ public class AihrKnowledgeShadowRolloutService {
update aihr_knowledge_generation set status = 'RETIRED', retire_after = now()
where tenant_id = ? and knowledge_id = ? and generation = ?
""", tenantId, knowledgeId, window.toGeneration());
jdbcTemplate.update("""
int restored = jdbcTemplate.update("""
update aihr_knowledge_generation set status = 'ACTIVE', retire_after = null
where tenant_id = ? and knowledge_id = ? and generation = ?
where tenant_id = ? and knowledge_id = ? and generation = ? and status = 'RETIRED'
""", tenantId, knowledgeId, window.fromGeneration());
if (restored != 1) {
throw new ServiceException("Rollback generation is no longer restorable", HttpStatus.CONFLICT);
}
jdbcTemplate.update("""
update aihr_knowledge_rollout_scope
set mode = 'SHADOW', active_generation = ?, candidate_generation = ?,
@@ -329,15 +346,42 @@ public class AihrKnowledgeShadowRolloutService {
private GenerationBuild buildStatus(String tenantId, long knowledgeId, long generation) {
List<GenerationBuild> rows = jdbcTemplate.query("""
select generation, collection_name, status, version_count, point_count
select generation, collection_name, status, version_count, point_count,
embedding_model, embedding_dimension, manifest_sha256
from aihr_knowledge_generation where tenant_id = ? and knowledge_id = ? and generation = ?
""", (rs, rowNum) -> new GenerationBuild(knowledgeId, rs.getLong("generation"),
rs.getString("collection_name"), rs.getString("status"), rs.getInt("version_count"),
rs.getInt("point_count"), 0), tenantId, knowledgeId, generation);
rs.getInt("point_count"), 0, rs.getString("embedding_model"),
rs.getObject("embedding_dimension", Integer.class), rs.getString("manifest_sha256")),
tenantId, knowledgeId, generation);
if (rows.isEmpty()) throw new ServiceException("Generation does not exist", HttpStatus.NOT_FOUND);
return rows.get(0);
}
private GenerationManifest generationManifest(String tenantId, long knowledgeId, long generation) {
List<GenerationManifest> rows = jdbcTemplate.query("""
select collection_name, status, embedding_model, embedding_dimension
from aihr_knowledge_generation
where tenant_id = ? and knowledge_id = ? and generation = ? limit 1
""", (rs, rowNum) -> new GenerationManifest(rs.getString("collection_name"),
rs.getString("status"), rs.getString("embedding_model"),
rs.getObject("embedding_dimension", Integer.class)), tenantId, knowledgeId, generation);
return rows.isEmpty() ? null : rows.get(0);
}
private String generationManifestSha256(String tenantId, long knowledgeId, long generation,
GenerationManifest manifest) {
List<String> members = jdbcTemplate.queryForList("""
select concat(asset_id, ':', version_id, ':', content_sha256, ':', point_count)
from aihr_generation_version
where tenant_id = ? and knowledge_id = ? and generation = ? and status = 'INCLUDED'
order by asset_id
""", String.class, tenantId, knowledgeId, generation);
return AihrKnowledgeLifecycleService.sha256(String.join("\n", List.of(
manifest.collection(), String.valueOf(generation), manifest.embeddingModel(),
String.valueOf(manifest.embeddingDimension()), String.join("\n", members))));
}
private RolloutTarget lockScope(String tenantId, long knowledgeId) {
return requireScope(tenantId, knowledgeId, true);
}
@@ -392,7 +436,8 @@ public class AihrKnowledgeShadowRolloutService {
}
public record GenerationBuild(long knowledgeId, long generation, String collection, String status,
int expectedVersions, int includedVersions, int outstandingJobs) {
int expectedVersions, int includedVersions, int outstandingJobs,
String embeddingModel, Integer embeddingDimension, String manifestSha256) {
}
public record ShadowComparison(long knowledgeId, long generation, String querySha256,
@@ -422,6 +467,10 @@ public class AihrKnowledgeShadowRolloutService {
Long candidateGeneration) {
}
private record GenerationManifest(String collection, String status, String embeddingModel,
Integer embeddingDimension) {
}
private record GateMetrics(int sampleCount, double top1Rate, double overlapAvg, double emptyRate) {
}
@@ -195,12 +195,14 @@ public class AihrKnowledgeVersionGovernanceService {
public List<GenerationSummary> generations(String tenantId, long knowledgeId) {
return jdbcTemplate.query("""
select generation, collection_name, base_generation, status, version_count, point_count,
manifest_sha256, create_time, ready_time, activated_time, retire_after, last_error
embedding_model, embedding_dimension, manifest_sha256,
create_time, ready_time, activated_time, retire_after, last_error
from aihr_knowledge_generation
where tenant_id = ? and knowledge_id = ? order by generation desc
""", (rs, rowNum) -> new GenerationSummary(rs.getLong("generation"),
rs.getString("collection_name"), rs.getObject("base_generation", Long.class),
rs.getString("status"), rs.getInt("version_count"), rs.getInt("point_count"),
rs.getString("embedding_model"), rs.getObject("embedding_dimension", Integer.class),
rs.getString("manifest_sha256"), rs.getObject("create_time", LocalDateTime.class),
rs.getObject("ready_time", LocalDateTime.class), rs.getObject("activated_time", LocalDateTime.class),
rs.getObject("retire_after", LocalDateTime.class), rs.getString("last_error")), tenantId, knowledgeId);
@@ -284,7 +286,8 @@ public class AihrKnowledgeVersionGovernanceService {
}
public record GenerationSummary(long generation, String collection, Long baseGeneration, String status,
int versionCount, int pointCount, String manifestSha256,
int versionCount, int pointCount, String embeddingModel,
Integer embeddingDimension, String manifestSha256,
LocalDateTime createTime, LocalDateTime readyTime,
LocalDateTime activatedTime, LocalDateTime retireAfter, String lastError) {
}
@@ -1,5 +1,6 @@
package org.dromara.aihr.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
@@ -2652,12 +2653,20 @@ public class AihrSopSeedService {
public int vectorizePublishedDocument(long knowledgeId, String docId) {
AihrKnowledgeRolloutService.ModeGeneration rollout = rolloutModeGeneration(knowledgeId);
return vectorizePublishedDocument(knowledgeId, docId, null,
firstNonBlank(rollout.governedCollection(), governedQdrantCollection()), rollout.activeGeneration());
return vectorizePublishedDocumentWithManifest(knowledgeId, docId, null,
firstNonBlank(rollout.governedCollection(), governedQdrantCollection()), rollout.activeGeneration()).count();
}
public int vectorizePublishedDocument(long knowledgeId, String docId, Long expectedVersionId,
String targetCollection, long targetGeneration) {
return vectorizePublishedDocumentWithManifest(knowledgeId, docId, expectedVersionId,
targetCollection, targetGeneration).count();
}
public VectorizationResult vectorizePublishedDocumentWithManifest(long knowledgeId, String docId,
Long expectedVersionId,
String targetCollection,
long targetGeneration) {
List<DocumentVectorRow> rows = jdbcTemplate.query("""
select i.name, f.content
from aihr_knowledge_fragment f
@@ -2677,8 +2686,9 @@ public class AihrSopSeedService {
List<String> fragments = rows.stream().map(DocumentVectorRow::content).toList();
EmbeddingData data = generateEmbeddings(fragments);
if (data.embeddings().isEmpty()) {
return 0;
throw new ServiceException("Published document embedding returned no vectors", 503);
}
int dimension = requireEmbeddingDimension(data.embeddings());
for (int index = 0; index < data.embeddings().size(); index++) {
jdbcTemplate.update("""
update aihr_knowledge_fragment
@@ -2694,7 +2704,7 @@ public class AihrSopSeedService {
} catch (Exception ex) {
throw new ServiceException("Vector index upsert failed").setDetailMessage(ex.getMessage());
}
return data.embeddings().size();
return new VectorizationResult(data.embeddings().size(), data.modelName(), dimension);
}
public void deleteVectorDocument(long knowledgeId, String docId) {
@@ -4565,6 +4575,132 @@ public class AihrSopSeedService {
}
}
public void deleteLegacyVectorDocument(long knowledgeId, String docId) {
deleteVectorDocument(knowledgeId, docId, legacyQdrantCollection(), null);
}
public List<Long> shadowServingFragmentIds(long knowledgeId, String queryText, int limit) {
int boundedLimit = Math.max(1, Math.min(limit, 20));
return dbHits("", queryText, boundedLimit, Set.of(knowledgeId)).stream()
.map(KnowledgeHit::fragmentId)
.filter(Objects::nonNull)
.toList();
}
public List<Long> shadowCandidateFragmentIds(long knowledgeId, String queryText, int limit,
String collection, long generation,
String embeddingModel, int embeddingDimension) {
if (knowledgeId <= 0 || isBlank(queryText) || generation <= 0 || isBlank(collection)
|| isBlank(embeddingModel) || embeddingDimension <= 0) {
throw new ServiceException("Candidate vector query manifest is incomplete", 409);
}
EmbeddingRuntime runtime = embeddingRuntimes().stream()
.filter(candidate -> embeddingModel.equals(candidate.modelName())
&& embeddingDimension == candidate.dimension())
.findFirst()
.orElseThrow(() -> new ServiceException(
"Candidate embedding model is no longer available", 409));
try {
List<String> embeddings = callEmbeddings(runtime, List.of(queryText));
if (embeddings.size() != 1 || requireEmbeddingDimension(embeddings) != embeddingDimension) {
throw new IllegalStateException("candidate query embedding dimension mismatch");
}
int boundedLimit = Math.max(1, Math.min(limit, 20));
ObjectNode body = objectMapper.createObjectNode();
body.set("query", objectMapper.readTree(embeddings.get(0)));
body.set("filter", qdrantProductionFilterForSpaces(Set.of(knowledgeId), null, null, generation));
body.put("limit", boundedLimit);
body.put("with_payload", true);
body.put("with_vector", false);
HttpResponse<String> response = qdrantRequest(
"POST", "/collections/" + collection + "/points/query", body);
if (!ok(response.statusCode())) {
throw new IllegalStateException("candidate qdrant query HTTP " + response.statusCode());
}
JsonNode points = objectMapper.readTree(response.body()).path("result").path("points");
if (!points.isArray()) {
throw new IllegalStateException("candidate qdrant response missing points");
}
List<Long> orderedFragmentIds = new ArrayList<>();
for (JsonNode point : points) {
JsonNode payload = point.path("payload");
long fragmentId = payload.path("fragment_id").asLong(0L);
if (fragmentId <= 0) {
throw new IllegalStateException("candidate qdrant point is missing fragment lineage");
}
orderedFragmentIds.add(fragmentId);
}
return validateCandidateFragmentLineage(knowledgeId, generation, orderedFragmentIds);
} catch (ServiceException ex) {
throw ex;
} catch (Exception ex) {
throw new ServiceException("Candidate vector query failed", 503).setDetailMessage(ex.getMessage());
}
}
private List<Long> validateCandidateFragmentLineage(long knowledgeId, long generation,
List<Long> orderedFragmentIds) {
if (orderedFragmentIds.isEmpty()) {
return List.of();
}
String placeholders = String.join(",", java.util.Collections.nCopies(orderedFragmentIds.size(), "?"));
List<Object> args = new ArrayList<>();
args.add(tenantId());
args.add(knowledgeId);
args.addAll(orderedFragmentIds);
Set<Long> authorized = new LinkedHashSet<>(jdbcTemplate.queryForList("""
select c.published_fragment_id
from aihr_chunk_revision c
join aihr_data_asset a on a.tenant_id = c.tenant_id and a.id = c.asset_id
and a.published_version_id = c.version_id
join aihr_dataset_membership d on d.tenant_id = c.tenant_id and d.version_id = c.version_id
and d.dataset_code = 'production' and d.status = 'ACTIVE'
join aihr_generation_version g on g.tenant_id = c.tenant_id
and g.knowledge_id = a.knowledge_id and g.generation = ?
and g.asset_id = a.id and g.version_id = c.version_id and g.status = 'INCLUDED'
where c.tenant_id = ? and a.knowledge_id = ?
and a.lifecycle_status = 'PUBLISHED' and a.trust_level = 'HUMAN_VERIFIED'
and (a.effective_from is null or a.effective_from <= current_date())
and (a.effective_to is null or a.effective_to >= current_date())
and c.published_fragment_id in (%s)
""".formatted(placeholders), Long.class,
prependGenerationArguments(args, generation).toArray()));
if (authorized.size() != new LinkedHashSet<>(orderedFragmentIds).size()) {
throw new ServiceException("Candidate vector lineage no longer matches the published generation", 409);
}
return orderedFragmentIds.stream().filter(authorized::contains).toList();
}
private static List<Object> prependGenerationArguments(List<Object> arguments, long generation) {
List<Object> ordered = new ArrayList<>();
ordered.add(generation);
ordered.addAll(arguments);
return ordered;
}
private int requireEmbeddingDimension(List<String> embeddings) {
Integer dimension = null;
for (String embedding : embeddings) {
try {
JsonNode vector = objectMapper.readTree(embedding);
if (!vector.isArray() || vector.isEmpty()) {
throw new IllegalStateException("embedding vector is empty or malformed");
}
if (dimension == null) {
dimension = vector.size();
} else if (dimension != vector.size()) {
throw new IllegalStateException("embedding vectors have inconsistent dimensions");
}
} catch (JsonProcessingException ex) {
throw new IllegalStateException("embedding vector is not valid JSON", ex);
}
}
if (dimension == null || dimension <= 0) {
throw new IllegalStateException("embedding vector dimension is unavailable");
}
return dimension;
}
private List<RolloutVectorScope> legacyVectorScopes(Set<Long> allowedKnowledgeIds) {
if (allowedKnowledgeIds != null && !allowedKnowledgeIds.isEmpty()) {
return allowedKnowledgeIds.stream().map(id -> new RolloutVectorScope(id, "LEGACY",
@@ -5357,6 +5493,9 @@ public class AihrSopSeedService {
private record EmbeddingData(String modelName, List<String> embeddings) {
}
public record VectorizationResult(int count, String embeddingModel, int embeddingDimension) {
}
private record VectorDbStats(int fragments, int ungovernedFragments, int embeddedFragments,
String embeddingModel, Integer embeddingDimension) {
}
@@ -281,7 +281,8 @@ class AihrDataQualityLifecycleSchemaTest {
assertThat(sql).contains("`aihr_shadow_comparison`", "top1_agreement", "overlap_at_5",
"`aihr_activation_gate`", "threshold_json", "evidence_json",
"`aihr_generation_activation`", "rollback_deadline", "ACTIVATE", "ROLLBACK");
"`aihr_generation_activation`", "rollback_deadline", "ACTIVATE", "ROLLBACK",
"embedding_model", "embedding_dimension");
assertThat(sql.toLowerCase()).doesNotContain("update aihr_knowledge_rollout_scope");
}
}
@@ -4,6 +4,8 @@ import org.dromara.common.core.exception.ServiceException;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import java.nio.file.Files;
import java.nio.file.Path;
import java.time.LocalDateTime;
import static org.assertj.core.api.Assertions.assertThat;
@@ -31,11 +33,12 @@ class AihrKnowledgeGovernanceSafetyTest {
@Test
void referencedOrActiveAssetsCannotEnterPhysicalPurge() {
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", false, 0, 0)).isTrue();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("PUBLISHED", "STANDARD", false, 0, 0)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", false, 1, 0)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", false, 0, 1)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", true, 0, 0)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", false, 0, 0, 0)).isTrue();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("PUBLISHED", "STANDARD", false, 0, 0, 0)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", false, 1, 0, 0)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", false, 0, 1, 0)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", true, 0, 0, 0)).isFalse();
assertThat(AihrKnowledgeRetentionService.canRequestPurge("DEPRECATED", "STANDARD", false, 0, 0, 1)).isFalse();
}
@Test
@@ -60,7 +63,25 @@ class AihrKnowledgeGovernanceSafetyTest {
existing.runKey(), null, "LOW_RISK", true, 50, null))).isTrue();
assertThat(AihrKnowledgeMigrationOrchestrationService.sameManifest(existing,
new AihrKnowledgeMigrationOrchestrationService.CreateRun(
existing.runKey(), null, "ALL", true, 50, null))).isFalse();
existing.runKey(), null, "ALL", true, 50, null))).isFalse();
}
@Test
void physicalPurgeUsesRetryableExternalCleanupAndAtomicDatabaseDeletion() throws Exception {
Path source = Path.of("src/main/java/org/dromara/aihr/knowledge/quality/AihrKnowledgeRetentionService.java");
if (!Files.exists(source)) {
source = Path.of("ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/knowledge/quality/AihrKnowledgeRetentionService.java");
}
String code = Files.readString(source);
assertThat(code).contains("prepareExecution", "completePurge", "markPurgeFailed",
"deleteLegacyVectorDocument", "VectorDeletionTarget", "transactionTemplate.executeWithoutResult",
"ACTIVE_GENERATION_REFERENCE", "INVALIDATED_BY_ASSET_PURGE", "manifest_sha256 = null");
assertThat(code).doesNotContain("private void safeUpdate");
Path rolloutSource = source.getParent().resolve("AihrKnowledgeShadowRolloutService.java");
String rolloutCode = Files.readString(rolloutSource);
assertThat(rolloutCode).contains("and status = 'RETIRED'", "Rollback generation is no longer restorable");
}
private static AihrKnowledgeMigrationOrchestrationService.MigrationRun migrationRun(
@@ -25,4 +25,19 @@ class AihrKnowledgeIndexOutboxServiceTest {
assertThat(AihrKnowledgeIndexOutboxService.isObsolete("UPSERT", 10, "PUBLISHED", 10)).isFalse();
assertThat(AihrKnowledgeIndexOutboxService.isObsolete("DELETE", 10, "DEPRECATED", 10)).isFalse();
}
@Test
@Tag("dev")
void generationManifestRejectsMixedModelsAndDimensions() {
assertThat(AihrKnowledgeIndexOutboxService.manifestCompatible(
null, null, "local-hash-v1", 256)).isTrue();
assertThat(AihrKnowledgeIndexOutboxService.manifestCompatible(
"local-hash-v1", 256, "local-hash-v1", 256)).isTrue();
assertThat(AihrKnowledgeIndexOutboxService.manifestCompatible(
"local-hash-v1", 256, "other-model", 256)).isFalse();
assertThat(AihrKnowledgeIndexOutboxService.manifestCompatible(
"local-hash-v1", 256, "local-hash-v1", 1024)).isFalse();
assertThat(AihrKnowledgeIndexOutboxService.manifestCompatible(
"local-hash-v1", 256, "", 0)).isFalse();
}
}
@@ -0,0 +1,86 @@
package org.dromara.aihr.knowledge.quality;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.dromara.aihr.service.AihrSopSeedService;
import org.dromara.common.core.exception.ServiceException;
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.sql.ResultSet;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.contains;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@Tag("dev")
class AihrKnowledgeShadowRolloutServiceTest {
@Test
void candidateGenerationMustBeReadyBeforeAComparison() throws Exception {
JdbcTemplate jdbc = rolloutJdbc("BUILDING");
AihrSopSeedService sop = mock(AihrSopSeedService.class);
AihrKnowledgeShadowRolloutService service = new AihrKnowledgeShadowRolloutService(
jdbc, new ObjectMapper(), sop);
ServiceException error = assertThrows(ServiceException.class,
() -> service.compare("000000", 10L,
new AihrKnowledgeShadowRolloutService.ComparisonCommand("收费标准", "MANUAL_SAMPLE")));
assertEquals("Shadow comparison requires a READY candidate generation", error.getMessage());
verifyNoInteractions(sop);
}
@Test
void strictCandidateVectorFailureCannotCreateActivationEvidence() throws Exception {
JdbcTemplate jdbc = rolloutJdbc("READY");
AihrSopSeedService sop = mock(AihrSopSeedService.class);
when(sop.shadowServingFragmentIds(10L, "收费标准", 5)).thenReturn(List.of(101L));
when(sop.shadowCandidateFragmentIds(10L, "收费标准", 5,
"aihr_knowledge_governed_v1", 2L, "local-hash-v1", 256))
.thenThrow(new ServiceException("Candidate vector query failed", 503));
AihrKnowledgeShadowRolloutService service = new AihrKnowledgeShadowRolloutService(
jdbc, new ObjectMapper(), sop);
assertThrows(ServiceException.class,
() -> service.compare("000000", 10L,
new AihrKnowledgeShadowRolloutService.ComparisonCommand("收费标准", "MANUAL_SAMPLE")));
verify(jdbc, never()).update(contains("insert into aihr_shadow_comparison"), any(Object[].class));
}
@SuppressWarnings({"rawtypes", "unchecked"})
private static JdbcTemplate rolloutJdbc(String generationStatus) throws Exception {
JdbcTemplate jdbc = mock(JdbcTemplate.class);
ResultSet scope = mock(ResultSet.class);
when(scope.getString("mode")).thenReturn("SHADOW");
when(scope.getString("governed_collection")).thenReturn("aihr_knowledge_governed_v1");
when(scope.getLong("active_generation")).thenReturn(1L);
when(scope.getObject("candidate_generation", Long.class)).thenReturn(2L);
when(jdbc.query(contains("aihr_knowledge_rollout_scope"), any(RowMapper.class), any(Object[].class)))
.thenAnswer(invocation -> {
RowMapper mapper = invocation.getArgument(1);
return List.of(mapper.mapRow(scope, 0));
});
ResultSet generation = mock(ResultSet.class);
when(generation.getString("collection_name")).thenReturn("aihr_knowledge_governed_v1");
when(generation.getString("status")).thenReturn(generationStatus);
when(generation.getString("embedding_model")).thenReturn("local-hash-v1");
when(generation.getObject("embedding_dimension", Integer.class)).thenReturn(256);
when(jdbc.query(contains("from aihr_knowledge_generation"), any(RowMapper.class), any(Object[].class)))
.thenAnswer(invocation -> {
RowMapper mapper = invocation.getArgument(1);
return List.of(mapper.mapRow(generation, 0));
});
return jdbc;
}
}
@@ -648,6 +648,10 @@ public class AihrSopSeedServiceTest {
assertTrue(code.contains("payload.put(\"trust_level\", \"HUMAN_VERIFIED\")"));
assertTrue(code.contains("payload.put(\"index_generation\", generation)"));
assertTrue(code.contains("qdrantProductionFilterForSpaces(knowledgeIds, null, category, generation)"));
assertTrue(code.contains("shadowCandidateFragmentIds"));
assertTrue(code.contains("qdrantProductionFilterForSpaces(Set.of(knowledgeId), null, null, generation)"));
assertTrue(code.contains("candidate qdrant point is missing fragment lineage"));
assertTrue(code.contains("Candidate vector query failed"));
assertTrue(code.contains("legacyQdrantCollection()"));
assertTrue(code.contains("governedQdrantCollection()"));
}
@@ -1,5 +1,25 @@
-- Shadow retrieval evidence and explicit generation activation gates.
SET @aihr_column_exists := (
SELECT COUNT(*) FROM information_schema.columns
WHERE table_schema = DATABASE() AND table_name = 'aihr_knowledge_generation'
AND column_name = 'embedding_model'
);
SET @aihr_shadow_ddl := IF(@aihr_column_exists = 0,
'ALTER TABLE `aihr_knowledge_generation` ADD COLUMN `embedding_model` varchar(100) DEFAULT NULL AFTER `point_count`',
'SELECT 1');
PREPARE aihr_shadow_stmt FROM @aihr_shadow_ddl; EXECUTE aihr_shadow_stmt; DEALLOCATE PREPARE aihr_shadow_stmt;
SET @aihr_column_exists := (
SELECT COUNT(*) FROM information_schema.columns
WHERE table_schema = DATABASE() AND table_name = 'aihr_knowledge_generation'
AND column_name = 'embedding_dimension'
);
SET @aihr_shadow_ddl := IF(@aihr_column_exists = 0,
'ALTER TABLE `aihr_knowledge_generation` ADD COLUMN `embedding_dimension` int DEFAULT NULL AFTER `embedding_model`',
'SELECT 1');
PREPARE aihr_shadow_stmt FROM @aihr_shadow_ddl; EXECUTE aihr_shadow_stmt; DEALLOCATE PREPARE aihr_shadow_stmt;
CREATE TABLE IF NOT EXISTS `aihr_shadow_comparison` (
`id` bigint NOT NULL AUTO_INCREMENT,
`tenant_id` varchar(20) NOT NULL,
@@ -78,3 +98,4 @@ SET @aihr_shadow_ddl := CONCAT('ALTER TABLE `aihr_generation_activation` CONVERT
PREPARE aihr_shadow_stmt FROM @aihr_shadow_ddl; EXECUTE aihr_shadow_stmt; DEALLOCATE PREPARE aihr_shadow_stmt;
SET @aihr_shadow_collation := NULL;
SET @aihr_shadow_ddl := NULL;
SET @aihr_column_exists := NULL;