diff --git a/minerva-api/src/main/java/io/minerva/api/queue/ConvertingResultConsumer.java b/minerva-api/src/main/java/io/minerva/api/queue/ConvertingResultConsumer.java index de82cde9..e3171275 100644 --- a/minerva-api/src/main/java/io/minerva/api/queue/ConvertingResultConsumer.java +++ b/minerva-api/src/main/java/io/minerva/api/queue/ConvertingResultConsumer.java @@ -1,10 +1,16 @@ package io.minerva.api.queue; -import io.minerva.api.model.data.node.*; +import io.minerva.api.model.data.node.Chunk; +import io.minerva.api.model.data.node.ChunkEmbedding; +import io.minerva.api.model.data.node.Document; +import io.minerva.api.model.data.node.DocumentPartition; import io.minerva.api.model.dto.document.DocumentChunkingDto; import io.minerva.api.model.status.DocumentStatus; import io.minerva.api.queue.model.ConvertingResult; -import io.minerva.api.repository.*; +import io.minerva.api.repository.ChunkEmbeddingRepository; +import io.minerva.api.repository.ChunkRepository; +import io.minerva.api.repository.DocumentRepository; +import io.minerva.api.repository.JobLockRepository; import io.minerva.api.service.*; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; @@ -143,24 +149,25 @@ public class ConvertingResultConsumer { private final DocumentSummarizationService documentSummarizationService; /** - * Repository for loading {@link Topic} nodes by ID. + * Repository for saving {@link ChunkEmbedding} nodes and linking them to chunks. */ - private final TopicRepository topicRepository; + private final ChunkEmbeddingRepository chunkEmbeddingRepository; /** - * Repository for loading {@link Clazz} nodes by ID. + * Service for generating AI-suggested questions from a document. */ - private final ClassRepository classRepository; + private final DocumentQuestionGenerationService documentQuestionGenerationService; /** - * Repository for saving {@link ChunkEmbedding} nodes and linking them to chunks. + * Repository used to acquire and release distributed locks that prevent duplicate topic/class description + * generation when multiple documents in the same topic or class complete concurrently. */ - private final ChunkEmbeddingRepository chunkEmbeddingRepository; + private final JobLockRepository jobLockRepository; /** - * Service for generating AI-suggested questions from a document. + * Number of minutes a topic- or class-level description generation lock is held before it is considered stale. */ - private final DocumentQuestionGenerationService documentQuestionGenerationService; + private static final int DESCRIPTION_LOCK_TIMEOUT_MINUTES = 10; /** * Ensures the consumer group exists and starts the background consumer loop. Invoked automatically by Spring after @@ -432,23 +439,7 @@ private void processSuccessfulResult(ConvertingResult result) { ); document.setParsedUrl(parsedUrl); } - document.setChunkCount(result.getNumChunks()); - document.setImageCount(result.getNumImages()); - document.setPageCount(result.getNumPages()); - document.setTableCount(result.getNumTables()); - document.setFormulaCount(result.getNumFormulas()); - - // Persist YouTube-specific metadata when available - if (result.getVideoTitle() != null && !result.getVideoTitle().isBlank()) { - document.setTitle(result.getVideoTitle()); - } - if (result.getThumbnailUrl() != null && !result.getThumbnailUrl().isBlank()) { - document.setThumbnailUrl(result.getThumbnailUrl()); - } - if (result.getVideoDuration() != null) { - document.setVideoDuration(result.getVideoDuration()); - } - + setDocumentMetadata(document, result); documentRepository.save(document); // Convert and filter chunks @@ -476,6 +467,16 @@ private void processSuccessfulResult(ConvertingResult result) { updatedDocument.setStatus(DocumentStatus.COMPLETED); documentRepository.save(updatedDocument); + // Trigger topic/class description generation once all documents in the topic are terminal. + // Runs asynchronously so the consumer loop is not blocked while waiting for LLM responses. + UUID completedTopicId = updatedDocument.getTopicId(); + UUID completedClassId = updatedDocument.getClassId(); + if (completedTopicId != null && completedClassId != null) { + CompletableFuture.runAsync( + () -> triggerDescriptionGenerationIfAllTerminal(completedTopicId, completedClassId) + ); + } + String reverseKey = ConvertingTaskPublisher.DOCUMENT_TASK_PREFIX + document.getId().toString(); Object taskId = redisTemplate.opsForValue().get(reverseKey); if (taskId != null) { @@ -485,22 +486,39 @@ private void processSuccessfulResult(ConvertingResult result) { log.info("Document {} processed successfully.", document.getId()); } + /** + * Enrich document node with metadata from {@link ConvertingResult} + * + * @param document The {@link Document} Node + * @param result {@link ConvertingResult} with Metadat + */ + private void setDocumentMetadata(Document document, ConvertingResult result) { + document.setChunkCount(result.getNumChunks()); + document.setImageCount(result.getNumImages()); + document.setPageCount(result.getNumPages()); + document.setTableCount(result.getNumTables()); + document.setFormulaCount(result.getNumFormulas()); + + // Persist YouTube-specific metadata when available + if (result.getVideoTitle() != null && !result.getVideoTitle().isBlank()) { + document.setTitle(result.getVideoTitle()); + } + if (result.getThumbnailUrl() != null && !result.getThumbnailUrl().isBlank()) { + document.setThumbnailUrl(result.getThumbnailUrl()); + } + if (result.getVideoDuration() != null) { + document.setVideoDuration(result.getVideoDuration()); + } + } + private CompletableFuture getSummarizationFuture(Document document, ConvertingResult result) { return CompletableFuture.runAsync( () -> { try { documentSummarizationService.summarizeDocument(document.getId(), result.getContent()); - Topic topic = topicRepository.findById(document.getTopicId()) - .orElseThrow(() -> new RuntimeException( - "Topic not found: " + document.getTopicId())); - topicDescriptionGenerationService.generateTopicDescription(topic.getId()); - Clazz clazz = classRepository.findById(document.getClassId()) - .orElseThrow(() -> new RuntimeException( - "Class not found: " + document.getClassId())); - classDescriptionGenerationService.generateClassDescription(clazz.getId()); } catch (Exception e) { log.error( - "Summarization/description generation failed for document {}: {}", + "Summarization failed for document {}: {}", document.getId(), e.getMessage(), e ); } @@ -508,6 +526,79 @@ private CompletableFuture getSummarizationFuture(Document document, Conver ); } + /** + * Triggers topic and class AI description generation if — and only if — every document in the same topic has + * reached a terminal state ({@code COMPLETED} or {@code FAILED}). + * + *

A {@link io.minerva.api.model.data.node.JobLock} is acquired for each level before generation begins, + * ensuring that exactly one API instance runs the generation even when multiple documents in the same topic finish + * concurrently. Each lock is released in a {@code finally} block. If the lock cannot be acquired (another instance + * is already generating), the call is silently skipped — the winning instance handles the work. + * + *

This method is called asynchronously (fire-and-forget) after the document status has been persisted, so + * the consumer loop is not blocked while waiting for LLM responses. + * + * @param topicId the ID of the topic whose document just reached a terminal state + * @param classId the ID of the class that owns the topic + */ + private void triggerDescriptionGenerationIfAllTerminal(UUID topicId, UUID classId) { + if (!documentRepository.areAllDocumentsInTopicTerminal(topicId)) { + log.debug( + "Skipping topic description generation for topic {} — not all documents are in terminal state", + topicId + ); + return; + } + + String topicLockName = "topic-desc:" + topicId; + boolean topicLockAcquired = false; + try { + topicLockAcquired = jobLockRepository.acquireLock(topicLockName, DESCRIPTION_LOCK_TIMEOUT_MINUTES); + if (topicLockAcquired) { + topicDescriptionGenerationService.generateTopicDescription(topicId); + } else { + log.debug( + "Topic description generation for topic {} already running on another instance, skipping", + topicId + ); + } + } catch (Exception e) { + log.error("Topic description generation failed for topic {}: {}", topicId, e.getMessage(), e); + } finally { + if (topicLockAcquired) { + jobLockRepository.releaseLock(topicLockName); + } + } + + if (!documentRepository.areAllDocumentsInClassTerminal(classId)) { + log.debug( + "Skipping class description generation for class {} — not all documents are in terminal state", + classId + ); + return; + } + + String classLockName = "class-desc:" + classId; + boolean classLockAcquired = false; + try { + classLockAcquired = jobLockRepository.acquireLock(classLockName, DESCRIPTION_LOCK_TIMEOUT_MINUTES); + if (classLockAcquired) { + classDescriptionGenerationService.generateClassDescription(classId); + } else { + log.debug( + "Class description generation for class {} already running on another instance, skipping", + classId + ); + } + } catch (Exception e) { + log.error("Class description generation failed for class {}: {}", classId, e.getMessage(), e); + } finally { + if (classLockAcquired) { + jobLockRepository.releaseLock(classLockName); + } + } + } + private CompletableFuture getEmbeddingFuture(Document document, List chunks) { return CompletableFuture.runAsync( () -> { @@ -544,7 +635,8 @@ private CompletableFuture getQuestionFuture(Document document, ConvertingR } /** - * Handles a failed conversion result by marking the document as {@link DocumentStatus#FAILED}. + * Handles a failed conversion result by marking the document as {@link DocumentStatus#FAILED} and, if all other + * documents in the same topic have also reached a terminal state, triggering topic/class description generation. * * @param result the failed conversion result; must not be {@code null} */ @@ -561,6 +653,16 @@ private void processFailedResult(ConvertingResult result) { document.setStatus(DocumentStatus.FAILED); documentRepository.save(document); + // Even when a document fails, description generation should still run once all documents in the topic + // have settled, so that partial content is not left without an AI description. + UUID failedTopicId = document.getTopicId(); + UUID failedClassId = document.getClassId(); + if (failedTopicId != null && failedClassId != null) { + CompletableFuture.runAsync( + () -> triggerDescriptionGenerationIfAllTerminal(failedTopicId, failedClassId) + ); + } + log.error("Document {} processing failed: {}", document.getId(), result.getError()); } diff --git a/minerva-api/src/main/java/io/minerva/api/repository/DocumentRepository.java b/minerva-api/src/main/java/io/minerva/api/repository/DocumentRepository.java index 9daa6be9..fc80e332 100644 --- a/minerva-api/src/main/java/io/minerva/api/repository/DocumentRepository.java +++ b/minerva-api/src/main/java/io/minerva/api/repository/DocumentRepository.java @@ -93,6 +93,44 @@ OPTIONAL MATCH (chunk)-[:HAS_CHUNK_EMBEDDING]->(chunkEmbedding:ChunkEmbedding) */ List findByUserId(@Param("userId") UUID userId); + /** + * Returns {@code true} when every document in the given topic has reached a terminal processing state + * ({@code COMPLETED} or {@code FAILED}) and at least one document exists. + * + *

Used by {@code ConvertingResultConsumer} to decide whether topic/class description generation + * should be triggered after a document reaches a terminal state. + * + * @param topicId the ID of the topic to check + * @return {@code true} if all documents are terminal and the topic is non-empty; {@code false} otherwise + */ + @Query( + """ + MATCH (d:Document {topicId: $topicId}) + WITH collect(d) AS docs + RETURN size(docs) > 0 AND all(doc IN docs WHERE doc.status IN ['COMPLETED', 'FAILED']) + """ + ) + boolean areAllDocumentsInTopicTerminal(@Param("topicId") UUID topicId); + + /** + * Returns {@code true} when every document across all topics in the given class has reached a terminal + * processing state ({@code COMPLETED} or {@code FAILED}) and at least one document exists. + * + *

Used by {@code ConvertingResultConsumer} to decide whether class description generation should be + * triggered after a topic's documents are all settled. + * + * @param classId the ID of the class to check + * @return {@code true} if all documents are terminal and the class is non-empty; {@code false} otherwise + */ + @Query( + """ + MATCH (d:Document {classId: $classId}) + WITH collect(d) AS docs + RETURN size(docs) > 0 AND all(doc IN docs WHERE doc.status IN ['COMPLETED', 'FAILED']) + """ + ) + boolean areAllDocumentsInClassTerminal(@Param("classId") UUID classId); + /** * Returns all documents for a topic together with their associated document summary (if any). *