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}