Skip to content
Open
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
110 changes: 106 additions & 4 deletions docs/testing/BINARY_S3_STORAGE.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,8 @@ table may not be seen at first read, and a runtime change is ignored until resta
overrides through `Config.setProperty` do re-read it, which is how tests switch modes. Tests
that mock `Config` statically must call `AssetStorageFeature.reset()` themselves.

With the flag off, the S3 cleanup job queues (`binaryAssetCleanup`, `binaryFieldCleanup`) are not
registered.
With the flag off, the S3 job queues (`binaryAssetCleanup`, `binaryFieldCleanup`,
`binaryAssetBackfill`) are not registered.

### Enabling the flag is not rollback-safe

Expand Down Expand Up @@ -202,6 +202,93 @@ database restore, and archives are kept until removed or expired by the bucket's
With the flag on, the `deleteAllVersionsandBackup` interceptor, previously a no-op, calls its
implementation and the all-version deletion hooks. With the flag off it stays a no-op.

## Migrating existing binaries (backfill)

With the flag on, an active administrator can copy existing binaries to S3 through the job API.
The processor checks administrator status when the job is queued and again when it runs.

```http
POST /api/v1/jobs/binaryAssetBackfill
Content-Type: application/json

{"batchSize":250}
```

Follow the returned `statusUrl` (`GET /api/v1/jobs/{jobId}/status`). `parameters.afterInode` is
the last committed batch cursor and `parameters.verifiedBinaries` counts the binaries verified in
completed batches. The batch size is a number of content inodes, from 1 to 1000. A successful
result includes `complete: true`; progress reaches 100% only when the scan finishes.

Each batch uses conditional, verified backfill of originals and metadata, and converts raw S3
objects to SHA-256 references without changing their database paths. It reads persisted binary
fields, including retired field definitions, and migrates recognized completed renditions for each
inode. Local sources are left in place. Only a fully verified batch advances the cursor, retries
reload the cursor from the job database, and a stale worker cannot overwrite newer progress. The
page size is applied in the SQL query, so each batch reads only its own rows. The job sends a
heartbeat after every inode, so a slow batch is not mistaken for an abandoned job.

Problems in the data itself do not stop the scan, because a retry cannot fix them: a referenced
binary or metadata record that exists neither locally nor in S3, a row whose JSON cannot be parsed,
and a legacy row whose content cannot be found. Each one is logged with its inode and field, the
rest of the row is still copied, and the job reports `skippedCount` and `skippedInodes` (the first
1000 inodes) in its parameters and result. Storage and database errors, including an S3 object
that conflicts with the local bytes, still fail the batch so the retry policy applies; the error
message names the inode.

`POST /api/v1/jobs/{jobId}/cancel` stops between batches. After a cancellation or exhausted
retries, submit a new job with the last persisted parameters, for example
`{"batchSize":250,"afterInode":"<parameters.afterInode>","verifiedBinaries":42}`, or omit
`afterInode` to verify everything again. With the flag off, enqueueing and running this job are
rejected. The scan covers content-referenced originals and metadata, not operational server files.

### Legacy Image and File values

Backfill, starter export and recovery archives use the same inventory of persisted references.
Typed `Binary` entries are included even when their field definitions were retired. An `Image` or
`File` entry holding a plain filename is included only when the exact
`<inode-shards>/<inode>/<field>/<filename>` object exists in the combined filesystem and S3
inventory, so linked asset identifiers, external URLs and unrelated text are never treated as
binaries. Listing failures propagate instead of being read as absence.

## Starter export and import

With the flag on, starter asset export lists persisted binary references, restores originals from
S3 and includes raw metadata alongside legacy local files. It holds a cache lease for one binary at
a time, so eviction can keep trimming the files it restores, and it skips a binary whose stored
metadata says it exceeds `maxSize` before restoring it. Restored originals still pass through the
local cache; streaming them from S3 straight into the archive is not implemented. A missing
referenced original or any other export error fails the export without finishing the archive: the
client receives a truncated ZIP with no central directory, which standard ZIP readers reject,
instead of a valid ZIP that silently lacks the remaining entries. Whether the HTTP transfer also
ends with an error depends on the servlet container. With the flag off, export behaves as before,
including logging and skipping a file that cannot be read.

Importing a starter publishes its binaries to S3 before the database commit, and cleanup of the
imported files runs only after both succeed. A referenced binary or metadata record that is in
neither the starter nor S3 is logged and skipped, as before, so starters exported without assets,
with `maxSize` or with `oldAssets=false` still import; a summary count is logged at the end. S3 and
database errors still fail the import, and so does an error importing rules, so a starter is
never left partly imported.

Importing over a populated database with the flag on performs a full replacement: existing
variants, workflows, templates, categories, rules and experiments are cleared before import,
foreign keys stay enforced, and caches are flushed afterwards. With the flag off, import behaves as
before.

A starter exported with the flag on cannot be fully restored into an instance with the flag off.
Content checked in with the flag on has its binary only at its `.revisions/<uuid>/` key, and the
archive stores it there; a flag-off instance resolves binaries by the legacy field path and ignores
`storageKey`, so those binaries are missing after import. Import such a starter only into an
instance with the flag on.

## Integrity checks

With the flag on, the file-asset integrity checker's repair copies stored binaries to the repaired
content before its corrected JSON is published, and removes the old sources in the same
repair transaction. The copies are uploaded before the commit, so a rollback listener deletes
them, and their metadata, if the repair rolls back; a failed deletion is logged and leaves an
unreferenced object behind.

## Local cache eviction

With the flag on, the local asset directory is a cache that `BinaryCacheEvictionJob` can trim.
Expand Down Expand Up @@ -300,7 +387,7 @@ docker run -d --rm --name binary-cleanup-postgres-test \
-e POSTGRES_DB=binary_storage_test postgres:16-alpine

./mvnw test -pl :dotcms-core -Dmaven.build.cache.enabled=false \
-Dtest=AssetStorageFeatureTest,AssetStorageFeatureLatchTest,S3StorageConfigurationTest,NoWebIdentityCredentialsProviderChainTest,BinaryS3StorageTest,BinaryAssetReferenceTest,BinaryCacheEvictionJobTest,BinaryFileSystemStorageTest,BinaryAssetStorageAPIImplTest,MetadataLocalCacheTest,BinaryAssetCleanupTransactionTest,BinaryAssetCleanupProcessorTest,ContentletBackupStorageGateTest,BinaryFieldCleanupProcessorTest,AssetJobEventSerializationTest \
-Dtest=AssetStorageFeatureTest,AssetStorageFeatureLatchTest,S3StorageConfigurationTest,NoWebIdentityCredentialsProviderChainTest,BinaryS3StorageTest,BinaryAssetReferenceTest,BinaryCacheEvictionJobTest,BinaryFileSystemStorageTest,BinaryAssetStorageAPIImplTest,MetadataLocalCacheTest,BinaryAssetCleanupTransactionTest,BinaryAssetCleanupProcessorTest,ContentletBackupStorageGateTest,BinaryFieldCleanupProcessorTest,AssetJobEventSerializationTest,BinaryAssetBackfillCheckpointTest,ImportStarterWorkflowCleanupTest,BinaryAssetBackfillProcessorTest,BinaryAssetBackfillTest,ExportStarterFailureTest \
-Ds3.test.endpoint=http://127.0.0.1:19002 \
-Ds3.test.jdbc=jdbc:postgresql://127.0.0.1:19003/binary_storage_test

Expand All @@ -313,7 +400,22 @@ stack:
```sh
./mvnw install -pl :dotcms-core --am -DskipTests -Ddocker.skip
./mvnw verify -pl :dotcms-integration -Dmaven.build.cache.enabled=false -Dcoreit.test.skip=false \
-Dit.test=BinaryAssetStorageIntegrationTest,ContentletBackupStorageTest,SharedAssetStorageIntegrationTest
-Dit.test=BinaryAssetStorageIntegrationTest,ContentletBackupStorageTest,SharedAssetStorageIntegrationTest,BinaryAssetStarterRestoreTest
```

These default to flag-off mode, where the S3 cases are skipped. To run the S3 cases, create a
bucket in the disposable MinIO (for example `mc mb local/s3-cms-it` inside the container) and add
`-Dit.test.forkcount=1 -Ds3.cms.enabled=true -DDOT_FEATURE_FLAG_S3_ASSET_STORAGE=true
-DDOT_BINARY_ASSET_STORAGE_TYPE=BINARY_CHAIN -DDOT_STORAGE_FILE_METADATA_DEFAULT_CHAIN=FILE_SYSTEM,S3`
plus the `DOT_STORAGE_FILE_METADATA_S3_BUCKET_NAME`, `_BUCKET_REGION`, `_ACCESS_KEY`,
`_SECRET_ACCESS_KEY` and `_ENDPOINT` properties.

`BinaryAssetStarterRestoreTest` is an opt-in starter acceptance test that runs in phases against
the same bucket. Add `-Ds3.starter.restore.enabled=true` and a persistent
`-Ds3.starter.archive=<absolute path under dotCMS/target>/starter.zip`, then run
`-Ds3.starter.phase=export`, followed by `-Ds3.starter.phase=restore` with
`-Dstarter.run.path` set to the same archive. `-Ds3.starter.phase=populated` (after a fresh
export, without `-Dstarter.run.path`) tests replacing a populated database; run it only in the
disposable harness, because it replaces that database's content.

CI does not yet provide the MinIO service, so `BinaryS3StorageTest` does not run there.
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,19 @@

import com.dotcms.content.business.json.ContentletJsonAPI;
import com.dotcms.content.elasticsearch.business.ContentletIndexAPI;
import com.dotcms.storage.AssetStorageFeature;
import com.dotcms.storage.FetchMetadataParams;
import com.dotcms.storage.FileMetadataAPI;
import com.dotcms.storage.StorageKey;
import com.dotcms.storage.StoragePersistenceProvider;
import com.dotcms.storage.binary.BinaryAssetCleanupProcessor;
import com.dotcms.storage.binary.BinaryAssetReference;
import com.dotmarketing.beans.Identifier;
import com.dotmarketing.business.APILocator;
import com.dotmarketing.business.DotStateException;
import com.dotmarketing.common.db.DotConnect;
import com.dotmarketing.db.DbConnectionFactory;
import com.dotmarketing.db.HibernateUtil;
import com.dotmarketing.exception.DotDataException;
import com.dotmarketing.exception.DotHibernateException;
import com.dotmarketing.exception.DotRuntimeException;
Expand All @@ -17,6 +25,7 @@
import com.dotmarketing.portlets.fileassets.business.FileAssetAPI;
import com.dotmarketing.portlets.structure.model.Structure;
import com.dotmarketing.util.Constants;
import com.dotmarketing.util.Config;
import com.dotmarketing.util.Logger;
import com.dotmarketing.util.UtilMethods;
import com.liferay.portal.model.User;
Expand Down Expand Up @@ -78,6 +87,25 @@ public boolean generateIntegrityResults(final String endpointId) throws Exceptio

@Override
public void executeFix(final String key) throws DotDataException, DotSecurityException {
if (!AssetStorageFeature.isEnabled()) {
executeFixInternal(key);
return;
}
final boolean localTransaction = HibernateUtil.startLocalTransactionIfNeeded();
try {
executeFixInternal(key);
if (localTransaction) {
HibernateUtil.commitTransaction();
}
} catch (DotDataException | DotSecurityException | RuntimeException failure) {
if (localTransaction) {
HibernateUtil.rollbackTransaction();
}
throw failure;
}
}

private void executeFixInternal(final String key) throws DotDataException, DotSecurityException {
DotConnect dc = new DotConnect();
// Get information from IR.
final String getResultsQuery = new StringBuilder("SELECT ")
Expand Down Expand Up @@ -200,6 +228,11 @@ private void fixContentletConflicts(final Map<String, Object> contentletData,
final ContentletAPI contentletAPI = APILocator.getContentletAPI();
final ContentletJsonAPI contentletJsonAPI = APILocator.getContentletJsonAPI();

if (AssetStorageFeature.isEnabled()) {
new DotConnect().setSQL("select inode from contentlet where inode in (?, ?) order by inode for update")
.addParam(localWorkingInode).addParam(localLiveInode).loadObjectResults();
}

User systemUser = APILocator.getUserAPI().getSystemUser();
Contentlet existingWorkingContentlet = contentletAPI.find(localWorkingInode, systemUser,
false);
Expand Down Expand Up @@ -249,6 +282,10 @@ private void fixContentletConflicts(final Map<String, Object> contentletData,
workingCopy.setLanguageId(languageId);
workingCopy.setModDate(new Date());

if (AssetStorageFeature.isEnabled() && structureTypeId == Structure.STRUCTURE_TYPE_FILEASSET) {
copyStoredBinaries(existingWorkingContentlet, workingCopy);
}

final String workingCopyJson = Try.of(() -> contentletJsonAPI.toJson(workingCopy))
.getOrElseThrow(() -> new DotRuntimeException(
String.format("Error converting contentlet to json from local working copy with inode [%s].", localWorkingInode)));
Expand Down Expand Up @@ -286,6 +323,10 @@ private void fixContentletConflicts(final Map<String, Object> contentletData,
liveCopy.setLanguageId(languageId);
liveCopy.setModDate(new Date());

if (AssetStorageFeature.isEnabled() && structureTypeId == Structure.STRUCTURE_TYPE_FILEASSET) {
copyStoredBinaries(existingLiveContentlet, liveCopy);
}

final String liveCopyJson = Try.of(() -> contentletJsonAPI.toJson(liveCopy))
.getOrElseThrow(() -> new DotDataException(
String.format("Error converting contentlet to json from local live copy with inode [%s].", localLiveInode)));
Expand Down Expand Up @@ -467,6 +508,79 @@ private void fixContentletConflicts(final Map<String, Object> contentletData,

// Remove the Lucene index for the old page
cleanIndex(existingWorkingContentlet, existingLiveContentlet);
if (AssetStorageFeature.isEnabled() && structureTypeId == Structure.STRUCTURE_TYPE_FILEASSET) {
BinaryAssetCleanupProcessor.enqueue(localWorkingInode);
if (UtilMethods.isSet(localLiveInode) && UtilMethods.isSet(remoteLiveInode)
&& !localLiveInode.equals(localWorkingInode)) {
BinaryAssetCleanupProcessor.enqueue(localLiveInode);
}
}
}

/**
* Copies each stored binary and its metadata to new revision keys owned by the repaired content,
* before its corrected JSON is published; source cleanup is committed with the repair. The copies
* are uploaded immediately, so each one registers a rollback listener that removes it again if
* the repair transaction rolls back.
*
* @param source the local content whose binaries are copied
* @param destination the repaired content that receives the copies
* @throws DotDataException if a source, its referenced metadata or a copy cannot be stored
*/
private void copyStoredBinaries(final Contentlet source, final Contentlet destination)
throws DotDataException {
final var metadata = APILocator.getFileMetadataAPI();
final var files = APILocator.getFileStorageAPI();
final String group = Config.getStringProperty(StoragePersistenceProvider.METADATA_GROUP_NAME,
FileMetadataAPI.DOT_METADATA);
for (final var field : source.getContentType().fields(com.dotcms.contenttype.model.field.BinaryField.class)) {
if (source.get(field.variable()) == null) {
continue;
}
try {
final File original = source.getBinary(field.variable());
if (original == null) {
throw new DotDataException("Missing repair source " + source.getInode() + "/" + field.variable());
}
File copied = APILocator.getBinaryAssetStorageAPI().storeRevision(destination.getInode(),
field.variable(), original.getName(), original);
final String revision = BinaryAssetReference.keyOf(copied, destination.getInode(), field.variable());
removeOnRollback(revision, () -> APILocator.getBinaryAssetStorageAPI()
.deleteBinaryPaths(destination.getInode(), field.variable(), List.of(revision)));
final var attributes = files.retrieveRawMetaData(new StorageKey.Builder().group(group)
.path(metadata.getFileName(source, field.variable()))
.storage(StoragePersistenceProvider.getStorageType()).build());
if (attributes != null) {
final String key = BinaryAssetReference.newMetadataKey(copied, destination.getInode(), field.variable());
final var metadataKey = new FetchMetadataParams.Builder().cache(false)
.storageKey(new StorageKey.Builder().group(group).path(key)
.storage(StoragePersistenceProvider.getStorageType()).build()).build();
removeOnRollback(key, () -> files.removeMetaData(metadataKey));
if (!files.setMetadata(metadataKey, attributes)) {
throw new DotDataException("Unable to copy repair metadata");
}
copied = BinaryAssetReference.withMetadata(copied, destination.getInode(), field.variable(), key);
} else if (BinaryAssetReference.metadataKeyOf(original) != null) {
throw new DotDataException("Missing referenced repair metadata");
}
destination.setBinary(field.variable(), copied);
} catch (IOException e) {
throw new DotDataException("Unable to copy repair source " + source.getInode(), e);
}
}
}

/**
* Registers a rollback listener that deletes an object a repair uploaded before its transaction
* rolled back. A failed deletion is logged rather than thrown, so it cannot mask the rollback.
*
* @param key the uploaded object's key, for the log message
* @param delete deletes the object
*/
private static void removeOnRollback(final String key, final io.vavr.CheckedRunnable delete) {
HibernateUtil.addRollbackListener(() -> Try.run(delete).onFailure(failure -> Logger.warn(
ContentFileAssetIntegrityChecker.class, "Unable to remove " + key
+ " after the integrity repair rolled back: " + failure.getMessage())));
}

private void fixFileAssetContainer(final String oldContentletIdentifier, final String newContentletIdentifier,
Expand Down Expand Up @@ -522,6 +636,13 @@ private Contentlet generateNewContentlet(Contentlet existingContentlet,
final String remoteInode) throws DotContentletStateException, DotRuntimeException,
DotSecurityException, DotDataException {

if (AssetStorageFeature.isEnabled() && structureTypeId == Structure.STRUCTURE_TYPE_FILEASSET) {
com.dotmarketing.business.CacheLocator.getContentletCache().remove(remoteInode);
final Contentlet repaired = APILocator.getContentletAPI().find(remoteInode, APILocator.systemUser(), false);
APILocator.getContentletIndexAPI().addContentToIndex(repaired);
return repaired;
}

// If its an asset file, move the asset to a new location
if (structureTypeId == Structure.STRUCTURE_TYPE_FILEASSET) {
moveInodeFolder(existingContentlet, remoteInode);
Expand Down
Loading
Loading