chibenwa commented on code in PR #3193: URL: https://github.com/apache/james-project/pull/3193#discussion_r4171726693
########## docs/modules/servers/partials/architecture/blobstore.adoc: ########## @@ -67,6 +74,60 @@ encrypt afterwards; reads decrypt first and decompress afterwards. This ordering preserves the benefit of compression, as encrypted payloads are generally not compressible. +[NOTE] +Chunked storage relies on partial HTTP Range requests (`Range: bytes=offset-(offset+limit-1)`) +against the underlying object storage. Range-read access requires unencrypted chunk storage +or server-side encryption (SSE-S3 / SSE-KMS / SSE-C); it is incompatible with client-side +`AESBlobStoreDAO` encryption because random byte ranges in AES ciphertext cannot be decrypted +in isolation. If server-side encryption with customer-provided keys (SSE-C) is used, ensure Review Comment: This is false. AES cannot be appled after chunking but can be applied before. It is the order of composition of the Blob stores that matters. Indeed AES needs to be applied before (each blob independently encrypted) ########## docs/modules/servers/partials/operate/webadmin.adoc: ########## @@ -3176,6 +3176,104 @@ Where: filter in later runs. - *gcedBlobCount* is the count of blobs that were garbage collected. +== Running blob object compaction + +NOTE: S3 Object Compaction is available only on the *Distributed Server* deployment using S3/MinIO Object Storage, without client-side encryption or whole-blob compression. + +WARNING: Once compaction has created multi-slot chunk objects in S3, enabling client-side encryption (AES) in `blobstore.properties` is unsupported and will render previously compacted chunk blobs unreadable. + +In large installations backed by S3 object storage, storing large numbers of small +blobs can lead to high S3 request billing and slow bucket listings. Object compaction +allows administrators to pack standalone blobs from historical generations into large +chunk objects, serving individual blobs via HTTP ranged reads. + +=== Initial Compaction + +To compact standalone blobs of a completed generation into chunk objects: + +.... +curl -XDELETE "http://ip:port/blobs?action=initial-compaction&generation=1" Review Comment: The verb "DELETE" looks weird, how about "POST" ? ########## mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraBlobIdUpdater.java: ########## @@ -0,0 +1,231 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.mailbox.cassandra.mail; + +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.bindMarker; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.selectFrom; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.update; +import static com.datastax.oss.driver.api.querybuilder.relation.Relation.column; +import static com.datastax.oss.driver.api.querybuilder.update.Assignment.setColumn; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.IMAP_UID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MAILBOX_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MESSAGE_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.BODY_CONTENT; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.HEADER_CONTENT; + +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; +import java.util.function.Predicate; + +import jakarta.inject.Inject; + +import org.apache.james.backends.cassandra.init.configuration.JamesExecutionProfiles; +import org.apache.james.backends.cassandra.utils.CassandraAsyncExecutor; +import org.apache.james.blob.api.BlobId; +import org.apache.james.blob.api.BlobIdUpdater; +import org.apache.james.mailbox.cassandra.table.CassandraMessageIdTable; +import org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table; +import org.apache.james.mailbox.cassandra.table.MessageIdToImapUid; +import org.apache.james.util.ReactorUtils; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.type.codec.TypeCodecs; + +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +public class CassandraBlobIdUpdater implements BlobIdUpdater { + public static class Factory implements BlobIdUpdater.Factory { + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final BlobId.Factory blobIdFactory; + private final DriverExecutionProfile batchProfile; + private final PreparedStatement selectAll; + private final PreparedStatement selectMessageV3; + private final PreparedStatement updateMessageV3Header; + private final PreparedStatement updateMessageV3Body; + private final PreparedStatement selectImapUidByMessageId; + private final PreparedStatement updateImapUidHeader; + private final PreparedStatement updateMessageIdTableHeader; + + @Inject + public Factory(CqlSession session, BlobId.Factory blobIdFactory) { + this.cassandraAsyncExecutor = new CassandraAsyncExecutor(session); + this.blobIdFactory = blobIdFactory; + this.batchProfile = JamesExecutionProfiles.getBatchProfile(session); + + this.selectAll = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(MESSAGE_ID, HEADER_CONTENT, BODY_CONTENT) + .build()); + + this.selectMessageV3 = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(HEADER_CONTENT, BODY_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Header = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Body = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(BODY_CONTENT, bindMarker(BODY_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.selectImapUidByMessageId = session.prepare(selectFrom(MessageIdToImapUid.TABLE_NAME) + .columns(MAILBOX_ID, IMAP_UID, HEADER_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateImapUidHeader = session.prepare(update(MessageIdToImapUid.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID)), + column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + + this.updateMessageIdTableHeader = session.prepare(update(CassandraMessageIdTable.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + } + + @Override + public Mono<BlobIdUpdater> forPredicate(Predicate<BlobId> generationCondition, + Consumer<BlobId> referencedBlobIdObserver) { + Map<BlobId, Set<UUID>> references = new ConcurrentHashMap<>(); + return cassandraAsyncExecutor.executeRows(selectAll.bind().setExecutionProfile(batchProfile)) + .doOnNext(row -> { + UUID messageId = row.get(MESSAGE_ID, TypeCodecs.TIMEUUID); + String headerStr = row.get(HEADER_CONTENT, TypeCodecs.TEXT); + String bodyStr = row.get(BODY_CONTENT, TypeCodecs.TEXT); + if (headerStr != null) { + BlobId headerId = blobIdFactory.parse(headerStr); + referencedBlobIdObserver.accept(headerId); + if (generationCondition.test(headerId)) { + references.computeIfAbsent(headerId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + if (bodyStr != null) { + BlobId bodyId = blobIdFactory.parse(bodyStr); + referencedBlobIdObserver.accept(bodyId); + if (generationCondition.test(bodyId)) { + references.computeIfAbsent(bodyId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + }) + .then(Mono.fromCallable(() -> new CassandraBlobIdUpdater( + cassandraAsyncExecutor, references, + selectMessageV3, updateMessageV3Header, updateMessageV3Body, + selectImapUidByMessageId, updateImapUidHeader, updateMessageIdTableHeader))); + } + } + + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final Map<BlobId, Set<UUID>> references; Review Comment: This won't fit in memory. We could envision a separated table IMO to store this... ########## server/apps/migration/core-data-jpa-to-pg/src/main/java/org/apache/james/JpaToPgCoreDataMigration.java: ########## @@ -162,6 +163,7 @@ static List<Module> chooseBlobStoreModules(MigrationConfiguration configuration) public static List<Module> chooseModules(BlobStoreConfiguration choosingConfiguration) { return ImmutableList.<Module>builder() .add(chooseBlobStoreDAOModule(choosingConfiguration.getImplementation())) + .add(chooseChunkedBlobStoreDAOModule(choosingConfiguration.getImplementation())) Review Comment: IMO we should always choose "non chunking" for JPA ########## mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraBlobIdUpdater.java: ########## @@ -0,0 +1,231 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.mailbox.cassandra.mail; + +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.bindMarker; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.selectFrom; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.update; +import static com.datastax.oss.driver.api.querybuilder.relation.Relation.column; +import static com.datastax.oss.driver.api.querybuilder.update.Assignment.setColumn; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.IMAP_UID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MAILBOX_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MESSAGE_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.BODY_CONTENT; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.HEADER_CONTENT; + +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; +import java.util.function.Predicate; + +import jakarta.inject.Inject; + +import org.apache.james.backends.cassandra.init.configuration.JamesExecutionProfiles; +import org.apache.james.backends.cassandra.utils.CassandraAsyncExecutor; +import org.apache.james.blob.api.BlobId; +import org.apache.james.blob.api.BlobIdUpdater; +import org.apache.james.mailbox.cassandra.table.CassandraMessageIdTable; +import org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table; +import org.apache.james.mailbox.cassandra.table.MessageIdToImapUid; +import org.apache.james.util.ReactorUtils; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.type.codec.TypeCodecs; + +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +public class CassandraBlobIdUpdater implements BlobIdUpdater { + public static class Factory implements BlobIdUpdater.Factory { + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final BlobId.Factory blobIdFactory; + private final DriverExecutionProfile batchProfile; + private final PreparedStatement selectAll; + private final PreparedStatement selectMessageV3; + private final PreparedStatement updateMessageV3Header; + private final PreparedStatement updateMessageV3Body; + private final PreparedStatement selectImapUidByMessageId; + private final PreparedStatement updateImapUidHeader; + private final PreparedStatement updateMessageIdTableHeader; + + @Inject + public Factory(CqlSession session, BlobId.Factory blobIdFactory) { + this.cassandraAsyncExecutor = new CassandraAsyncExecutor(session); + this.blobIdFactory = blobIdFactory; + this.batchProfile = JamesExecutionProfiles.getBatchProfile(session); + + this.selectAll = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(MESSAGE_ID, HEADER_CONTENT, BODY_CONTENT) + .build()); + + this.selectMessageV3 = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(HEADER_CONTENT, BODY_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Header = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Body = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(BODY_CONTENT, bindMarker(BODY_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.selectImapUidByMessageId = session.prepare(selectFrom(MessageIdToImapUid.TABLE_NAME) + .columns(MAILBOX_ID, IMAP_UID, HEADER_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateImapUidHeader = session.prepare(update(MessageIdToImapUid.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID)), + column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + + this.updateMessageIdTableHeader = session.prepare(update(CassandraMessageIdTable.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + } + + @Override + public Mono<BlobIdUpdater> forPredicate(Predicate<BlobId> generationCondition, + Consumer<BlobId> referencedBlobIdObserver) { + Map<BlobId, Set<UUID>> references = new ConcurrentHashMap<>(); + return cassandraAsyncExecutor.executeRows(selectAll.bind().setExecutionProfile(batchProfile)) + .doOnNext(row -> { + UUID messageId = row.get(MESSAGE_ID, TypeCodecs.TIMEUUID); + String headerStr = row.get(HEADER_CONTENT, TypeCodecs.TEXT); + String bodyStr = row.get(BODY_CONTENT, TypeCodecs.TEXT); + if (headerStr != null) { + BlobId headerId = blobIdFactory.parse(headerStr); + referencedBlobIdObserver.accept(headerId); + if (generationCondition.test(headerId)) { + references.computeIfAbsent(headerId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + if (bodyStr != null) { + BlobId bodyId = blobIdFactory.parse(bodyStr); + referencedBlobIdObserver.accept(bodyId); + if (generationCondition.test(bodyId)) { + references.computeIfAbsent(bodyId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + }) + .then(Mono.fromCallable(() -> new CassandraBlobIdUpdater( + cassandraAsyncExecutor, references, + selectMessageV3, updateMessageV3Header, updateMessageV3Body, + selectImapUidByMessageId, updateImapUidHeader, updateMessageIdTableHeader))); + } + } + + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final Map<BlobId, Set<UUID>> references; + private final PreparedStatement selectMessageV3; + private final PreparedStatement updateMessageV3Header; + private final PreparedStatement updateMessageV3Body; + private final PreparedStatement selectImapUidByMessageId; + private final PreparedStatement updateImapUidHeader; + private final PreparedStatement updateMessageIdTableHeader; + + public CassandraBlobIdUpdater(CassandraAsyncExecutor cassandraAsyncExecutor, + Map<BlobId, Set<UUID>> references, + PreparedStatement selectMessageV3, + PreparedStatement updateMessageV3Header, + PreparedStatement updateMessageV3Body, + PreparedStatement selectImapUidByMessageId, + PreparedStatement updateImapUidHeader, + PreparedStatement updateMessageIdTableHeader) { + this.cassandraAsyncExecutor = cassandraAsyncExecutor; + this.references = references; + this.selectMessageV3 = selectMessageV3; + this.updateMessageV3Header = updateMessageV3Header; + this.updateMessageV3Body = updateMessageV3Body; + this.selectImapUidByMessageId = selectImapUidByMessageId; + this.updateImapUidHeader = updateImapUidHeader; + this.updateMessageIdTableHeader = updateMessageIdTableHeader; + } + + @Override + public Mono<Void> replaceReferences(BlobId oldId, BlobId newId) { + Set<UUID> messageUuids = references.remove(oldId); + if (messageUuids == null || messageUuids.isEmpty()) { + return Mono.empty(); + } + String oldIdStr = oldId.asString(); + String newIdStr = newId.asString(); + + return Flux.fromIterable(messageUuids) + .flatMap(messageUuid -> updateMessageReferences(messageUuid, oldIdStr, newIdStr), ReactorUtils.DEFAULT_CONCURRENCY) + .then(); + } + + private Mono<Void> updateMessageReferences(UUID messageUuid, String oldIdStr, String newIdStr) { Review Comment: This method is too complex and deserves to be extracted in small parts ########## mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraBlobIdUpdater.java: ########## @@ -0,0 +1,231 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.mailbox.cassandra.mail; + +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.bindMarker; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.selectFrom; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.update; +import static com.datastax.oss.driver.api.querybuilder.relation.Relation.column; +import static com.datastax.oss.driver.api.querybuilder.update.Assignment.setColumn; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.IMAP_UID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MAILBOX_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MESSAGE_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.BODY_CONTENT; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.HEADER_CONTENT; + +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; +import java.util.function.Predicate; + +import jakarta.inject.Inject; + +import org.apache.james.backends.cassandra.init.configuration.JamesExecutionProfiles; +import org.apache.james.backends.cassandra.utils.CassandraAsyncExecutor; +import org.apache.james.blob.api.BlobId; +import org.apache.james.blob.api.BlobIdUpdater; +import org.apache.james.mailbox.cassandra.table.CassandraMessageIdTable; +import org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table; +import org.apache.james.mailbox.cassandra.table.MessageIdToImapUid; +import org.apache.james.util.ReactorUtils; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.type.codec.TypeCodecs; + +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +public class CassandraBlobIdUpdater implements BlobIdUpdater { + public static class Factory implements BlobIdUpdater.Factory { + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final BlobId.Factory blobIdFactory; + private final DriverExecutionProfile batchProfile; + private final PreparedStatement selectAll; + private final PreparedStatement selectMessageV3; + private final PreparedStatement updateMessageV3Header; + private final PreparedStatement updateMessageV3Body; + private final PreparedStatement selectImapUidByMessageId; + private final PreparedStatement updateImapUidHeader; + private final PreparedStatement updateMessageIdTableHeader; + + @Inject + public Factory(CqlSession session, BlobId.Factory blobIdFactory) { + this.cassandraAsyncExecutor = new CassandraAsyncExecutor(session); + this.blobIdFactory = blobIdFactory; + this.batchProfile = JamesExecutionProfiles.getBatchProfile(session); + + this.selectAll = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(MESSAGE_ID, HEADER_CONTENT, BODY_CONTENT) + .build()); + + this.selectMessageV3 = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(HEADER_CONTENT, BODY_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Header = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Body = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(BODY_CONTENT, bindMarker(BODY_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.selectImapUidByMessageId = session.prepare(selectFrom(MessageIdToImapUid.TABLE_NAME) + .columns(MAILBOX_ID, IMAP_UID, HEADER_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateImapUidHeader = session.prepare(update(MessageIdToImapUid.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID)), + column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + + this.updateMessageIdTableHeader = session.prepare(update(CassandraMessageIdTable.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + } + + @Override + public Mono<BlobIdUpdater> forPredicate(Predicate<BlobId> generationCondition, + Consumer<BlobId> referencedBlobIdObserver) { + Map<BlobId, Set<UUID>> references = new ConcurrentHashMap<>(); + return cassandraAsyncExecutor.executeRows(selectAll.bind().setExecutionProfile(batchProfile)) + .doOnNext(row -> { + UUID messageId = row.get(MESSAGE_ID, TypeCodecs.TIMEUUID); + String headerStr = row.get(HEADER_CONTENT, TypeCodecs.TEXT); + String bodyStr = row.get(BODY_CONTENT, TypeCodecs.TEXT); + if (headerStr != null) { + BlobId headerId = blobIdFactory.parse(headerStr); + referencedBlobIdObserver.accept(headerId); + if (generationCondition.test(headerId)) { + references.computeIfAbsent(headerId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + if (bodyStr != null) { + BlobId bodyId = blobIdFactory.parse(bodyStr); + referencedBlobIdObserver.accept(bodyId); + if (generationCondition.test(bodyId)) { + references.computeIfAbsent(bodyId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + }) + .then(Mono.fromCallable(() -> new CassandraBlobIdUpdater( + cassandraAsyncExecutor, references, + selectMessageV3, updateMessageV3Header, updateMessageV3Body, + selectImapUidByMessageId, updateImapUidHeader, updateMessageIdTableHeader))); + } + } + + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final Map<BlobId, Set<UUID>> references; + private final PreparedStatement selectMessageV3; + private final PreparedStatement updateMessageV3Header; + private final PreparedStatement updateMessageV3Body; + private final PreparedStatement selectImapUidByMessageId; + private final PreparedStatement updateImapUidHeader; + private final PreparedStatement updateMessageIdTableHeader; + + public CassandraBlobIdUpdater(CassandraAsyncExecutor cassandraAsyncExecutor, + Map<BlobId, Set<UUID>> references, + PreparedStatement selectMessageV3, + PreparedStatement updateMessageV3Header, + PreparedStatement updateMessageV3Body, + PreparedStatement selectImapUidByMessageId, + PreparedStatement updateImapUidHeader, + PreparedStatement updateMessageIdTableHeader) { + this.cassandraAsyncExecutor = cassandraAsyncExecutor; + this.references = references; + this.selectMessageV3 = selectMessageV3; + this.updateMessageV3Header = updateMessageV3Header; + this.updateMessageV3Body = updateMessageV3Body; + this.selectImapUidByMessageId = selectImapUidByMessageId; + this.updateImapUidHeader = updateImapUidHeader; + this.updateMessageIdTableHeader = updateMessageIdTableHeader; + } + + @Override + public Mono<Void> replaceReferences(BlobId oldId, BlobId newId) { + Set<UUID> messageUuids = references.remove(oldId); + if (messageUuids == null || messageUuids.isEmpty()) { + return Mono.empty(); + } + String oldIdStr = oldId.asString(); + String newIdStr = newId.asString(); + + return Flux.fromIterable(messageUuids) + .flatMap(messageUuid -> updateMessageReferences(messageUuid, oldIdStr, newIdStr), ReactorUtils.DEFAULT_CONCURRENCY) + .then(); + } + + private Mono<Void> updateMessageReferences(UUID messageUuid, String oldIdStr, String newIdStr) { Review Comment: This method is too complex and deserves to be extracted in small parts ########## server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/InitialBlobCompactionTask.java: ########## @@ -0,0 +1,165 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.blob.compaction; + +import java.time.Clock; +import java.time.Instant; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.james.task.Task; +import org.apache.james.task.TaskExecutionDetails; +import org.apache.james.task.TaskType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.common.base.Preconditions; + +public class InitialBlobCompactionTask implements Task { + public static final TaskType TASK_TYPE = TaskType.of("InitialBlobCompactionTask"); + private static final Logger LOGGER = LoggerFactory.getLogger(InitialBlobCompactionTask.class); + + public record AdditionalInformation( + Instant timestamp, + String bucketName, + long generation, + Optional<Integer> family, + long packedBlobs, + long packedBytes, + long chunksWritten, + long freedBytes + ) implements TaskExecutionDetails.AdditionalInformation { + + public AdditionalInformation { + Preconditions.checkNotNull(timestamp, "'timestamp' must not be null"); + Preconditions.checkNotNull(bucketName, "'bucketName' must not be null"); + Preconditions.checkArgument(generation >= 0, "'generation' must not be negative"); + Preconditions.checkNotNull(family, "'family' must not be null"); + family.ifPresent(f -> Preconditions.checkArgument(f > 0, "'family' must be strictly positive")); + Preconditions.checkArgument(packedBlobs >= 0, "'packedBlobs' must not be negative"); + Preconditions.checkArgument(packedBytes >= 0, "'packedBytes' must not be negative"); + Preconditions.checkArgument(chunksWritten >= 0, "'chunksWritten' must not be negative"); + Preconditions.checkArgument(freedBytes >= 0, "'freedBytes' must not be negative"); + } + + public Instant getTimestamp() { + return timestamp; + } + + public String getBucketName() { + return bucketName; + } + + public long getGeneration() { + return generation; + } + + public Optional<Integer> getFamily() { + return family; + } + + public long getPackedBlobs() { + return packedBlobs; + } + + public long getPackedBytes() { + return packedBytes; + } + + public long getChunksWritten() { + return chunksWritten; + } + + public long getFreedBytes() { + return freedBytes; + } + } + + private final BlobCompactionAlgorithm algorithm; + private final CompactionRequest request; + private final Clock clock; + private final AtomicReference<CompactionResult> currentResult; + + public InitialBlobCompactionTask(BlobCompactionAlgorithm algorithm, CompactionRequest request, Clock clock) { + this.algorithm = Preconditions.checkNotNull(algorithm, "'algorithm' must not be null"); + this.request = Preconditions.checkNotNull(request, "'request' must not be null"); + this.clock = Preconditions.checkNotNull(clock, "'clock' must not be null"); + this.currentResult = new AtomicReference<>(CompactionResult.NONE); + } + + @Override + public Result run() { + try { + CompactionResult result = algorithm.initialCompact(request).block(); + if (result != null) { + currentResult.set(result); + } + return Result.COMPLETED; + } catch (Exception e) { + LOGGER.error("Error while running InitialBlobCompactionTask for generation {}", request.generation(), e); + return Result.PARTIAL; + } + } + + @Override + public TaskType type() { + return TASK_TYPE; + } + + @Override + public Optional<TaskExecutionDetails.AdditionalInformation> details() { + CompactionResult res = currentResult.get(); + return Optional.of(new AdditionalInformation( + clock.instant(), + request.bucketName().asString(), + request.generation(), + request.family(), + res.packedBlobs(), + res.packedBytes(), + res.chunksWritten(), + res.freedBytes() + )); + } + + public CompactionRequest getRequest() { + return request; + } + + public Clock getClock() { + return clock; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o instanceof InitialBlobCompactionTask that) { + return Objects.equals(request, that.request); + } + return false; + } + + @Override + public int hashCode() { + return Objects.hash(request); + } Review Comment: Do we need equals and hashcode for tasks ? ########## server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/ChunkFormat.java: ########## @@ -0,0 +1,315 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.blob.compaction; + +import java.io.ByteArrayOutputStream; +import java.io.DataOutputStream; +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.Map; +import java.util.zip.CRC32C; + +import org.apache.james.blob.api.BlobStoreDAO.BlobMetadata; +import org.apache.james.blob.api.BlobStoreDAO.BlobMetadataName; +import org.apache.james.blob.api.BlobStoreDAO.BlobMetadataValue; +import org.apache.james.blob.api.ObjectStoreIOException; + +import com.google.common.base.Joiner; +import com.google.common.base.Preconditions; +import com.google.common.base.Splitter; +import com.google.common.collect.ImmutableList; +import com.google.common.io.ByteStreams; + +/** + * Binary chunk layout for S3 multi-slot compaction. + * + * A chunk consists of sequential slot payloads followed by a trailing footer: + * [FORMAT_BYTE (1B)] [SLOT_1] [SLOT_2] ... [SLOT_N] [FOOTER_TEXT] [FOOTER_LEN (4B)] [FOOTER_POS (8B)] + */ +public class ChunkFormat { Review Comment: Long complicated methods in here too. Code extraction needed... ########## mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraBlobIdUpdater.java: ########## @@ -0,0 +1,231 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.mailbox.cassandra.mail; + +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.bindMarker; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.selectFrom; +import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.update; +import static com.datastax.oss.driver.api.querybuilder.relation.Relation.column; +import static com.datastax.oss.driver.api.querybuilder.update.Assignment.setColumn; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.IMAP_UID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MAILBOX_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MESSAGE_ID; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.BODY_CONTENT; +import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.HEADER_CONTENT; + +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; +import java.util.function.Predicate; + +import jakarta.inject.Inject; + +import org.apache.james.backends.cassandra.init.configuration.JamesExecutionProfiles; +import org.apache.james.backends.cassandra.utils.CassandraAsyncExecutor; +import org.apache.james.blob.api.BlobId; +import org.apache.james.blob.api.BlobIdUpdater; +import org.apache.james.mailbox.cassandra.table.CassandraMessageIdTable; +import org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table; +import org.apache.james.mailbox.cassandra.table.MessageIdToImapUid; +import org.apache.james.util.ReactorUtils; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.type.codec.TypeCodecs; + +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +public class CassandraBlobIdUpdater implements BlobIdUpdater { + public static class Factory implements BlobIdUpdater.Factory { + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final BlobId.Factory blobIdFactory; + private final DriverExecutionProfile batchProfile; + private final PreparedStatement selectAll; + private final PreparedStatement selectMessageV3; + private final PreparedStatement updateMessageV3Header; + private final PreparedStatement updateMessageV3Body; + private final PreparedStatement selectImapUidByMessageId; + private final PreparedStatement updateImapUidHeader; + private final PreparedStatement updateMessageIdTableHeader; + + @Inject + public Factory(CqlSession session, BlobId.Factory blobIdFactory) { + this.cassandraAsyncExecutor = new CassandraAsyncExecutor(session); + this.blobIdFactory = blobIdFactory; + this.batchProfile = JamesExecutionProfiles.getBatchProfile(session); + + this.selectAll = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(MESSAGE_ID, HEADER_CONTENT, BODY_CONTENT) + .build()); + + this.selectMessageV3 = session.prepare(selectFrom(CassandraMessageV3Table.TABLE_NAME) + .columns(HEADER_CONTENT, BODY_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Header = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateMessageV3Body = session.prepare(update(CassandraMessageV3Table.TABLE_NAME) + .set(setColumn(BODY_CONTENT, bindMarker(BODY_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.selectImapUidByMessageId = session.prepare(selectFrom(MessageIdToImapUid.TABLE_NAME) + .columns(MAILBOX_ID, IMAP_UID, HEADER_CONTENT) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID))) + .build()); + + this.updateImapUidHeader = session.prepare(update(MessageIdToImapUid.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID)), + column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + + this.updateMessageIdTableHeader = session.prepare(update(CassandraMessageIdTable.TABLE_NAME) + .set(setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT))) + .where(column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)), + column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID))) + .build()); + } + + @Override + public Mono<BlobIdUpdater> forPredicate(Predicate<BlobId> generationCondition, + Consumer<BlobId> referencedBlobIdObserver) { + Map<BlobId, Set<UUID>> references = new ConcurrentHashMap<>(); + return cassandraAsyncExecutor.executeRows(selectAll.bind().setExecutionProfile(batchProfile)) + .doOnNext(row -> { + UUID messageId = row.get(MESSAGE_ID, TypeCodecs.TIMEUUID); + String headerStr = row.get(HEADER_CONTENT, TypeCodecs.TEXT); + String bodyStr = row.get(BODY_CONTENT, TypeCodecs.TEXT); + if (headerStr != null) { + BlobId headerId = blobIdFactory.parse(headerStr); + referencedBlobIdObserver.accept(headerId); + if (generationCondition.test(headerId)) { + references.computeIfAbsent(headerId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + if (bodyStr != null) { + BlobId bodyId = blobIdFactory.parse(bodyStr); + referencedBlobIdObserver.accept(bodyId); + if (generationCondition.test(bodyId)) { + references.computeIfAbsent(bodyId, k -> ConcurrentHashMap.newKeySet()).add(messageId); + } + } + }) + .then(Mono.fromCallable(() -> new CassandraBlobIdUpdater( + cassandraAsyncExecutor, references, + selectMessageV3, updateMessageV3Header, updateMessageV3Body, + selectImapUidByMessageId, updateImapUidHeader, updateMessageIdTableHeader))); + } + } + + private final CassandraAsyncExecutor cassandraAsyncExecutor; + private final Map<BlobId, Set<UUID>> references; + private final PreparedStatement selectMessageV3; + private final PreparedStatement updateMessageV3Header; + private final PreparedStatement updateMessageV3Body; + private final PreparedStatement selectImapUidByMessageId; + private final PreparedStatement updateImapUidHeader; + private final PreparedStatement updateMessageIdTableHeader; + + public CassandraBlobIdUpdater(CassandraAsyncExecutor cassandraAsyncExecutor, + Map<BlobId, Set<UUID>> references, + PreparedStatement selectMessageV3, + PreparedStatement updateMessageV3Header, + PreparedStatement updateMessageV3Body, + PreparedStatement selectImapUidByMessageId, + PreparedStatement updateImapUidHeader, + PreparedStatement updateMessageIdTableHeader) { + this.cassandraAsyncExecutor = cassandraAsyncExecutor; + this.references = references; + this.selectMessageV3 = selectMessageV3; + this.updateMessageV3Header = updateMessageV3Header; + this.updateMessageV3Body = updateMessageV3Body; + this.selectImapUidByMessageId = selectImapUidByMessageId; + this.updateImapUidHeader = updateImapUidHeader; + this.updateMessageIdTableHeader = updateMessageIdTableHeader; + } + + @Override + public Mono<Void> replaceReferences(BlobId oldId, BlobId newId) { + Set<UUID> messageUuids = references.remove(oldId); + if (messageUuids == null || messageUuids.isEmpty()) { + return Mono.empty(); + } + String oldIdStr = oldId.asString(); + String newIdStr = newId.asString(); + + return Flux.fromIterable(messageUuids) + .flatMap(messageUuid -> updateMessageReferences(messageUuid, oldIdStr, newIdStr), ReactorUtils.DEFAULT_CONCURRENCY) + .then(); + } + + private Mono<Void> updateMessageReferences(UUID messageUuid, String oldIdStr, String newIdStr) { Review Comment: ```suggestion private Mono<Void> updateMessageReferences(UUID messageUuid, BlobId oldIdStr, BlobId newIdStr) { ``` Strong types ########## server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/ChunkedBlobStoreDAO.java: ########## @@ -0,0 +1,225 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.blob.compaction; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.util.Collection; +import java.util.List; + +import jakarta.inject.Inject; +import jakarta.inject.Named; + +import org.apache.james.blob.api.BlobId; +import org.apache.james.blob.api.BlobStoreDAO; +import org.apache.james.blob.api.BucketName; +import org.apache.james.blob.api.ObjectNotFoundException; +import org.apache.james.blob.api.ObjectStoreIOException; +import org.reactivestreams.Publisher; + +import com.google.common.base.Preconditions; + +import reactor.core.publisher.Mono; + +public class ChunkedBlobStoreDAO implements BlobStoreDAO { + public static final String RAW = "raw"; + + private final BlobStoreDAO rawStore; + + @Inject + public ChunkedBlobStoreDAO(@Named(RAW) BlobStoreDAO rawStore) { + this.rawStore = Preconditions.checkNotNull(rawStore, "'rawStore' must not be null"); + } + + @Override + public InputStreamBlob read(BucketName bucketName, BlobId blobId) throws ObjectStoreIOException, ObjectNotFoundException { + if (!ChunkId.isChunkRef(blobId)) { + return rawStore.read(bucketName, blobId); + } + ChunkId chunkId = ChunkId.parseChunkOrSlotRef(blobId.asString()); + if (!chunkId.isSlotRef()) { + return rawStore.read(bucketName, chunkId.chunkBlobId()); + } + BytesBlob bytes = Mono.from(readChunkSlot(bucketName, chunkId)).block(); + if (bytes == null) { + throw new ObjectNotFoundException("Blob not found: " + blobId.asString()); + } + return InputStreamBlob.of(new ByteArrayInputStream(bytes.payload()), bytes.metadata()); + } + + @Override + public Publisher<InputStreamBlob> readReactive(BucketName bucketName, BlobId blobId) { + if (!ChunkId.isChunkRef(blobId)) { + return rawStore.readReactive(bucketName, blobId); + } + ChunkId chunkId = ChunkId.parseChunkOrSlotRef(blobId.asString()); + if (!chunkId.isSlotRef()) { + return rawStore.readReactive(bucketName, chunkId.chunkBlobId()); + } + return Mono.from(readChunkSlot(bucketName, chunkId)) + .map(bytesBlob -> InputStreamBlob.of(new ByteArrayInputStream(bytesBlob.payload()), bytesBlob.metadata())); + } + + @Override + public Publisher<BytesBlob> readBytes(BucketName bucketName, BlobId blobId) { + if (!ChunkId.isChunkRef(blobId)) { + return rawStore.readBytes(bucketName, blobId); + } + ChunkId chunkId = ChunkId.parseChunkOrSlotRef(blobId.asString()); + if (!chunkId.isSlotRef()) { + return rawStore.readBytes(bucketName, chunkId.chunkBlobId()); + } + return readChunkSlot(bucketName, chunkId); + } + + private Mono<BytesBlob> readChunkSlot(BucketName bucketName, ChunkId slotRef) { Review Comment: This method is way to complex and deserve method extraction refactoring ########## server/container/guice/distributed/src/main/java/org/apache/james/modules/blobstore/BlobCompactionModule.java: ########## @@ -0,0 +1,145 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.modules.blobstore; + +import java.time.Clock; +import java.util.Optional; + +import org.apache.james.blob.api.BlobIdUpdater; +import org.apache.james.blob.api.BlobStoreDAO; +import org.apache.james.blob.compaction.BlobCompactionAlgorithm; +import org.apache.james.blob.compaction.BlobCompactionDTOModules; +import org.apache.james.blob.compaction.CompactionConfiguration; +import org.apache.james.modules.blobstore.BlobStoreConfiguration.BlobStoreImplName; +import org.apache.james.server.task.json.dto.AdditionalInformationDTO; +import org.apache.james.server.task.json.dto.AdditionalInformationDTOModule; +import org.apache.james.server.task.json.dto.TaskDTO; +import org.apache.james.server.task.json.dto.TaskDTOModule; +import org.apache.james.task.Task; +import org.apache.james.task.TaskExecutionDetails; +import org.apache.james.webadmin.dto.DTOModuleInjections; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.inject.AbstractModule; +import com.google.inject.Injector; +import com.google.inject.Key; +import com.google.inject.Provider; +import com.google.inject.Provides; +import com.google.inject.Singleton; +import com.google.inject.multibindings.ProvidesIntoSet; +import com.google.inject.name.Named; +import com.google.inject.name.Names; + +public class BlobCompactionModule extends AbstractModule { + private static final Logger LOGGER = LoggerFactory.getLogger(BlobCompactionModule.class); + + @Provides + @Singleton + public CompactionConfiguration compactionConfiguration() { + return CompactionConfiguration.DEFAULT; + } + + @Provides + @Singleton + public Optional<BlobCompactionAlgorithm> optionalBlobCompactionAlgorithm(Injector injector) { Review Comment: Pass the wanted fields as parameter instead of `injector` ########## server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/GCBlobCompactionTask.java: ########## @@ -0,0 +1,158 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.blob.compaction; + +import java.time.Clock; +import java.time.Instant; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.james.task.Task; +import org.apache.james.task.TaskExecutionDetails; +import org.apache.james.task.TaskType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.common.base.Preconditions; + +public class GCBlobCompactionTask implements Task { + public static final TaskType TASK_TYPE = TaskType.of("GCBlobCompactionTask"); + private static final Logger LOGGER = LoggerFactory.getLogger(GCBlobCompactionTask.class); + + public record AdditionalInformation( + Instant timestamp, + String bucketName, + long generation, + Optional<Integer> family, + long deadPurged, + long mergedChunks, + long freedBytes + ) implements TaskExecutionDetails.AdditionalInformation { + + public AdditionalInformation { + Preconditions.checkNotNull(timestamp, "'timestamp' must not be null"); + Preconditions.checkNotNull(bucketName, "'bucketName' must not be null"); + Preconditions.checkArgument(generation >= 0, "'generation' must not be negative"); + Preconditions.checkNotNull(family, "'family' must not be null"); + family.ifPresent(f -> Preconditions.checkArgument(f > 0, "'family' must be strictly positive")); + Preconditions.checkArgument(deadPurged >= 0, "'deadPurged' must not be negative"); + Preconditions.checkArgument(mergedChunks >= 0, "'mergedChunks' must not be negative"); + Preconditions.checkArgument(freedBytes >= 0, "'freedBytes' must not be negative"); + } + + public Instant getTimestamp() { + return timestamp; + } + + public String getBucketName() { + return bucketName; + } + + public long getGeneration() { + return generation; + } + + public Optional<Integer> getFamily() { + return family; + } + + public long getDeadPurged() { + return deadPurged; + } + + public long getMergedChunks() { + return mergedChunks; + } + + public long getFreedBytes() { + return freedBytes; + } + } + + private final BlobCompactionAlgorithm algorithm; + private final CompactionRequest request; + private final Clock clock; + private final AtomicReference<CompactionResult> currentResult; + + public GCBlobCompactionTask(BlobCompactionAlgorithm algorithm, CompactionRequest request, Clock clock) { + this.algorithm = Preconditions.checkNotNull(algorithm, "'algorithm' must not be null"); + this.request = Preconditions.checkNotNull(request, "'request' must not be null"); + this.clock = Preconditions.checkNotNull(clock, "'clock' must not be null"); + this.currentResult = new AtomicReference<>(CompactionResult.NONE); + } + + @Override + public Result run() { + try { + CompactionResult result = algorithm.gcCompact(request).block(); + if (result != null) { + currentResult.set(result); + } + return Result.COMPLETED; + } catch (Exception e) { + LOGGER.error("Error while running GCBlobCompactionTask for generation {}", request.generation(), e); + return Result.PARTIAL; + } + } + + @Override + public TaskType type() { + return TASK_TYPE; + } + + @Override + public Optional<TaskExecutionDetails.AdditionalInformation> details() { + CompactionResult res = currentResult.get(); + return Optional.of(new AdditionalInformation( + clock.instant(), + request.bucketName().asString(), + request.generation(), + request.family(), + res.deadPurged(), + res.mergedChunks(), + res.freedBytes() + )); + } + + public CompactionRequest getRequest() { + return request; + } + + public Clock getClock() { + return clock; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o instanceof GCBlobCompactionTask that) { + return Objects.equals(request, that.request); + } + return false; + } + + @Override + public int hashCode() { + return Objects.hash(request); + } Review Comment: Not needed for tasks IMO ########## server/protocols/webadmin/webadmin-data/src/main/java/org/apache/james/webadmin/routes/BlobRoutes.java: ########## @@ -85,14 +106,133 @@ public String getBasePath() { @Override public void define(Service service) { - TaskFromRequest gcUnreferencedTaskRequest = this::gcUnreferenced; - service.delete(BASE_PATH, gcUnreferencedTaskRequest.asRoute(taskManager), jsonTransformer); + TaskFromRequest deleteTaskRequest = this::delete; + service.delete(BASE_PATH, deleteTaskRequest.asRoute(taskManager), jsonTransformer); + } + + public Task delete(Request request) { + String action = request.queryParams("action"); + if (action == null) { + action = request.queryParams("scope"); + } + if ("gc".equals(action) || "unreferenced".equals(action)) { + return gcUnreferenced(request); + } + if ("initial-compaction".equals(action)) { + return initialCompact(request); + } + if ("re-compaction".equals(action) || "recompaction".equals(action) || "gc-compaction".equals(action)) { + return gcCompact(request); + } + if ("compaction".equals(action)) { + return compact(request); + } + if (request.queryParams("action") != null) { + throw new IllegalArgumentException("'action' query parameter is invalid. Supported actions: 'unreferenced', 'gc', 'initial-compaction', 're-compaction', 'compaction'"); + } + throw new IllegalArgumentException("'scope' is missing or must be 'unreferenced'"); + } + + public Task initialCompact(Request request) { + BlobCompactionAlgorithm algorithm = blobCompactionAlgorithm + .orElseThrow(() -> new IllegalArgumentException("Blob compaction is not configured or not supported on this server (requires S3 blobstore without client-side encryption or whole-blob compression)")); + return new InitialBlobCompactionTask(algorithm, buildCompactionRequest(request), clock); + } + + public Task gcCompact(Request request) { + BlobCompactionAlgorithm algorithm = blobCompactionAlgorithm + .orElseThrow(() -> new IllegalArgumentException("Blob compaction is not configured or not supported on this server (requires S3 blobstore without client-side encryption or whole-blob compression)")); + return new GCBlobCompactionTask(algorithm, buildCompactionRequest(request), clock); + } + + public Task compact(Request request) { + BlobCompactionAlgorithm algorithm = blobCompactionAlgorithm + .orElseThrow(() -> new IllegalArgumentException("Blob compaction is not configured or not supported on this server (requires S3 blobstore without client-side encryption or whole-blob compression)")); + return new BlobCompactionTask(algorithm, buildCompactionRequest(request), clock); + } + + private CompactionRequest buildCompactionRequest(Request request) { Review Comment: Code extraction... -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
