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
48 changes: 43 additions & 5 deletions docs/testing/BINARY_S3_STORAGE.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ overrides through `Config.setProperty` do re-read it, which is how tests switch
that mock `Config` statically must call `AssetStorageFeature.reset()` themselves.

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

### Enabling the flag is not rollback-safe

Expand Down Expand Up @@ -79,8 +79,8 @@ The storage chain (`ChainableStoragePersistenceAPI`) and its providers change as
`BinaryAssetStorageAPI` (`APILocator.getBinaryAssetStorageAPI()`) manages binary assets and
completed renditions through the storage chain. Binary assets use the `binary-assets` group under
the asset root, laid out as `{inode[0]}/{inode[1]}/{inode}/{field}/{fileName}`. Renditions use the
`generated-assets` group under the `dotGenerated` root. With the flag on, keys in both groups keep
their mixed-case names; other groups are still lowercased.
`generated-assets` group under the `dotGenerated` root. With the flag on, keys in these groups keep
their mixed-case names, as do publishing bundles (below); other groups are still lowercased.

With the flag on, S3 keeps one immutable copy of each set of bytes for these two groups, at
`asset-blobs/sha256/<chars-1-2>/<chars-3-4>/<chars-5-6>/<chars-7-8>/<sha256>`
Expand Down Expand Up @@ -289,6 +289,44 @@ repair transaction. The copies are uploaded before the commit, so a rollback lis
them, and their metadata, if the repair rolls back; a failed deletion is logged and leaves an
unreferenced object behind.

## Publishing bundles

With the flag on, completed push-publishing archives (`<bundle-id>.tar.gz`, including their
manifest) are stored durably in the `publishing-bundles` group, with the local bundle directory as
a cache (`BundleArchiveStorage`). Bundle ids keep their case and are checked so an archive cannot
resolve outside the bundle directory. Generated and received bundles are written to a private
staging file and published to S3 before they replace the last complete local archive, so an
incomplete upload never becomes visible. A received file's bundle id is the part of its name
before the first `.tar.gz`, which is the rule the receiving endpoints and `BundlePublisher` use to
find the archive again. Reading a bundle's manifest or payload restores the archive from S3 only
when its bytes are needed. Storage failures are reported as failures, not as a missing bundle, for
every publishing and retry decision. The bundle pages are the exception: whether to show a
download link or a retry button is checked once per bundle, and a storage failure there is logged
as a warning and shown as "not generated", so an S3 outage does not stop the page from rendering.

With the flag on, static publishing fails the bundle when a File Asset's binary is missing, so it
can be retried. This matches the flag-off behavior, where the copy of the missing file fails.

Deleting a bundle records a `bundleArchiveCleanup` job after the deletion commits. The job deletes
the S3 object and keeps the local bytes if remote deletion fails, so the existing queue can retry.
If a bundle row with the same id exists again when the job runs, for example because a receiver got
the same bundle again, the job keeps the archive and finishes successfully. On one node, storing an
archive and the job's row check and delete hold the same lock, so the job cannot delete an archive
stored between its check and its delete. That lock does not cover other cluster nodes, or a
receiver whose new bundle row is not yet committed when its archive is stored. In those cases the
new archive can still be deleted, and the receive fails with an error and can be retried.

Durable archives are kept for the same time as local ones. When `BinaryCleanupJob` deletes files
older than `CLEANUP_BUNDLES_OLDER_THAN_DAYS` (default 4) from the bundle directory, with the flag on
it also deletes S3 archives whose S3 last-modified time is older than that, together with their
local copy and extraction directory. A bundle older than that can no longer be retried or
downloaded, as with the flag off. A value below 1 turns off both cleanups. An archive that cannot be
deleted is logged and tried again on the next run. If several nodes run the cleanup job, each one
lists the group and deletes the same expired archives, which is harmless because deleting an archive
that is already gone succeeds.

With the flag off, bundles are written, read and deleted in the bundle directory as before.

## Local cache eviction

With the flag on, the local asset directory is a cache that `BinaryCacheEvictionJob` can trim.
Expand Down Expand Up @@ -387,7 +425,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,BinaryAssetBackfillCheckpointTest,ImportStarterWorkflowCleanupTest,BinaryAssetBackfillProcessorTest,BinaryAssetBackfillTest,ExportStarterFailureTest \
-Dtest=AssetStorageFeatureTest,AssetStorageFeatureLatchTest,S3StorageConfigurationTest,NoWebIdentityCredentialsProviderChainTest,BinaryS3StorageTest,BinaryAssetReferenceTest,BinaryCacheEvictionJobTest,BinaryFileSystemStorageTest,BinaryAssetStorageAPIImplTest,MetadataLocalCacheTest,BinaryAssetCleanupTransactionTest,BinaryAssetCleanupProcessorTest,ContentletBackupStorageGateTest,BinaryFieldCleanupProcessorTest,AssetJobEventSerializationTest,BinaryAssetBackfillCheckpointTest,ImportStarterWorkflowCleanupTest,BinaryAssetBackfillProcessorTest,BinaryAssetBackfillTest,ExportStarterFailureTest,BundleArchiveStorageTest,FileAssetBundlerTest \
-Ds3.test.endpoint=http://127.0.0.1:19002 \
-Ds3.test.jdbc=jdbc:postgresql://127.0.0.1:19003/binary_storage_test

Expand All @@ -400,7 +438,7 @@ 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,BinaryAssetStarterRestoreTest
-Dit.test=BinaryAssetStorageIntegrationTest,ContentletBackupStorageTest,SharedAssetStorageIntegrationTest,BinaryAssetStarterRestoreTest,PublishingArchiveStorageTest
```

These default to flag-off mode, where the S3 cases are skipped. To run the S3 cases, create a
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import com.dotcms.publishing.PublisherConfig;
import com.dotcms.publishing.PublisherConfig.Operation;
import com.dotcms.publishing.output.BundleOutput;
import com.dotcms.storage.binary.BinaryAssetStorageAPI;
import com.dotmarketing.beans.Host;
import com.dotmarketing.business.APILocator;
import com.dotmarketing.business.UserAPI;
Expand Down Expand Up @@ -323,6 +324,21 @@ private void writeFileAsset(final BundleOutput output, final FileAsset fileAsset
}
}

/**
* Writes one File Asset into the bundle for one language: its XML descriptor (except for static
* publishing) and a copy of its binary file, or removes them when an incremental live-only bundle
* no longer includes the asset. Files already in the bundle with the same modification date are
* left alone.
*
* @param host the {@link Host} the File Asset lives in
* @param languageId the language folder to write the asset under
* @param output the bundle being written
* @param fileAssetWrapper the File Asset with its version info and identifier
* @throws IOException if the asset's binary file is missing or cannot be copied, which
* fails the bundle so it can be retried
* @throws DotDataException if the binary file cannot be looked up
* @throws DotSecurityException if the asset cannot be read
*/
private void writeFileToDisk(Host host, String languageId, BundleOutput output, FileAssetWrapper fileAssetWrapper)
throws IOException, DotDataException, DotSecurityException {

Expand Down Expand Up @@ -371,13 +387,27 @@ private void writeFileToDisk(Host host, String languageId, BundleOutput output,
else {
//only write if changed
if(!output.exists(filePath) || output.lastModified(filePath) != cal.getTimeInMillis()){
File oldAsset = new File(APILocator.getFileAssetAPI().getRealAssetPathIgnoreExtensionCase(fileAssetWrapper.getAsset().getInode(), fileAssetWrapper.getAsset().getUnderlyingFileName()));
if(output.exists(filePath)) {
output.delete(filePath);
}
// With S3 asset storage on, the cache lease keeps the binary from being evicted
// between finding it and copying it. A null resource is allowed and never closed.
try (BinaryAssetStorageAPI.CacheLease lease = com.dotcms.storage.AssetStorageFeature.isEnabled()
? APILocator.getBinaryAssetStorageAPI().acquireCacheLease() : null) {
final File oldAsset = com.dotcms.storage.AssetStorageFeature.isEnabled()
? APILocator.getBinaryAssetStorageAPI().getBinaryFile(fileAssetWrapper.getAsset().getInode(), FileAssetAPI.BINARY_FIELD)
: new File(APILocator.getFileAssetAPI().getRealAssetPathIgnoreExtensionCase(fileAssetWrapper.getAsset().getInode(), fileAssetWrapper.getAsset().getUnderlyingFileName()));
if (oldAsset == null) {
// Fail the bundle, as a missing source file does with the flag off, so it can be retried.
throw new IOException(String.format(
"Binary file not found for File Asset inode '%s' (underlying name: '%s')",
fileAssetWrapper.getAsset().getInode(),
fileAssetWrapper.getAsset().getUnderlyingFileName()));
}
if(output.exists(filePath)) {
output.delete(filePath);
}

FileUtil.copyFile(oldAsset, output.getFile(filePath), true);
output.setLastModified(filePath, cal.getTimeInMillis());
FileUtil.copyFile(oldAsset, output.getFile(filePath), true);
output.setLastModified(filePath, cal.getTimeInMillis());
}
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -413,7 +413,7 @@ public void retry ( HttpServletRequest request, HttpServletResponse response ) t
Verify if the bundle exist and was created correctly..., meaning, if there is not a .tar.gz file is because
something happened on the creation of the bundle.
*/
File bundleFile = new File( ConfigUtils.getBundlePath() + File.separator + basicConfig.getId() + ".tar.gz" );
File bundleFile = com.dotcms.publishing.output.TarGzipBundleOutput.getBundleTarGzipFile(basicConfig.getId());
if ( !bundleFile.exists() ) {
Logger.warn( this.getClass(), "No Push Publish Bundle with id: " + bundleId + " found." );
appendMessage( responseMessage, "publisher_retry.error.not.found", bundleId, true );
Expand Down Expand Up @@ -520,7 +520,9 @@ public void downloadBundle ( HttpServletRequest request, HttpServletResponse res

ArrayList<File> list = new ArrayList<>( 1 );
list.add( bundleRoot );
File bundle = new File( bundleRoot + File.separator + ".." + File.separator + config.getId() + ".tar.gz" );
File bundle = com.dotcms.storage.AssetStorageFeature.isEnabled()
? com.dotcms.publishing.output.TarGzipBundleOutput.getBundleTarGzipFile(config.getId())
: new File( bundleRoot + File.separator + ".." + File.separator + config.getId() + ".tar.gz" );
if ( !bundle.exists() ) {
response.sendError( 500, "No Bundle Found" );
return;
Expand Down Expand Up @@ -634,7 +636,14 @@ public void downloadUnpushedBundle ( HttpServletRequest request, HttpServletResp
//Clean the just created bundle because on each download we will generate a new bundle file with a new id in order to avoid conflicts with ids
final File bundleRoot = BundlerUtil.getBundleRoot( bundleId );
final File compressedBundle = new File( ConfigUtils.getBundlePath() + File.separator + bundleId + ".tar.gz" );
if ( compressedBundle.exists() ) {
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
try {
com.dotcms.publishing.output.BundleArchiveStorage.getInstance().delete(bundleId);
} catch (DotDataException | RuntimeException e) {
Logger.error(this, "Unable to remove generated download bundle " + bundleId, e);
}
}
else if ( compressedBundle.exists() ) {
compressedBundle.delete();
if ( bundleRoot.exists() ) {
com.liferay.util.FileUtil.deltree( bundleRoot );
Expand Down Expand Up @@ -689,7 +698,11 @@ public void uploadBundle ( HttpServletRequest request, HttpServletResponse respo
status = PublishAuditAPI.getInstance().updateAuditTable( endpointId, endpointId, bundleFolder );

// Write file on FS
FileUtil.writeToFile( bundle, bundlePath + bundleName );
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
com.dotcms.publishing.output.BundleArchiveStorage.getInstance().receive(bundleName, bundle);
} else {
FileUtil.writeToFile( bundle, bundlePath + bundleName );
}

if ( !status.getStatus().equals( Status.PUBLISHING_BUNDLE ) ) {
PushPublisherJob.triggerPushPublisherJob(bundleName, status);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -88,10 +88,15 @@ public void setFilterKey(final String filterKey) {
}

/**
* Checks if the bundle was already generated based on the id: BUNDLE_ID.tar.gz
* Checks if the bundle was already generated based on the id: BUNDLE_ID.tar.gz. The bundle
* pages use this to decide whether to show download links. With S3 asset storage on, a storage
* failure is logged and reported as false so the page still renders.
* @return boolean - true if the bundle exists.
*/
public boolean bundleTgzExists() {
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
return com.dotcms.publishing.output.BundleArchiveStorage.getInstance().existsForDisplay(id);
}

return Try.of(()->new File( ConfigUtils.getBundlePath() + File.separator + id + ".tar.gz" ).exists()).getOrElse(false);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,9 @@ public Bundle getBundleById(String id) throws DotDataException {
@WrapInTransaction
@Override
public void deleteBundle(String id) throws DotDataException {
if (UtilMethods.isSet(id)) {
com.dotcms.publishing.output.BundleArchiveCleanupProcessor.enqueue(id);
}
bundleFactory.deleteBundle(id);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,9 @@ public PublisherConfig process ( final PublishStatus status ) throws DotPublishi
File bundleRoot = BundlerUtil.getBundleRoot(this.config.getName(), false);
final List<File> list = new ArrayList<>(1);
list.add(bundleRoot);
File bundleFile = new File(bundleRoot + ".tar.gz");
File bundleFile = com.dotcms.storage.AssetStorageFeature.isEnabled()
? com.dotcms.publishing.output.TarGzipBundleOutput.getBundleTarGzipFile(this.config.getId())
: new File(bundleRoot + ".tar.gz");

List<Environment> environments = APILocator.getEnvironmentAPI().findEnvironmentsByBundleId(this.config.getId());

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -199,7 +199,9 @@ public PublisherConfig process ( final PublishStatus status ) throws DotPublishi
// Extract file to a directory
InputStream bundleIS = null;
try {
bundleIS = Files.newInputStream(Paths.get(bundlePath + bundleName));
bundleIS = com.dotcms.storage.AssetStorageFeature.isEnabled()
? Files.newInputStream(com.dotcms.publishing.output.TarGzipBundleOutput.getBundleTarGzipFile(bundleID).toPath())
: Files.newInputStream(Paths.get(bundlePath + bundleName));
untar(bundleIS, folderOut.getAbsolutePath() + File.separator + bundleName, bundleName);
} finally {
CloseUtils.closeQuietly(bundleIS);
Expand Down
17 changes: 17 additions & 0 deletions dotCMS/src/main/java/com/dotcms/publishing/BundlerUtil.java
Original file line number Diff line number Diff line change
Expand Up @@ -136,12 +136,23 @@ public static boolean isRetryable(final String bundleId) {
}
}

/**
* Tells {@link #isRetryable(String)} whether the bundle's files are present. With S3 asset
* storage on, the archive check is the display check, which reports a storage failure as
* false instead of throwing, because the audit detail page calls this while rendering.
*
* @param bundleId the bundle id
* @return true if the bundle's files are present
*/
private static boolean bundleExists(String bundleId) {
final PublisherConfig basicConfig = new PublisherConfig();
basicConfig.setId(bundleId);
final File bundleRoot = BundlerUtil.getBundleRoot( basicConfig.getName(), false );

final File bundleStaticFile = new File(bundleRoot.getAbsolutePath() + PublisherConfig.STATIC_SUFFIX);
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
return bundleStaticFile.exists() || com.dotcms.publishing.output.BundleArchiveStorage.getInstance().existsForDisplay(bundleId);
}
if ( !bundleStaticFile.exists() ) {
return true;
}
Expand Down Expand Up @@ -210,6 +221,9 @@ public static void writeBundleXML(final PublisherConfig config, final BundleOutp
try (final OutputStream outputStream = output.addFile(bundleXmlFilePath)) {
objectToXML(config, outputStream);
} catch ( IOException e ) {
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
throw new DotRuntimeException("Unable to include bundle descriptor", e);
}
Logger.error( BundlerUtil.class, e.getMessage(), e );
}
}
Expand Down Expand Up @@ -522,6 +536,9 @@ public static String sanitizeBundleName(String bundleName) throws DotPublisherEx
}

public static boolean tarGzipExists(final String bundleId) {
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
return com.dotcms.publishing.output.BundleArchiveStorage.getInstance().exists(bundleId);
}
final File bundleTarGzip = TarGzipBundleOutput.getBundleTarGzipFile(bundleId);
return bundleTarGzip.exists();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,9 @@ public PublishStatus publish ( PublisherConfig config, PublishStatus status) thr
} else {
addBundleXMLIntoBundle(config, output);
}
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
output.complete();
}
} else {
Logger.info(this, "Retrying bundle: " + config.getId()
+ ", we don't need to run bundlers again");
Expand Down Expand Up @@ -221,6 +224,9 @@ private void addManifestIntoBundleOutput(final BundleOutput output,
manifestBuilder.close();
output.copyFile(manifestFile, File.separator + ManifestBuilder.MANIFEST_NAME);
} catch (final IOException e) {
if (com.dotcms.storage.AssetStorageFeature.isEnabled()) {
throw new com.dotmarketing.exception.DotRuntimeException("Unable to include bundle manifest", e);
}
Logger.error(PublisherAPIImpl.class, "Error trying to copy the manifest file: " +
e.getMessage());
}
Expand Down Expand Up @@ -567,4 +573,4 @@ private FilterDescriptor createFilterFromFile(final Path path) {
}
}

}
}
Loading
Loading