diff --git a/conversation/backend/src/main/java/org/entcore/conversation/Conversation.java b/conversation/backend/src/main/java/org/entcore/conversation/Conversation.java index 7e8cf7ccfb..103a0ca717 100644 --- a/conversation/backend/src/main/java/org/entcore/conversation/Conversation.java +++ b/conversation/backend/src/main/java/org/entcore/conversation/Conversation.java @@ -32,6 +32,7 @@ import org.entcore.conversation.controllers.ApiController; import org.entcore.conversation.controllers.ConversationController; import org.entcore.conversation.controllers.TaskController; +import org.entcore.conversation.cron.PurgeMessages; import org.entcore.conversation.service.ConversationService; import org.entcore.conversation.service.impl.ConversationRepositoryEvents; import org.entcore.conversation.service.impl.ConversationStorage; @@ -96,8 +97,10 @@ public Future initConversation(StorageFactory storageFactory) { // Delete Orphans final String deleteOrphanCron = config.getString("deleteOrphanCron"); final DeleteOrphan deleteOrphan = new DeleteOrphan(storage); - // Enable delete orphan task to be triggered via API - addController(new TaskController(deleteOrphan)); + final String purgeMessagesCron = config.getString("purgeMessagesCron"); + final PurgeMessages purgeMessages = new PurgeMessages(conversationService); + // Enable delete orphan and purge old messages tasks to be triggered via API + addController(new TaskController(deleteOrphan, purgeMessages)); // Schedule delete orphan task from cron expression if (deleteOrphanCron != null) { try { @@ -106,6 +109,14 @@ public Future initConversation(StorageFactory storageFactory) { log.error("Invalid cron expression.", e); } } + // Schedule purge old messages task from cron expression + if (purgeMessagesCron != null) { + try { + new CronTrigger(vertx, purgeMessagesCron).schedule(purgeMessages); + } catch (ParseException e) { + log.error("Invalid purge messages cron expression.", e); + } + } return Future.succeededFuture(); } diff --git a/conversation/backend/src/main/java/org/entcore/conversation/controllers/ApiController.java b/conversation/backend/src/main/java/org/entcore/conversation/controllers/ApiController.java index 9f93589d63..f473fd5507 100644 --- a/conversation/backend/src/main/java/org/entcore/conversation/controllers/ApiController.java +++ b/conversation/backend/src/main/java/org/entcore/conversation/controllers/ApiController.java @@ -22,19 +22,13 @@ import java.util.List; import java.util.Optional; -import fr.wseduc.security.MfaProtected; -import io.vertx.core.json.JsonArray; - -import org.entcore.common.http.filter.IgnoreCsrf; import org.entcore.common.http.filter.ResourceFilter; -import static fr.wseduc.webutils.request.RequestUtils.bodyToJson; import static org.entcore.common.http.response.DefaultResponseHandler.arrayResponseHandler; import static org.entcore.common.http.response.DefaultResponseHandler.defaultResponseHandler; import static org.entcore.common.user.UserUtils.getAuthenticatedUserInfos; import static org.entcore.common.utils.StringUtils.isEmpty; -import org.entcore.common.http.filter.SuperAdminFilter; import org.entcore.conversation.filters.FoldersFilter; import org.entcore.conversation.filters.MessageUserFilter; import org.entcore.conversation.filters.SystemOrUserFolderFilter; @@ -192,35 +186,6 @@ public void count(final HttpServerRequest request) { }); } - @Get("api/purge/list") - @SecuredAction(value = "", type = ActionType.RESOURCE) - @ResourceFilter(SuperAdminFilter.class) - @MfaProtected() - public void purgeList(final HttpServerRequest request) { - conversationService.getMessagesToPurge() - .onSuccess( messages -> renderJson(request, messages)) - .onFailure( throwable -> renderJson(request, new JsonObject().put("error", throwable.getMessage()), 400)); - } - - @Post("api/purge/messages") - @SecuredAction(value = "", type = ActionType.RESOURCE) - @ResourceFilter(SuperAdminFilter.class) - @MfaProtected() - @IgnoreCsrf - public void purgeMessages(final HttpServerRequest request) { - bodyToJson(request, body -> { - JsonArray ids = body.getJsonArray("id"); - if (ids == null || ids.isEmpty()) { - badRequest(request); - return; - } - - conversationService.purgeMessages(ids.getList()) - .onSuccess(unused1 -> renderJson(request, new JsonObject().put("status", "ok"))) - .onFailure(throwable -> renderJson(request, new JsonObject().put("error", ((Throwable) throwable).getMessage()), 500)); - }); - } - /** Utility method to read a query param and convert it to an Integer. */ private Integer parseQueryParam(final HttpServerRequest request, String param, final Integer defaultValue) { final String paramValue = getOrElse(request.params().get(param), "" + defaultValue, false); diff --git a/conversation/backend/src/main/java/org/entcore/conversation/controllers/TaskController.java b/conversation/backend/src/main/java/org/entcore/conversation/controllers/TaskController.java index 426ffb68e4..cf471b4074 100644 --- a/conversation/backend/src/main/java/org/entcore/conversation/controllers/TaskController.java +++ b/conversation/backend/src/main/java/org/entcore/conversation/controllers/TaskController.java @@ -7,15 +7,18 @@ import io.vertx.core.http.HttpServerRequest; import io.vertx.core.impl.logging.Logger; import io.vertx.core.impl.logging.LoggerFactory; +import org.entcore.conversation.cron.PurgeMessages; import org.entcore.conversation.service.impl.DeleteOrphan; public class TaskController extends BaseController { protected static final Logger log = LoggerFactory.getLogger(TaskController.class); private final DeleteOrphan deleteOrphan; + private final PurgeMessages purgeMessages; - public TaskController(DeleteOrphan deleteOrphan) { + public TaskController(DeleteOrphan deleteOrphan, PurgeMessages purgeMessages) { this.deleteOrphan = deleteOrphan; + this.purgeMessages = purgeMessages; } @Post("api/internal/purge/orphans") @@ -25,4 +28,12 @@ public void deleteOrphans(final HttpServerRequest request) { deleteOrphan.handle(0L); render(request, null, 202); } + + @Post("api/internal/purge/messages") + @SecuredAction(value = "", type = ActionType.RESOURCE) + public void purgeMessages(final HttpServerRequest request) { + log.info("Triggered purge old messages task"); + this.purgeMessages.handle(0L); + render(request, null, 202); + } } diff --git a/conversation/backend/src/main/java/org/entcore/conversation/cron/PurgeMessages.java b/conversation/backend/src/main/java/org/entcore/conversation/cron/PurgeMessages.java new file mode 100644 index 0000000000..225b992851 --- /dev/null +++ b/conversation/backend/src/main/java/org/entcore/conversation/cron/PurgeMessages.java @@ -0,0 +1,30 @@ +package org.entcore.conversation.cron; + +import io.vertx.core.Handler; +import io.vertx.core.impl.logging.Logger; +import io.vertx.core.impl.logging.LoggerFactory; +import org.entcore.conversation.service.ConversationService; + +/** + * Tâche cron de purge des anciens messages de la messagerie. + * Peut être déclenchée par le CRON vertx (hors-Kube) ou par l'endpoint + * api/internal/purge/messages du TaskController (job Kube). + */ +public class PurgeMessages implements Handler { + + private static final Logger log = LoggerFactory.getLogger(PurgeMessages.class); + + private final ConversationService conversationService; + + public PurgeMessages(ConversationService conversationService) { + this.conversationService = conversationService; + } + + @Override + public void handle(Long event) { + log.info("Starting old messages purge process"); + conversationService.purgeMessages() + .onSuccess(v -> log.info("Old messages purge process completed successfully")) + .onFailure(err -> log.error("Old messages purge process failed: " + err.getMessage(), err)); + } +} diff --git a/conversation/backend/src/main/java/org/entcore/conversation/service/ConversationService.java b/conversation/backend/src/main/java/org/entcore/conversation/service/ConversationService.java index 701d12f095..62ddb4d6b4 100644 --- a/conversation/backend/src/main/java/org/entcore/conversation/service/ConversationService.java +++ b/conversation/backend/src/main/java/org/entcore/conversation/service/ConversationService.java @@ -201,6 +201,5 @@ void findVisibleRecipients(String parentMessageId, UserInfos user, void forwardAttachments(String forwardId, String messageId, UserInfos user, Handler> result); // Purge - Future getMessagesToPurge(); - Future purgeMessages(final List messagesId); + Future purgeMessages(); } diff --git a/conversation/backend/src/main/java/org/entcore/conversation/service/impl/SqlConversationService.java b/conversation/backend/src/main/java/org/entcore/conversation/service/impl/SqlConversationService.java index ea1017a8e0..076bd6c777 100644 --- a/conversation/backend/src/main/java/org/entcore/conversation/service/impl/SqlConversationService.java +++ b/conversation/backend/src/main/java/org/entcore/conversation/service/impl/SqlConversationService.java @@ -1958,52 +1958,60 @@ public void forwardAttachments(String forwardId, String messageId, UserInfos use /* Purge */ @Override - public Future getMessagesToPurge() { - final int months = Math.max(36, Config.getConf().getInteger("purge-grace-period", 0)); - final int timeout = Config.getConf().getInteger("purge-query-timeout", 300000); - - String query = - "SELECT DISTINCT um.message_id " + - "FROM conversation.usermessages um " + - "JOIN conversation.messages m ON um.message_id = m.id " + - "WHERE um.folder_id IS NULL AND m.date > (EXTRACT(EPOCH FROM NOW())::bigint - (" + months + " * 31 * 86400));"; - - DeliveryOptions deliveryOptions = new DeliveryOptions(); - deliveryOptions.setSendTimeout(timeout); - - Promise promise = Promise.promise(); - sql.prepared(query, new JsonArray(), deliveryOptions, SqlResult.validResultHandler(result -> { - if (result.isRight() && result.right().getValue() != null) { - promise.complete(result.right().getValue()); - } else { - log.error("An error occurred getting messages list to purge:", result.left().getValue()); - promise.fail("conversation.purge.list.error"); - } - })); - - return promise.future(); - } - - @Override - public Future purgeMessages(final List messagesId) { - final int timeout = Config.getConf().getInteger("purge-query-timeout", 300000); - - SqlStatementsBuilder builder = new SqlStatementsBuilder(); - JsonArray values = new JsonArray(); - builder.prepared("DELETE FROM conversation.usermessages WHERE message_id IN " + generateInVars(messagesId, values), values); + public Future purgeMessages() { + final JsonObject purgeConfig = Config.getConf().getJsonObject("purge-messages", new JsonObject()); + + final int months = Math.max(36, purgeConfig.getInteger("grace-period", 0)); + final int timeout = purgeConfig.getInteger("query-timeout", 300000); + final int batchSize = purgeConfig.getInteger("batch-size", 5000); + final int maxBatches = purgeConfig.getInteger("max-batches", 100); + + // Purge par lots des liens usermessages de messages âgés de plus de `months` mois + // et rangés dans aucun dossier (folder_id IS NULL). Chaque lot est une transaction + // autonome afin de relâcher les verrous entre les lots (cf. SUPPORT-4770). + final String query = + "DELETE FROM " + userMessageTable + " um " + + "USING " + messageTable + " m " + + "WHERE um.message_id = m.id " + + "AND um.folder_id IS NULL " + + "AND m.date < (EXTRACT(EPOCH FROM NOW())::bigint - (" + months + " * 31 * 86400)) * 1000 " + + "AND (um.user_id, um.message_id) IN (" + + "SELECT um2.user_id, um2.message_id FROM " + userMessageTable + " um2 " + + "JOIN " + messageTable + " m2 ON um2.message_id = m2.id " + + "WHERE um2.folder_id IS NULL " + + "AND m2.date < (EXTRACT(EPOCH FROM NOW())::bigint - (" + months + " * 31 * 86400)) * 1000 " + + "LIMIT ?);"; + + log.info("Starting old messages purge (retention=" + months + " months, batch-size=" + batchSize + ", max-batches=" + maxBatches + ")"); + return purgeMessagesBatch(query, batchSize, timeout, maxBatches, 0, 0); + } + + private Future purgeMessagesBatch(final String query, final int batchSize, final int timeout, final int maxBatches, final int batchNumber, final int totalDeleted) { + if (batchNumber >= maxBatches) { + log.info("Old messages purge stopped after reaching max batches (" + maxBatches + "), total deleted: " + totalDeleted + " usermessages"); + return Future.succeededFuture(); + } DeliveryOptions deliveryOptions = new DeliveryOptions(); deliveryOptions.setSendTimeout(timeout); - Promise promise = Promise.promise(); - sql.transaction(builder.build(), deliveryOptions, SqlResult.validResultsHandler(result -> { - if (result.isRight() && result.right().getValue() != null) { - promise.complete(result.right().getValue()); + Promise promise = Promise.promise(); + sql.prepared(query, new JsonArray().add(batchSize), deliveryOptions, result -> { + if ("ok".equals(result.body().getString("status"))) { + int deleted = result.body().getInteger("rows", 0); + if (deleted == 0) { + log.info("Old messages purge completed, total deleted: " + totalDeleted + " usermessages"); + promise.complete(); + } else { + log.info("Old messages purge batch " + batchNumber + " deleted " + deleted + " usermessages"); + purgeMessagesBatch(query, batchSize, timeout, maxBatches, batchNumber + 1, totalDeleted + deleted).onComplete(promise); + } } else { - log.error("An error occurred purging messages:", result.left().getValue()); - promise.fail("conversation.purge.error"); + String error = result.body().getString("message", "Unknown error"); + log.error("An error occurred purging old messages (batch " + batchNumber + "): " + error); + promise.fail(error); } - })); + }); return promise.future(); } diff --git a/conversation/backend/src/main/resources/template.j2 b/conversation/backend/src/main/resources/template.j2 index 453fce24a0..ad03175e23 100644 --- a/conversation/backend/src/main/resources/template.j2 +++ b/conversation/backend/src/main/resources/template.j2 @@ -29,6 +29,13 @@ "batch-size": {{ conversationDeleteOrphanBatchSize | default('500') }}, "max-batches": {{ conversationDeleteOrphanMaxBatches | default('10') }} }, + "purgeMessagesCron" : "{{ conversationPurgeMessagesCron | default('0 0 3 ? * SUN *') }}", + "purge-messages": { + "grace-period": {{ conversationPurgeGracePeriod | default('36') }}, + "query-timeout": {{ conversationPurgeQueryTimeout | default('300000') }}, + "batch-size": {{ conversationPurgeMessagesBatchSize | default('5000') }}, + "max-batches": {{ conversationPurgeMessagesMaxBatches | default('100') }} + }, {% endif %} "fileAnalyzer": { "enabled": {{ fileAnalyzerOn | default('true') }}, @@ -43,7 +50,6 @@ }, "sql": true, "db-schema": "conversation", - "purge-grace-period": {{ conversationPurgeGracePeriod | default('36') }}, {% if mailToExercizer is defined and mailToExercizer == 'true' %} "mail-to-exercizer": { "langfuse_public_key": "{{ langfuse_public_key }}",