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
Expand Up @@ -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;
Expand Down Expand Up @@ -96,8 +97,10 @@ public Future<Void> 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 {
Expand All @@ -106,6 +109,14 @@ public Future<Void> 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();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
@@ -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<Long> {

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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,5 @@ void findVisibleRecipients(String parentMessageId, UserInfos user,
void forwardAttachments(String forwardId, String messageId, UserInfos user, Handler<Either<String, JsonObject>> result);

// Purge
Future<JsonArray> getMessagesToPurge();
Future<JsonArray> purgeMessages(final List<String> messagesId);
Future<Void> purgeMessages();
}
Original file line number Diff line number Diff line change
Expand Up @@ -1958,52 +1958,60 @@
/* Purge */

@Override
public Future<JsonArray> 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<JsonArray> 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<JsonArray> purgeMessages(final List<String> 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<Void> 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<Void> 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");

Check failure on line 1991 in conversation/backend/src/main/java/org/entcore/conversation/service/impl/SqlConversationService.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Define a constant instead of duplicating this literal " usermessages" 3 times.

See more on https://sonarcloud.io/project/issues?id=edificeio_entcore&issues=AZ-TPk_5zGMUC9_w9yLZ&open=AZ-TPk_5zGMUC9_w9yLZ&pullRequest=1088
return Future.succeededFuture();
}

DeliveryOptions deliveryOptions = new DeliveryOptions();
deliveryOptions.setSendTimeout(timeout);

Promise<JsonArray> promise = Promise.promise();
sql.transaction(builder.build(), deliveryOptions, SqlResult.validResultsHandler(result -> {
if (result.isRight() && result.right().getValue() != null) {
promise.complete(result.right().getValue());
Promise<Void> 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();
}
Expand Down
8 changes: 7 additions & 1 deletion conversation/backend/src/main/resources/template.j2
Original file line number Diff line number Diff line change
Expand Up @@ -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') }},
Expand All @@ -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 }}",
Expand Down