Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand All @@ -485,29 +486,119 @@ 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<Void> 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
);
}
}
);
}

/**
* 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}).
*
* <p>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.
*
* <p>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<Void> getEmbeddingFuture(Document document, List<Chunk> chunks) {
return CompletableFuture.runAsync(
() -> {
Expand Down Expand Up @@ -544,7 +635,8 @@ private CompletableFuture<Void> 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}
*/
Expand All @@ -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());
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,44 @@ OPTIONAL MATCH (chunk)-[:HAS_CHUNK_EMBEDDING]->(chunkEmbedding:ChunkEmbedding)
*/
List<Document> 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.
*
* <p>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.
*
* <p>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).
*
Expand Down
Loading