001/*
002 * This library is part of OpenCms -
003 * The Open Source Content Management System
004 *
005 * Copyright (c) Alkacon Software GmbH & Co. KG (https://www.alkacon.com)
006 *
007 * This library is free software; you can redistribute it and/or
008 * modify it under the terms of the GNU Lesser General Public
009 * License as published by the Free Software Foundation; either
010 * version 2.1 of the License, or (at your option) any later version.
011 *
012 * This library is distributed in the hope that it will be useful,
013 * but WITHOUT ANY WARRANTY; without even the implied warranty of
014 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
015 * Lesser General Public License for more details.
016 *
017 * For further information about Alkacon Software GmbH & Co. KG, please see the
018 * company website: https://www.alkacon.com
019 *
020 * For further information about OpenCms, please see the
021 * project website: https://www.opencms.org
022 *
023 * You should have received a copy of the GNU Lesser General Public
024 * License along with this library; if not, write to the Free Software
025 * Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA  02111-1307  USA
026 */
027
028package org.opencms.loader.imagecache;
029
030import org.opencms.configuration.CmsImageCacheConfiguration;
031import org.opencms.db.storage.s3.CmsGenericS3Client;
032import org.opencms.db.storage.s3.CmsS3DeleteResult;
033import org.opencms.loader.CmsS3ImageCache;
034import org.opencms.loader.imagecache.CmsImageCacheCapabilities.Capability;
035
036import java.time.Duration;
037import java.time.Instant;
038import java.util.ArrayList;
039import java.util.List;
040import java.util.concurrent.ExecutionException;
041import java.util.concurrent.ExecutorService;
042import java.util.concurrent.Executors;
043import java.util.concurrent.Future;
044import java.util.concurrent.Semaphore;
045import java.util.concurrent.TimeUnit;
046
047/**
048 * Maintenance adapter for the S3 image cache.<p>
049 */
050public class CmsS3ImageCacheMaintenance implements I_CmsImageCacheMaintenance {
051
052    /** Result for one delete batch. */
053    private static final class BatchResult {
054
055        /** The batch entries. */
056        private final List<CmsImageCacheEntry> m_entries;
057
058        /** A whole-request failure. */
059        private final Exception m_failure;
060
061        /** The per-object result. */
062        private final CmsS3DeleteResult m_result;
063
064        /** Creates a batch result. */
065        BatchResult(List<CmsImageCacheEntry> entries, CmsS3DeleteResult result, Exception failure) {
066
067            m_entries = entries;
068            m_result = result;
069            m_failure = failure;
070        }
071    }
072
073    /** Rate limiter shared by concurrent S3 copy workers. */
074    private static final class CopyRateLimiter {
075
076        /** Minimum interval between copy starts. */
077        private final long m_intervalNanos;
078
079        /** The next reserved copy start. */
080        private long m_nextPermitNanos;
081
082        /** Creates a limiter. */
083        CopyRateLimiter(int maxCopiesPerSecond) {
084
085            m_intervalNanos = Math.max(1L, TimeUnit.SECONDS.toNanos(1) / maxCopiesPerSecond);
086        }
087
088        /** Waits until the next copy may start. */
089        void acquire() throws InterruptedException {
090
091            long delay;
092            synchronized (this) {
093                long now = System.nanoTime();
094                long scheduled = Math.max(now, m_nextPermitNanos);
095                delay = scheduled - now;
096                m_nextPermitNanos = scheduled + m_intervalNanos;
097            }
098            if (delay > 0) {
099                TimeUnit.NANOSECONDS.sleep(delay);
100            }
101        }
102    }
103
104    /** Result for one renewal. */
105    private static final class RenewalResult {
106
107        /** The renewed entry. */
108        private final CmsImageCacheEntry m_entry;
109
110        /** The renewal failure. */
111        private final Exception m_failure;
112
113        /** Whether renewal was skipped. */
114        private final boolean m_skipped;
115
116        /** Creates a renewal result. */
117        RenewalResult(CmsImageCacheEntry entry, boolean skipped, Exception failure) {
118
119            m_entry = entry;
120            m_skipped = skipped;
121            m_failure = failure;
122        }
123    }
124
125    /** The supported capabilities. */
126    private static final CmsImageCacheCapabilities CAPABILITIES = CmsImageCacheCapabilities.of(
127        Capability.LIST_ENTRIES,
128        Capability.ENTRY_TIMESTAMPS,
129        Capability.DELETE_ENTRIES,
130        Capability.CLEAR,
131        Capability.RENEW_ENTRIES);
132
133    /** The S3 cache. */
134    private final CmsS3ImageCache m_cache;
135
136    /** The number of concurrent copy requests. */
137    private final int m_copyConcurrency;
138
139    /** The number of concurrent delete requests. */
140    private final int m_deleteConcurrency;
141
142    /** The maximum number of keys per delete request. */
143    private final int m_deleteBatchSize;
144
145    /** Rate limiter shared by all renewal requests handled by this adapter. */
146    private final CopyRateLimiter m_copyRateLimiter;
147
148    /** The maximum number of waiting renewal tasks. */
149    private final int m_renewalQueueCapacity;
150
151    /**
152     * Creates a maintenance adapter.<p>
153     *
154     * @param cache the S3 cache
155     */
156    public CmsS3ImageCacheMaintenance(CmsS3ImageCache cache) {
157
158        this(
159            cache,
160            CmsImageCacheConfiguration.DEFAULT_S3_DELETE_BATCH_SIZE,
161            CmsImageCacheConfiguration.DEFAULT_S3_DELETE_CONCURRENCY,
162            CmsImageCacheConfiguration.DEFAULT_S3_COPY_CONCURRENCY,
163            CmsImageCacheConfiguration.DEFAULT_S3_MAX_COPIES_PER_SECOND,
164            CmsImageCacheConfiguration.DEFAULT_S3_RENEWAL_QUEUE_CAPACITY);
165    }
166
167    /**
168     * Creates a maintenance adapter.<p>
169     *
170     * @param cache the S3 cache
171     * @param deleteBatchSize the maximum number of keys per delete request
172     * @param deleteConcurrency the number of concurrent delete requests
173     */
174    public CmsS3ImageCacheMaintenance(CmsS3ImageCache cache, int deleteBatchSize, int deleteConcurrency) {
175
176        this(
177            cache,
178            deleteBatchSize,
179            deleteConcurrency,
180            CmsImageCacheConfiguration.DEFAULT_S3_COPY_CONCURRENCY,
181            CmsImageCacheConfiguration.DEFAULT_S3_MAX_COPIES_PER_SECOND,
182            CmsImageCacheConfiguration.DEFAULT_S3_RENEWAL_QUEUE_CAPACITY);
183    }
184
185    /**
186     * Creates a maintenance adapter.<p>
187     *
188     * @param cache the S3 cache
189     * @param deleteBatchSize the maximum number of keys per delete request
190     * @param deleteConcurrency the number of concurrent delete requests
191     * @param copyConcurrency the number of concurrent copy requests
192     * @param maxCopiesPerSecond the maximum number of copy requests started per second
193     * @param renewalQueueCapacity the maximum number of waiting renewal tasks
194     */
195    public CmsS3ImageCacheMaintenance(
196        CmsS3ImageCache cache,
197        int deleteBatchSize,
198        int deleteConcurrency,
199        int copyConcurrency,
200        int maxCopiesPerSecond,
201        int renewalQueueCapacity) {
202
203        if ((deleteBatchSize < 1) || (deleteBatchSize > CmsGenericS3Client.MAX_DELETE_OBJECTS)) {
204            throw new IllegalArgumentException(
205                "S3 delete batch size must be between 1 and " + CmsGenericS3Client.MAX_DELETE_OBJECTS + ".");
206        }
207        if (deleteConcurrency < 1) {
208            throw new IllegalArgumentException("S3 delete concurrency must be positive.");
209        }
210        if (copyConcurrency < 1) {
211            throw new IllegalArgumentException("S3 copy concurrency must be positive.");
212        }
213        if (maxCopiesPerSecond < 1) {
214            throw new IllegalArgumentException("S3 maximum copies per second must be positive.");
215        }
216        if (renewalQueueCapacity < 1) {
217            throw new IllegalArgumentException("S3 renewal queue capacity must be positive.");
218        }
219        if (((long)copyConcurrency + renewalQueueCapacity) > Integer.MAX_VALUE) {
220            throw new IllegalArgumentException("S3 copy concurrency and renewal queue capacity are too large.");
221        }
222        m_cache = cache;
223        m_deleteBatchSize = deleteBatchSize;
224        m_deleteConcurrency = deleteConcurrency;
225        m_copyConcurrency = copyConcurrency;
226        m_copyRateLimiter = new CopyRateLimiter(maxCopiesPerSecond);
227        m_renewalQueueCapacity = renewalQueueCapacity;
228    }
229
230    /**
231     * @see org.opencms.loader.imagecache.I_CmsImageCacheMaintenance#execute(org.opencms.loader.imagecache.CmsImageCacheMaintenanceRequest)
232     */
233    @Override
234    public CmsImageCacheMaintenanceResult execute(CmsImageCacheMaintenanceRequest request) throws Exception {
235
236        long start = System.nanoTime();
237        int requested = request.getOperation() == CmsImageCacheMaintenanceRequest.Operation.CLEAR
238        ? 1
239        : request.getEntries().size();
240        CmsImageCacheMaintenanceResult.Builder result = new CmsImageCacheMaintenanceResult.Builder(
241            request.getOperation(),
242            requested);
243        switch (request.getOperation()) {
244            case CLEAR:
245                try {
246                    m_cache.clear(m_deleteBatchSize);
247                    result.addSuccess();
248                } catch (Exception e) {
249                    result.addFailure(getBackendId(), e);
250                }
251                break;
252            case DELETE:
253                deleteEntries(request.getEntries(), result);
254                break;
255            case RENEW:
256                renewEntries(request.getEntries(), request.getRenewalTime(), result);
257                break;
258            default:
259                throw new IllegalArgumentException(
260                    "Unsupported S3 image cache maintenance request: " + request.getOperation());
261        }
262        return result.build(Duration.ofNanos(System.nanoTime() - start));
263    }
264
265    /**
266     * @see org.opencms.loader.imagecache.I_CmsImageCacheMaintenance#getBackendId()
267     */
268    @Override
269    public String getBackendId() {
270
271        return "s3";
272    }
273
274    /**
275     * @see org.opencms.loader.imagecache.I_CmsImageCacheMaintenance#getCapabilities()
276     */
277    @Override
278    public CmsImageCacheCapabilities getCapabilities() {
279
280        return CAPABILITIES;
281    }
282
283    /**
284     * @see org.opencms.loader.imagecache.I_CmsImageCacheMaintenance#getEntry(java.lang.String)
285     */
286    @Override
287    public CmsImageCacheEntry getEntry(String key) throws Exception {
288
289        org.opencms.db.storage.s3.CmsS3ObjectMetadata metadata;
290        try {
291            metadata = m_cache.getMetadata(key);
292        } catch (org.opencms.db.storage.CmsStorageBlobNotFoundException e) {
293            return null;
294        }
295        return new CmsImageCacheEntry(
296            metadata.getKey(),
297            metadata.getLength(),
298            metadata.getLastModified(),
299            metadata.getRevision());
300    }
301
302    /**
303     * @see org.opencms.loader.imagecache.I_CmsImageCacheMaintenance#visitEntries(org.opencms.loader.imagecache.I_CmsImageCacheMaintenanceEntryVisitor)
304     */
305    @Override
306    public void visitEntries(I_CmsImageCacheMaintenanceEntryVisitor visitor) throws Exception {
307
308        m_cache.visitEntriesWithMetadata(
309            metadata -> visitor.visit(
310                new CmsImageCacheEntry(
311                    metadata.getKey(),
312                    metadata.getLength(),
313                    metadata.getLastModified(),
314                    metadata.getRevision())));
315    }
316
317    /** Adds a batch outcome to the maintenance metrics. */
318    private void addBatchResult(BatchResult batch, CmsImageCacheMaintenanceResult.Builder result) {
319
320        for (CmsImageCacheEntry entry : batch.m_entries) {
321            if (batch.m_failure != null) {
322                result.addFailure(entry.getKey(), batch.m_failure);
323                continue;
324            }
325            Exception failure = batch.m_result.getFailures().get(normalizeKey(entry.getKey()));
326            if (failure == null) {
327                result.addSuccess();
328            } else {
329                result.addFailure(entry.getKey(), failure);
330            }
331        }
332    }
333
334    /** Adds one renewal outcome to the maintenance metrics. */
335    private void addRenewalResult(RenewalResult renewal, CmsImageCacheMaintenanceResult.Builder result) {
336
337        if (renewal.m_failure != null) {
338            result.addFailure(renewal.m_entry.getKey(), renewal.m_failure);
339        } else if (renewal.m_skipped) {
340            result.addSkipped();
341        } else {
342            result.addSuccess();
343        }
344    }
345
346    /** Executes one S3 delete batch and retains whole-request failures as per-entry failures. */
347    private BatchResult deleteBatch(List<CmsImageCacheEntry> entries) {
348
349        List<String> keys = new ArrayList<String>(entries.size());
350        for (CmsImageCacheEntry entry : entries) {
351            keys.add(entry.getKey());
352        }
353        try {
354            return new BatchResult(entries, m_cache.deleteBatch(keys), null);
355        } catch (Exception e) {
356            return new BatchResult(entries, null, e);
357        }
358    }
359
360    /** Deletes selected entries using concurrent S3 multi-object requests. */
361    private void deleteEntries(List<CmsImageCacheEntry> entries, CmsImageCacheMaintenanceResult.Builder result)
362    throws Exception {
363
364        List<List<CmsImageCacheEntry>> batches = new ArrayList<List<CmsImageCacheEntry>>();
365        for (int start = 0; start < entries.size(); start += m_deleteBatchSize) {
366            int end = Math.min(entries.size(), start + m_deleteBatchSize);
367            batches.add(new ArrayList<CmsImageCacheEntry>(entries.subList(start, end)));
368        }
369        if ((m_deleteConcurrency == 1) || (batches.size() < 2)) {
370            for (List<CmsImageCacheEntry> batch : batches) {
371                addBatchResult(deleteBatch(batch), result);
372            }
373            return;
374        }
375        ExecutorService executor = Executors.newFixedThreadPool(Math.min(m_deleteConcurrency, batches.size()));
376        try {
377            List<Future<BatchResult>> futures = new ArrayList<Future<BatchResult>>(batches.size());
378            for (List<CmsImageCacheEntry> batch : batches) {
379                futures.add(executor.submit(() -> deleteBatch(batch)));
380            }
381            for (Future<BatchResult> future : futures) {
382                try {
383                    addBatchResult(future.get(), result);
384                } catch (ExecutionException e) {
385                    Throwable cause = e.getCause();
386                    if (cause instanceof Exception) {
387                        throw (Exception)cause;
388                    }
389                    throw e;
390                }
391            }
392        } catch (InterruptedException e) {
393            Thread.currentThread().interrupt();
394            throw e;
395        } finally {
396            executor.shutdownNow();
397        }
398    }
399
400    /** Normalizes a cache key for matching a batch result. */
401    private String normalizeKey(String key) {
402
403        String result = key;
404        while (result.startsWith("/")) {
405            result = result.substring(1);
406        }
407        return result;
408    }
409
410    /** Renews selected entries with bounded concurrency, queueing and copy rate. */
411    private void renewEntries(
412        List<CmsImageCacheEntry> entries,
413        Instant renewalTime,
414        CmsImageCacheMaintenanceResult.Builder result)
415    throws Exception {
416
417        if (entries.isEmpty()) {
418            return;
419        }
420        if ((m_copyConcurrency == 1) || (entries.size() < 2)) {
421            for (CmsImageCacheEntry entry : entries) {
422                addRenewalResult(renewEntry(entry, renewalTime, m_copyRateLimiter), result);
423            }
424            return;
425        }
426        int maximumOutstanding = m_copyConcurrency + m_renewalQueueCapacity;
427        Semaphore outstanding = new Semaphore(maximumOutstanding);
428        ExecutorService executor = Executors.newFixedThreadPool(Math.min(m_copyConcurrency, entries.size()));
429        try {
430            List<Future<RenewalResult>> futures = new ArrayList<Future<RenewalResult>>(entries.size());
431            for (CmsImageCacheEntry entry : entries) {
432                outstanding.acquire();
433                try {
434                    futures.add(executor.submit(() -> {
435                        try {
436                            return renewEntry(entry, renewalTime, m_copyRateLimiter);
437                        } finally {
438                            outstanding.release();
439                        }
440                    }));
441                } catch (RuntimeException e) {
442                    outstanding.release();
443                    throw e;
444                }
445            }
446            for (Future<RenewalResult> future : futures) {
447                try {
448                    addRenewalResult(future.get(), result);
449                } catch (ExecutionException e) {
450                    Throwable cause = e.getCause();
451                    if (cause instanceof Exception) {
452                        throw (Exception)cause;
453                    }
454                    throw e;
455                }
456            }
457        } catch (InterruptedException e) {
458            Thread.currentThread().interrupt();
459            throw e;
460        } finally {
461            executor.shutdownNow();
462        }
463    }
464
465    /** Renews one entry with a conditional S3 self-copy. */
466    private RenewalResult renewEntry(CmsImageCacheEntry entry, Instant renewalTime, CopyRateLimiter rateLimiter)
467    throws InterruptedException {
468
469        if ((entry.getRevision() == null) || entry.getRevision().trim().isEmpty()) {
470            return new RenewalResult(entry, true, null);
471        }
472        rateLimiter.acquire();
473        try {
474            boolean renewed = m_cache.renew(entry.getKey(), entry.getRevision(), renewalTime);
475            return new RenewalResult(entry, !renewed, null);
476        } catch (Exception e) {
477            return new RenewalResult(entry, false, e);
478        }
479    }
480}