HesandaLiyanage commented on code in PR #3193:
URL: https://github.com/apache/james-project/pull/3193#discussion_r4070313249


##########
server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/BlobCompactionAlgorithm.java:
##########
@@ -0,0 +1,617 @@
+/****************************************************************
+ * 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.IOException;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.Set;
+
+import org.apache.james.blob.api.BlobId;
+import org.apache.james.blob.api.BlobReferenceSource;
+import org.apache.james.blob.api.BlobStoreDAO;
+import org.apache.james.blob.api.ObjectStoreIOException;
+import 
org.apache.james.blob.compaction.BlobReferenceMappingSource.BlobIdMessageIdMapping;
+import org.apache.james.blob.compaction.ChunkFormat.BlobSlotContent;
+import org.apache.james.blob.compaction.ChunkFormat.ChunkWriteResult;
+import org.apache.james.blob.compaction.ChunkFormat.SlotRange;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.google.common.base.Preconditions;
+
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Schedulers;
+
+/**
+ * Executes object compaction and garbage collection for chunked blob storage 
in Apache James.
+ *
+ * <h3>Crash Safety and Liveness Invariants</h3>
+ * The compaction algorithms follow a strict step ordering to ensure crash 
safety and liveness:
+ * <ul>
+ *   <li><b>Initial Compaction:</b>
+ *     <ol>
+ *       <li>Read candidates and accumulate slots into an immutable chunk.</li>
+ *       <li>Save the new chunk to raw storage (unreferenced yet by 
source-of-truth tables).</li>
+ *       <li>Update source-of-truth table references (Cassandra) to the new 
chunk-slot references.</li>
+ *       <li>Delete original standalone objects from raw storage.</li>
+ *     </ol>
+ *     If interrupted before step 3, the saved chunk is an orphan with 0 
references and will be
+ *     safely reclaimed by the next {@code gc-compact} run. The original blobs 
and references remain untouched.
+ *     If interrupted after step 3 but before step 4, the old standalone blobs 
have 0 references and will be
+ *     cleaned up by standard GC.
+ *   </li>
+ *   <li><b>GC-Compact (Rewrite / Merge / Purge):</b>
+ *     <ol>
+ *       <li>Read existing chunk objects and inspect slot references against 
the BloomFilter/mapping.</li>
+ *       <li>For chunks with dead slots (or pairs of small chunks to merge), 
assemble a new chunk with only live slots.</li>
+ *       <li>Save the new chunk to raw storage.</li>
+ *       <li>Update source-of-truth table references to the new chunk slot 
references.</li>
+ *       <li>Delete the old chunk(s) from raw storage.</li>
+ *     </ol>
+ *     If interrupted before step 4, the newly created chunk is an orphan with 
no references and will be
+ *     purged on the subsequent GC-compact pass. If interrupted after step 4, 
the old chunk has 0 references
+ *     and will be purged as an orphan chunk on the next GC-compact pass.
+ *   </li>
+ *   <li><b>Orphan Chunk Purging:</b>
+ *     Chunks with 0 live references (100% dead slots) are identified as 
orphan chunks and deleted immediately.
+ *     This guarantees self-healing and liveness by construction across 
process crashes.
+ *   </li>
+ * </ul>
+ *
+ * <h3>Memory Bounds and Operational Characteristics</h3>
+ * <ul>
+ *   <li><b>Candidate Payload Streaming:</b> {@link 
#initialCompact(CompactionRequest)} streams candidate blob identifiers
+ *       and partitions them into windows of {@value 
#DEFAULT_CANDIDATE_BATCH_SIZE} blobs. Candidate payloads are fetched
+ *       and packed chunk-by-chunk. Payloads are persisted and freed 
window-by-window, ensuring that candidate byte arrays
+ *       are never held in heap for the entire generation simultaneously. 
Candidate payload heap usage is bounded by
+ *       {@code O(min(candidateBatchSize * avgBlobSize, 
chunkTargetSize))}.</li>
+ *   <li><b>Reference Mapping Memory Ceiling:</b> {@code 
loadReferenceMapping()} materializes all live blob-to-messageId
+ *       mappings for the generation into an in-memory multimap. Memory 
consumption is {@code O(liveGenerationReferences)}
+ *       at approximately ~200 bytes per reference (~200MB heap for 1 million 
live references; ~2GB heap for 10 million).
+ *       Because {@link BlobReferenceMappingSource} currently exposes a full 
stream without partition-paged query capabilities,
+ *       this table is loaded per compaction pass. High-scale deployments 
exceeding tens of millions of live references per
+ *       generation can introduce partition-paged reference lookups in future 
iterations.</li>
+ *   <li><b>GC Compaction Memory Bounds:</b> {@link 
#gcCompact(CompactionRequest)} discovers chunks by reading only trailing
+ *       64KB footers via HTTP ranged reads (metadata-only). Orphan chunks 
(100% dead slots) are deleted with 0 payload bytes read.
+ *       During chunk purge or merge, surviving live slots are streamed 
individually via HTTP ranged reads, strictly bounding
+ *       GC payload heap usage to {@code O(maxSlotSize)} (~1MB).</li>
+ * </ul>
+ */
+public class BlobCompactionAlgorithm {
+    public static final int DEFAULT_CANDIDATE_BATCH_SIZE = 1000;
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(BlobCompactionAlgorithm.class);
+
+    private final BlobStoreDAO blobStoreDAO;
+    private final BlobStoreDAO rawStore;
+    private final BlobReferenceSource referenceSource;
+    private final BlobReferenceMappingSource mappingSource;
+    private final BlobIdUpdater blobIdUpdater;
+
+    public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO,
+                                  BlobStoreDAO rawStore,
+                                  BlobReferenceSource referenceSource,
+                                  BlobReferenceMappingSource mappingSource,
+                                  BlobIdUpdater blobIdUpdater) {
+        this.blobStoreDAO = Preconditions.checkNotNull(blobStoreDAO, 
"'blobStoreDAO' must not be null");
+        this.rawStore = Preconditions.checkNotNull(rawStore, "'rawStore' must 
not be null");
+        this.referenceSource = Preconditions.checkNotNull(referenceSource, 
"'referenceSource' must not be null");
+        this.mappingSource = Preconditions.checkNotNull(mappingSource, 
"'mappingSource' must not be null");
+        this.blobIdUpdater = Preconditions.checkNotNull(blobIdUpdater, 
"'blobIdUpdater' must not be null");
+    }
+
+    public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO,
+                                  BlobStoreDAO rawStore,
+                                  BlobReferenceMappingSource mappingSource,
+                                  BlobIdUpdater blobIdUpdater) {
+        this(blobStoreDAO, rawStore,
+            () -> 
Flux.from(mappingSource.listBlobIdMessageIdMappings()).map(BlobIdMessageIdMapping::blobId),
+            mappingSource, blobIdUpdater);
+    }
+
+    public Mono<CompactionResult> compact(CompactionRequest request) {
+        return initialCompact(request)
+            .flatMap(initialResult -> 
gcCompact(request).map(initialResult::combine));
+    }
+
+    public Mono<CompactionResult> initialCompact(CompactionRequest request) {
+        Preconditions.checkNotNull(request, "'request' must not be null");
+
+        return loadReferenceMapping()
+            .flatMap(mapping -> {
+                if (mapping.isEmpty()) {
+                    LOGGER.info("No blob references found in mapping source; 
skipping initial compaction for generation {}", request.generation());
+                    return Mono.just(CompactionResult.NONE);
+                }
+
+                return Flux.from(rawStore.listBlobs(request.bucketName()))
+                    .filter(blobId -> 
matchesGenerationAndFamily(blobId.asString(), request.generation(), 
request.family()))

Review Comment:
   Done. Added generation prefix pushdown when listing blobs using 
`BlobStoreDAO.listBlobs(bucket, prefix)` with fallback to post-filtering when 
prefix listing is unsupported.



##########
server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/BlobCompactionAlgorithm.java:
##########
@@ -0,0 +1,617 @@
+/****************************************************************
+ * 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.IOException;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.Set;
+
+import org.apache.james.blob.api.BlobId;
+import org.apache.james.blob.api.BlobReferenceSource;
+import org.apache.james.blob.api.BlobStoreDAO;
+import org.apache.james.blob.api.ObjectStoreIOException;
+import 
org.apache.james.blob.compaction.BlobReferenceMappingSource.BlobIdMessageIdMapping;
+import org.apache.james.blob.compaction.ChunkFormat.BlobSlotContent;
+import org.apache.james.blob.compaction.ChunkFormat.ChunkWriteResult;
+import org.apache.james.blob.compaction.ChunkFormat.SlotRange;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.google.common.base.Preconditions;
+
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Schedulers;
+
+/**
+ * Executes object compaction and garbage collection for chunked blob storage 
in Apache James.
+ *
+ * <h3>Crash Safety and Liveness Invariants</h3>
+ * The compaction algorithms follow a strict step ordering to ensure crash 
safety and liveness:
+ * <ul>
+ *   <li><b>Initial Compaction:</b>
+ *     <ol>
+ *       <li>Read candidates and accumulate slots into an immutable chunk.</li>
+ *       <li>Save the new chunk to raw storage (unreferenced yet by 
source-of-truth tables).</li>
+ *       <li>Update source-of-truth table references (Cassandra) to the new 
chunk-slot references.</li>
+ *       <li>Delete original standalone objects from raw storage.</li>
+ *     </ol>
+ *     If interrupted before step 3, the saved chunk is an orphan with 0 
references and will be
+ *     safely reclaimed by the next {@code gc-compact} run. The original blobs 
and references remain untouched.
+ *     If interrupted after step 3 but before step 4, the old standalone blobs 
have 0 references and will be
+ *     cleaned up by standard GC.
+ *   </li>
+ *   <li><b>GC-Compact (Rewrite / Merge / Purge):</b>
+ *     <ol>
+ *       <li>Read existing chunk objects and inspect slot references against 
the BloomFilter/mapping.</li>
+ *       <li>For chunks with dead slots (or pairs of small chunks to merge), 
assemble a new chunk with only live slots.</li>
+ *       <li>Save the new chunk to raw storage.</li>
+ *       <li>Update source-of-truth table references to the new chunk slot 
references.</li>
+ *       <li>Delete the old chunk(s) from raw storage.</li>
+ *     </ol>
+ *     If interrupted before step 4, the newly created chunk is an orphan with 
no references and will be
+ *     purged on the subsequent GC-compact pass. If interrupted after step 4, 
the old chunk has 0 references
+ *     and will be purged as an orphan chunk on the next GC-compact pass.
+ *   </li>
+ *   <li><b>Orphan Chunk Purging:</b>
+ *     Chunks with 0 live references (100% dead slots) are identified as 
orphan chunks and deleted immediately.
+ *     This guarantees self-healing and liveness by construction across 
process crashes.
+ *   </li>
+ * </ul>
+ *
+ * <h3>Memory Bounds and Operational Characteristics</h3>
+ * <ul>
+ *   <li><b>Candidate Payload Streaming:</b> {@link 
#initialCompact(CompactionRequest)} streams candidate blob identifiers
+ *       and partitions them into windows of {@value 
#DEFAULT_CANDIDATE_BATCH_SIZE} blobs. Candidate payloads are fetched
+ *       and packed chunk-by-chunk. Payloads are persisted and freed 
window-by-window, ensuring that candidate byte arrays
+ *       are never held in heap for the entire generation simultaneously. 
Candidate payload heap usage is bounded by
+ *       {@code O(min(candidateBatchSize * avgBlobSize, 
chunkTargetSize))}.</li>
+ *   <li><b>Reference Mapping Memory Ceiling:</b> {@code 
loadReferenceMapping()} materializes all live blob-to-messageId
+ *       mappings for the generation into an in-memory multimap. Memory 
consumption is {@code O(liveGenerationReferences)}
+ *       at approximately ~200 bytes per reference (~200MB heap for 1 million 
live references; ~2GB heap for 10 million).
+ *       Because {@link BlobReferenceMappingSource} currently exposes a full 
stream without partition-paged query capabilities,
+ *       this table is loaded per compaction pass. High-scale deployments 
exceeding tens of millions of live references per
+ *       generation can introduce partition-paged reference lookups in future 
iterations.</li>
+ *   <li><b>GC Compaction Memory Bounds:</b> {@link 
#gcCompact(CompactionRequest)} discovers chunks by reading only trailing
+ *       64KB footers via HTTP ranged reads (metadata-only). Orphan chunks 
(100% dead slots) are deleted with 0 payload bytes read.
+ *       During chunk purge or merge, surviving live slots are streamed 
individually via HTTP ranged reads, strictly bounding
+ *       GC payload heap usage to {@code O(maxSlotSize)} (~1MB).</li>
+ * </ul>
+ */
+public class BlobCompactionAlgorithm {
+    public static final int DEFAULT_CANDIDATE_BATCH_SIZE = 1000;
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(BlobCompactionAlgorithm.class);
+
+    private final BlobStoreDAO blobStoreDAO;
+    private final BlobStoreDAO rawStore;
+    private final BlobReferenceSource referenceSource;
+    private final BlobReferenceMappingSource mappingSource;
+    private final BlobIdUpdater blobIdUpdater;
+
+    public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO,
+                                  BlobStoreDAO rawStore,
+                                  BlobReferenceSource referenceSource,
+                                  BlobReferenceMappingSource mappingSource,
+                                  BlobIdUpdater blobIdUpdater) {
+        this.blobStoreDAO = Preconditions.checkNotNull(blobStoreDAO, 
"'blobStoreDAO' must not be null");
+        this.rawStore = Preconditions.checkNotNull(rawStore, "'rawStore' must 
not be null");
+        this.referenceSource = Preconditions.checkNotNull(referenceSource, 
"'referenceSource' must not be null");
+        this.mappingSource = Preconditions.checkNotNull(mappingSource, 
"'mappingSource' must not be null");
+        this.blobIdUpdater = Preconditions.checkNotNull(blobIdUpdater, 
"'blobIdUpdater' must not be null");
+    }
+
+    public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO,
+                                  BlobStoreDAO rawStore,
+                                  BlobReferenceMappingSource mappingSource,
+                                  BlobIdUpdater blobIdUpdater) {
+        this(blobStoreDAO, rawStore,
+            () -> 
Flux.from(mappingSource.listBlobIdMessageIdMappings()).map(BlobIdMessageIdMapping::blobId),
+            mappingSource, blobIdUpdater);
+    }
+
+    public Mono<CompactionResult> compact(CompactionRequest request) {
+        return initialCompact(request)
+            .flatMap(initialResult -> 
gcCompact(request).map(initialResult::combine));
+    }
+
+    public Mono<CompactionResult> initialCompact(CompactionRequest request) {
+        Preconditions.checkNotNull(request, "'request' must not be null");
+
+        return loadReferenceMapping()
+            .flatMap(mapping -> {
+                if (mapping.isEmpty()) {
+                    LOGGER.info("No blob references found in mapping source; 
skipping initial compaction for generation {}", request.generation());
+                    return Mono.just(CompactionResult.NONE);
+                }
+
+                return Flux.from(rawStore.listBlobs(request.bucketName()))
+                    .filter(blobId -> 
matchesGenerationAndFamily(blobId.asString(), request.generation(), 
request.family()))
+                    .filter(blobId -> !ChunkId.isChunkRef(blobId))
+                    .filter(mapping::containsKey)
+                    .window(DEFAULT_CANDIDATE_BATCH_SIZE)

Review Comment:
   Done. Windowing is now based on cumulative byte size up to `chunkTargetSize` 
rather than pure item count.



-- 
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]

Reply via email to