From 2d174bcddfd711168290b9406102dbfbc1a69cfb Mon Sep 17 00:00:00 2001 From: fabrizzio-dotCMS Date: Wed, 5 Aug 2026 13:37:43 -0600 Subject: [PATCH 1/3] fix(opensearch): surface a shadow store that stopped accepting writes (#36222) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit In dual-write phases an OpenSearch write failure is fire-and-forget, and the reindex listener logged one WARN per failed item with nothing else. When the OS user loses a permission every document of every batch is rejected with the same message, which produced two problems observed while QA-ing #36222 (TC-056): ~900 identical warnings in one minute (a real reindex emits hundreds of thousands, burying every other line), and no signal at all that the shadow store had stopped receiving writes — so advancing the migration phase would silently switch reads to a frozen index. - BulkProcessorListener: aggregate a batch's failures into one line per distinct message with its count, keep the per-item detail at DEBUG, and log a single actionable ERROR when the batch was rejected in full for a systemic cause (rate-limited via REINDEX_SHADOW_FAILURE_ESCALATION_MINUTES, default 10). Partial failures and unclassifiable causes stay silent: those are per-document problems, not migration blockers. The phase is NOT reset here — a bulk rejection can be scoped to one index, so halting stays the operator's call. - OSIndexAPIImpl: classify a bare failure message, not only an exception chain (a rejected bulk item carries text, never a status code), and recognise the security plugin's own wording (security_exception / no permissions for). - IndexConfigHelper: neutral seam so the reindex listener escalates without importing the OpenSearch adapter's classifier. Co-Authored-By: Claude Opus 5 (1M context) --- .../content/index/IndexConfigHelper.java | 25 +++ .../index/opensearch/OSIndexAPIImpl.java | 106 ++++++---- .../common/reindex/BulkProcessorListener.java | 138 ++++++++++++- .../OSIndexAPIImplConnectionClassifyTest.java | 29 +++ ...ulkProcessorListenerShadowFailureTest.java | 185 ++++++++++++++++++ 5 files changed, 441 insertions(+), 42 deletions(-) create mode 100644 dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java diff --git a/dotCMS/src/main/java/com/dotcms/content/index/IndexConfigHelper.java b/dotCMS/src/main/java/com/dotcms/content/index/IndexConfigHelper.java index 75d5ba472aa1..c62006732beb 100644 --- a/dotCMS/src/main/java/com/dotcms/content/index/IndexConfigHelper.java +++ b/dotCMS/src/main/java/com/dotcms/content/index/IndexConfigHelper.java @@ -1,9 +1,12 @@ package com.dotcms.content.index; +import com.dotcms.content.index.opensearch.OSIndexAPIImpl; +import com.dotcms.content.index.opensearch.OSIndexAPIImpl.ConnectionFailureKind; import com.dotcms.content.index.opensearch.OSIndexProperty; import com.dotcms.featureflag.FeatureFlagName; import com.dotmarketing.util.Config; import com.dotmarketing.util.Logger; +import java.util.Optional; /** * Central helper for reading index-layer configuration properties. @@ -65,6 +68,28 @@ static void logShadowWriteFailure(final Class clazz, } } + /** + * Returns the operator-facing remediation for a shadow-write failure whose cause is + * systemic — a connection, TLS or permission problem that keeps rejecting every write + * until it is fixed — or {@link Optional#empty()} when the cause cannot be classified, which is + * the signature of a per-document problem (a mapping conflict, a malformed field) that must not + * be escalated as a migration blocker. + * + *

Exists as a neutral seam so callers outside the index layer — the reindex bulk listener, + * which deliberately carries no vendor imports — can escalate a systemic rejection without + * reaching into the OpenSearch adapter's classifier themselves (issue #36222 follow-up).

+ * + * @param failureMessage vendor-reported failure text, or {@code null} + * @return {@code KIND (remediation)} for a systemic cause, empty otherwise + */ + static Optional systemicFailureRemediation(final String failureMessage) { + final ConnectionFailureKind kind = OSIndexAPIImpl.classifyFailureMessage(failureMessage); + if (kind == ConnectionFailureKind.UNKNOWN) { + return Optional.empty(); + } + return Optional.of(kind.name() + " (" + kind.remediation() + ")"); + } + // ------------------------------------------------------------------------- // Migration phase // ------------------------------------------------------------------------- diff --git a/dotCMS/src/main/java/com/dotcms/content/index/opensearch/OSIndexAPIImpl.java b/dotCMS/src/main/java/com/dotcms/content/index/opensearch/OSIndexAPIImpl.java index a6a5748e8301..aa252a2da505 100644 --- a/dotCMS/src/main/java/com/dotcms/content/index/opensearch/OSIndexAPIImpl.java +++ b/dotCMS/src/main/java/com/dotcms/content/index/opensearch/OSIndexAPIImpl.java @@ -723,48 +723,80 @@ public String remediation() { */ public static ConnectionFailureKind classifyConnectionError(final Throwable error) { for (Throwable t = error; t != null; t = t.getCause()) { - final String type = t.getClass().getName().toLowerCase(Locale.ROOT); - final String msg = t.getMessage() == null - ? "" : t.getMessage().toLowerCase(Locale.ROOT); - - // TLS / scheme mismatch: SSL handshake errors, or a connection closed with no response - // (the classic symptom of speaking http:// to an https-only port). - if (type.contains("sslexception") - || type.contains("sslhandshake") - || type.contains("connectionclosedexception") - || msg.contains("unrecognized ssl message") - || msg.contains("plaintext") - || msg.contains("ssl") - || msg.contains("certificat") // certificate / certification path (PKIX) - || msg.contains("pkix") - || msg.contains("connection closed")) { - return ConnectionFailureKind.TLS_SCHEME_MISMATCH; - } - // Authentication / authorization. The status code must not be matched as a bare - // substring: dotCMS wrapper messages embed the physical index name, whose - // _yyyyMMddHHmmss timestamp regularly contains the digits 401/403 (e.g. an index - // created at 12:04:03 → working_20260728120403.os), which would misreport an unrelated - // failure — a settings-parse or orphan-delete error — as a permission problem and send - // the operator to change DOT_DOTCMS_CLUSTER_ID for nothing (issue #36222). - if (AUTH_STATUS_CODE.matcher(msg).find() - || msg.contains("unauthorized") || msg.contains("forbidden")) { - return ConnectionFailureKind.AUTH_FORBIDDEN; - } - // Host/port unreachable. - if (type.contains("connectexception") - || type.contains("unknownhostexception") - || type.contains("sockettimeoutexception") - || type.contains("nohttpresponseexception") - || msg.contains("connection refused") - || msg.contains("timed out") - || msg.contains("timeout") - || msg.contains("unknown host")) { - return ConnectionFailureKind.UNREACHABLE; + final ConnectionFailureKind kind = classify(t.getClass().getName(), t.getMessage()); + if (kind != ConnectionFailureKind.UNKNOWN) { + return kind; } } return ConnectionFailureKind.UNKNOWN; } + /** + * Classifies a bare failure message, for rejections that never reach the caller as an + * exception. The motivating case is a bulk write: the OpenSearch client reports a per-item + * rejection as text on the bulk item (e.g. {@code security_exception: no permissions for + * [indices:data/write/bulk[s]] …}) rather than throwing, so a systemic permission problem is + * indistinguishable from a per-document problem unless the message itself is classified + * (issue #36222 follow-up). + * + * @param failureMessage the vendor-reported failure text, or {@code null} + * @return the inferred {@link ConnectionFailureKind}; never {@code null} + */ + public static ConnectionFailureKind classifyFailureMessage(final String failureMessage) { + return classify("", failureMessage); + } + + /** + * Matching core shared by {@link #classifyConnectionError(Throwable)} (once per link of the + * exception chain) and {@link #classifyFailureMessage(String)}. + * + * @param typeName fully-qualified exception type name, or {@code ""} when classifying a message + * @param message the failure message, or {@code null} + */ + private static ConnectionFailureKind classify(final String typeName, final String message) { + final String type = typeName == null ? "" : typeName.toLowerCase(Locale.ROOT); + final String msg = message == null ? "" : message.toLowerCase(Locale.ROOT); + + // TLS / scheme mismatch: SSL handshake errors, or a connection closed with no response + // (the classic symptom of speaking http:// to an https-only port). + if (type.contains("sslexception") + || type.contains("sslhandshake") + || type.contains("connectionclosedexception") + || msg.contains("unrecognized ssl message") + || msg.contains("plaintext") + || msg.contains("ssl") + || msg.contains("certificat") // certificate / certification path (PKIX) + || msg.contains("pkix") + || msg.contains("connection closed")) { + return ConnectionFailureKind.TLS_SCHEME_MISMATCH; + } + // Authentication / authorization. The status code must not be matched as a bare + // substring: dotCMS wrapper messages embed the physical index name, whose + // _yyyyMMddHHmmss timestamp regularly contains the digits 401/403 (e.g. an index + // created at 12:04:03 → working_20260728120403.os), which would misreport an unrelated + // failure — a settings-parse or orphan-delete error — as a permission problem and send + // the operator to change DOT_DOTCMS_CLUSTER_ID for nothing (issue #36222). + // The security plugin's own wording is matched too: a rejected bulk item never carries a + // status code, only "security_exception: no permissions for [action] and User [...]". + if (AUTH_STATUS_CODE.matcher(msg).find() + || msg.contains("unauthorized") || msg.contains("forbidden") + || msg.contains("security_exception") || msg.contains("no permissions for")) { + return ConnectionFailureKind.AUTH_FORBIDDEN; + } + // Host/port unreachable. + if (type.contains("connectexception") + || type.contains("unknownhostexception") + || type.contains("sockettimeoutexception") + || type.contains("nohttpresponseexception") + || msg.contains("connection refused") + || msg.contains("timed out") + || msg.contains("timeout") + || msg.contains("unknown host")) { + return ConnectionFailureKind.UNREACHABLE; + } + return ConnectionFailureKind.UNKNOWN; + } + // ========================================================================= // Alias management diff --git a/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java b/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java index c3420f8491b1..9a743647d590 100644 --- a/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java +++ b/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java @@ -12,13 +12,17 @@ import com.dotmarketing.exception.DotDataException; import com.dotmarketing.util.Config; import com.dotmarketing.util.Logger; +import com.dotmarketing.util.UtilMethods; import com.liferay.util.StringPool; import io.vavr.control.Try; import java.util.ArrayList; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; /** * {@link IndexBulkListener} that handles the business logic before/after reindexing content. @@ -36,6 +40,22 @@ public class BulkProcessorListener implements IndexBulkListener { static final List RESERVED_IDS = List.of(Host.SYSTEM_HOST); + /** + * Config key: minutes between two systemic-rejection escalations + * ({@link #systemicFailureEscalation(IndexTag, int, Map)}). Default: 10. + */ + static final String SHADOW_ESCALATION_INTERVAL_KEY = + "REINDEX_SHADOW_FAILURE_ESCALATION_MINUTES"; + + /** Stand-in used when the vendor reports a failed item without any message. */ + static final String NO_FAILURE_MESSAGE = "(no failure message reported)"; + + /** + * Epoch millis of the last systemic-rejection escalation. Static on purpose: a fresh listener is + * created for every batch, so per-instance state could not rate-limit anything. + */ + private static final AtomicLong LAST_SYSTEMIC_ESCALATION = new AtomicLong(0L); + private volatile long contentletsIndexed; private int lastBatchSize; @@ -110,11 +130,12 @@ public void beforeBulk(final long executionId, final int actionCount) { @Override public void afterBulk(final long executionId, final List results) { if (shadow) { - // OS shadow — fire-and-forget; log individual failures for observability only - results.stream() - .filter(IndexBulkItemResult::failed) - .forEach(r -> logShadowWriteFailure(this.getClass(), - "[OS] Index failure (fire-and-forget): " + r.failureMessage(), null)); + // OS shadow — fire-and-forget, but summarised: a systemic rejection (permissions, TLS, + // unreachable cluster) repeats verbatim for every document in the batch, so logging one + // line per failed item buried the log — a real reindex emits hundreds of thousands of + // identical entries, hiding every other line including the actionable one below + // (observed on issue #36222, TC-056: ~900 identical WARNs in one minute). + logShadowBatchFailures(results); return; } Logger.debug(this.getClass(), "Bulk process completed"); @@ -165,6 +186,113 @@ public void afterBulk(final long executionId, final Throwable failure) { workingRecords.values().forEach(idx -> handleFailure(idx, msg)); } + /** + * Logs the failed items of a shadow bulk batch as one aggregated entry per distinct failure + * message, keeps the per-item detail at {@code DEBUG}, and escalates once when the whole batch + * was rejected for a systemic reason. + * + * @param results every item result of the completed batch, successful ones included + */ + private void logShadowBatchFailures(final List results) { + final Map failuresByMessage = summarizeFailures(results); + if (failuresByMessage.isEmpty()) { + return; + } + final String tag = "[" + provider.name() + "] "; + failuresByMessage.forEach((message, count) -> logShadowWriteFailure(this.getClass(), + tag + "Index failure (fire-and-forget): " + count + " of " + results.size() + + " item(s) in this batch — " + message, null)); + + if (Logger.isDebugEnabled(this.getClass())) { + results.stream() + .filter(IndexBulkItemResult::failed) + .forEach(result -> Logger.debug(this.getClass(), tag + "Failed item id=" + + result.id() + ": " + result.failureMessage())); + } + + final Optional escalation = + systemicFailureEscalation(provider, results.size(), failuresByMessage); + if (escalation.isPresent() && shouldEscalateNow()) { + Logger.error(this.getClass(), escalation.get()); + } + } + + /** + * Counts the failed items of a batch grouped by failure message, preserving first-seen order. + * A systemic cause produces a single entry whose count equals the batch size; a per-document + * cause produces several entries, or one entry that covers only part of the batch. + * + * @param results every item result of the completed batch + * @return failure message → number of items that failed with it; empty when nothing failed + */ + static Map summarizeFailures(final List results) { + final Map failuresByMessage = new LinkedHashMap<>(); + for (final IndexBulkItemResult result : results) { + if (result.failed()) { + final String message = result.failureMessage(); + failuresByMessage.merge(UtilMethods.isSet(message) ? message : NO_FAILURE_MESSAGE, + 1L, Long::sum); + } + } + return failuresByMessage; + } + + /** + * Builds the actionable message for a shadow batch that was rejected in full for a + * systemic reason — the state where the shadow store silently stops receiving writes and + * diverges from the authoritative one, which must never be promoted by advancing the migration + * phase (issue #36222 follow-up: index creation was already covered by + * {@code ContentletIndexAPIImpl.handleOsBootstrapFailure}, the write path was not). + * + *

Deliberately silent unless every item failed: a partial failure is a per-document + * problem, and unclassifiable messages (mapping conflicts) are not migration blockers either. + * The phase is not reset here — a bulk rejection can be scoped to a single index, so the + * decision to halt is left to the operator, with this line and the readiness report as the + * signal.

+ * + * @param provider the shadow provider that rejected the batch + * @param batchSize total items in the batch, successful ones included + * @param failuresByMessage output of {@link #summarizeFailures(List)} + * @return the escalation message, or empty when this batch does not warrant one + */ + static Optional systemicFailureEscalation(final IndexTag provider, final int batchSize, + final Map failuresByMessage) { + final long failed = failuresByMessage.values().stream().mapToLong(Long::longValue).sum(); + if (batchSize <= 0 || failed < batchSize) { + return Optional.empty(); + } + return failuresByMessage.entrySet().stream() + .max(Map.Entry.comparingByValue()) + .flatMap(dominant -> IndexConfigHelper + .systemicFailureRemediation(dominant.getKey()) + .map(remediation -> "[" + provider.name() + "] EVERY document in this bulk" + + " batch (" + batchSize + " item(s)) was rejected by " + + provider.name() + " — likelyCause=" + remediation + + ". The shadow store is NOT receiving writes and is diverging from" + + " the authoritative store, so it must not be promoted: advancing" + + " the migration phase would serve reads from a stale index." + + " Fix the cause above, then run a full reindex to resynchronise." + + " Rejection: " + dominant.getKey())); + } + + /** + * Rate-limits the systemic escalation to one entry per + * {@value #SHADOW_ESCALATION_INTERVAL_KEY} minutes. Static because a fresh listener is built + * for every batch, so per-instance state would not throttle anything. + * + * @return {@code true} when this caller won the right to log the escalation now + */ + private static boolean shouldEscalateNow() { + final long intervalMs = TimeUnit.MINUTES.toMillis( + Config.getIntProperty(SHADOW_ESCALATION_INTERVAL_KEY, 10)); + final long now = System.currentTimeMillis(); + final long previous = LAST_SYSTEMIC_ESCALATION.get(); + if (previous != 0L && now - previous < intervalMs) { + return false; + } + return LAST_SYSTEMIC_ESCALATION.compareAndSet(previous, now); + } + static String getMatchingReservedIdIfAny(final String id) { for (final String reservedId : RESERVED_IDS) { if (id.contains(reservedId)) { diff --git a/dotCMS/src/test/java/com/dotcms/content/index/opensearch/OSIndexAPIImplConnectionClassifyTest.java b/dotCMS/src/test/java/com/dotcms/content/index/opensearch/OSIndexAPIImplConnectionClassifyTest.java index da4386db2c1d..f1bb7689e48f 100644 --- a/dotCMS/src/test/java/com/dotcms/content/index/opensearch/OSIndexAPIImplConnectionClassifyTest.java +++ b/dotCMS/src/test/java/com/dotcms/content/index/opensearch/OSIndexAPIImplConnectionClassifyTest.java @@ -83,6 +83,35 @@ public void securityExceptionStatus403_isAuthForbidden() { + " reason=no permissions for [indices:admin/create]] status: 403")))); } + /** + * A rejected bulk item never carries a status code: the client reports the cluster's + * refusal as plain text on the item. Without matching the security plugin's own wording, a + * permission problem that rejects every shadow write would be classified as + * {@link ConnectionFailureKind#UNKNOWN} and never escalated (issue #36222 follow-up). + */ + @Test + public void bulkItemSecurityException_withoutStatusCode_isAuthForbidden() { + assertEquals(ConnectionFailureKind.AUTH_FORBIDDEN, + OSIndexAPIImpl.classifyFailureMessage( + "security_exception: no permissions for [indices:data/write/bulk[s]," + + " indices:data/write/index] and User [name=non-admin," + + " backend_roles=[], requestedTenant=null]")); + } + + @Test + public void bulkItemMappingFailure_isNotClassified() { + // A per-document problem must stay UNKNOWN: it is not a migration blocker, and reporting it + // as one would escalate every malformed contentlet. + assertEquals(ConnectionFailureKind.UNKNOWN, + OSIndexAPIImpl.classifyFailureMessage( + "mapper_parsing_exception: failed to parse field [myNumber] of type [long]")); + } + + @Test + public void nullFailureMessage_isNotClassified() { + assertEquals(ConnectionFailureKind.UNKNOWN, OSIndexAPIImpl.classifyFailureMessage(null)); + } + /** * A dotCMS wrapper message embeds the physical index name, and a {@code _yyyyMMddHHmmss} * timestamp regularly contains the digits 403 or 401 (here 12:04:03). Matching the status code diff --git a/dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java b/dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java new file mode 100644 index 000000000000..0cc44a759eb6 --- /dev/null +++ b/dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java @@ -0,0 +1,185 @@ +package com.dotmarketing.common.reindex; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +import com.dotcms.content.index.IndexTag; +import com.dotcms.content.index.domain.IndexBulkItemResult; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import org.junit.Test; + +/** + * Unit tests for how {@link BulkProcessorListener} reports failures of the OpenSearch shadow + * bulk leg in dual-write phases (issue #36222 follow-up). + * + *

What this covers

+ *

The listener used to log one WARN per failed item and nothing else. When the OS user loses a + * permission, every document in every batch is rejected with the same message, which produced + * hundreds of thousands of identical lines and — worse — no signal at all that the shadow store had + * stopped receiving writes altogether. QA hit exactly that on TC-056: ~900 identical warnings in one + * minute, and a reindex that looked successful while OpenSearch received nothing.

+ * + *

Both helpers under test are pure functions, so no container, cluster or config is needed. The + * classification of the rejection text itself lives in the OpenSearch adapter and is exercised by + * {@code OSIndexAPIImplConnectionClassifyTest}; here we only assert the escalation policy built on + * top of it.

+ * + *
+ *   ./mvnw test -pl :dotcms-core -Dmaven.build.cache.enabled=false \
+ *       -Dtest=BulkProcessorListenerShadowFailureTest
+ * 
+ * + * @author Fabrizzio Araya + */ +public class BulkProcessorListenerShadowFailureTest { + + /** Verbatim rejection emitted by the OpenSearch security plugin for a denied bulk write. */ + private static final String SECURITY_EXCEPTION = + "security_exception: no permissions for [indices:data/write/bulk[s]," + + " indices:data/write/index] and User [name=non-admin, backend_roles=[]," + + " requestedTenant=null]"; + + /** A per-document problem: the batch is healthy, this one contentlet is not. */ + private static final String MAPPER_PARSING_EXCEPTION = + "mapper_parsing_exception: failed to parse field [myNumber] of type [long]"; + + private static IndexBulkItemResult failed(final String id, final String message) { + return IndexBulkItemResult.builder().id(id).failed(true).failureMessage(message).build(); + } + + private static IndexBulkItemResult succeeded(final String id) { + return IndexBulkItemResult.builder().id(id).failed(false).build(); + } + + private static List allFailedWith(final int count, final String message) { + return IntStream.range(0, count) + .mapToObj(i -> failed("id-" + i, message)) + .collect(Collectors.toList()); + } + + // ---- summarizeFailures ---------------------------------------------------------------------- + + @Test + public void identicalFailures_areCollapsedIntoOneEntryWithItsCount() { + // The whole point: one log line for 185 rejections, not 185 lines. + final Map summary = + BulkProcessorListener.summarizeFailures(allFailedWith(185, SECURITY_EXCEPTION)); + + assertEquals("Identical rejections must collapse to a single entry", 1, summary.size()); + assertEquals(Long.valueOf(185L), summary.get(SECURITY_EXCEPTION)); + } + + @Test + public void successfulItems_areNotCounted() { + final List results = List.of( + succeeded("ok-1"), + failed("bad-1", MAPPER_PARSING_EXCEPTION), + succeeded("ok-2")); + + final Map summary = BulkProcessorListener.summarizeFailures(results); + + assertEquals(1, summary.size()); + assertEquals(Long.valueOf(1L), summary.get(MAPPER_PARSING_EXCEPTION)); + } + + @Test + public void distinctFailures_areKeptApart() { + final List results = List.of( + failed("bad-1", SECURITY_EXCEPTION), + failed("bad-2", MAPPER_PARSING_EXCEPTION), + failed("bad-3", SECURITY_EXCEPTION)); + + final Map summary = BulkProcessorListener.summarizeFailures(results); + + assertEquals(2, summary.size()); + assertEquals(Long.valueOf(2L), summary.get(SECURITY_EXCEPTION)); + assertEquals(Long.valueOf(1L), summary.get(MAPPER_PARSING_EXCEPTION)); + } + + @Test + public void missingFailureMessage_doesNotProduceANullKey() { + final Map summary = + BulkProcessorListener.summarizeFailures(List.of(failed("bad-1", null))); + + assertEquals(Long.valueOf(1L), + summary.get(BulkProcessorListener.NO_FAILURE_MESSAGE)); + } + + @Test + public void healthyBatch_producesNothing() { + assertTrue(BulkProcessorListener.summarizeFailures( + List.of(succeeded("ok-1"), succeeded("ok-2"))).isEmpty()); + } + + // ---- systemicFailureEscalation ------------------------------------------------------------- + + @Test + public void wholeBatchRejectedByPermissions_escalates() { + final List results = allFailedWith(185, SECURITY_EXCEPTION); + + final Optional escalation = BulkProcessorListener.systemicFailureEscalation( + IndexTag.OS, results.size(), BulkProcessorListener.summarizeFailures(results)); + + assertTrue("A batch rejected in full for a permission problem must escalate", + escalation.isPresent()); + final String message = escalation.get(); + assertTrue("The cause must be named so the operator knows what to fix: " + message, + message.contains("AUTH_FORBIDDEN")); + assertTrue("The batch size must be reported: " + message, message.contains("185")); + assertTrue("The consequence — do not promote a diverging store — must be stated: " + message, + message.contains("must not be promoted")); + assertTrue("The verbatim rejection must be kept for support: " + message, + message.contains(SECURITY_EXCEPTION)); + } + + @Test + public void partiallyRejectedBatch_doesNotEscalate() { + // One rejected document out of many is a content problem, not a migration blocker — even + // when its message would classify as systemic on its own. + final List results = List.of( + succeeded("ok-1"), + succeeded("ok-2"), + failed("bad-1", SECURITY_EXCEPTION)); + + assertFalse(BulkProcessorListener.systemicFailureEscalation(IndexTag.OS, results.size(), + BulkProcessorListener.summarizeFailures(results)).isPresent()); + } + + @Test + public void wholeBatchRejectedByAMappingProblem_doesNotEscalate() { + // Every item failing is not enough: a batch of documents that share a broken field is a + // content problem, and halting the operator on it would be noise. + final List results = allFailedWith(10, MAPPER_PARSING_EXCEPTION); + + assertFalse("An unclassifiable cause must not be reported as a systemic failure", + BulkProcessorListener.systemicFailureEscalation(IndexTag.OS, results.size(), + BulkProcessorListener.summarizeFailures(results)).isPresent()); + } + + @Test + public void mixedCauses_escalateOnTheDominantOne() { + // Realistic shape of a denied batch: the permission rejection dominates, one document also + // happens to be malformed. The escalation must name the permission problem. + final List results = List.of( + failed("bad-1", SECURITY_EXCEPTION), + failed("bad-2", SECURITY_EXCEPTION), + failed("bad-3", MAPPER_PARSING_EXCEPTION)); + + final Optional escalation = BulkProcessorListener.systemicFailureEscalation( + IndexTag.OS, results.size(), BulkProcessorListener.summarizeFailures(results)); + + assertTrue(escalation.isPresent()); + assertTrue(escalation.get().contains("AUTH_FORBIDDEN")); + } + + @Test + public void emptyBatch_doesNotEscalate() { + assertFalse(BulkProcessorListener.systemicFailureEscalation( + IndexTag.OS, 0, Map.of()).isPresent()); + } +} From e767de2d69a5b7aa93542b352febfe2866f77492 Mon Sep 17 00:00:00 2001 From: fabrizzio-dotCMS Date: Wed, 5 Aug 2026 14:27:27 -0600 Subject: [PATCH 2/3] refactor(opensearch): throttle the shadow escalation with Logger.errorEvery (#36222) Replaces the ad-hoc AtomicLong throttle and its config key (REINDEX_SHADOW_FAILURE_ESCALATION_MINUTES) with the existing rate-limited logging mechanism: Logger.warnEvery already keys an ExpiringMap by message, so this adds the missing ERROR-level sibling instead of reinventing it. The throttle identity is now the dominant rejection message, so a change of cause is reported immediately rather than waiting out the interval. Co-Authored-By: Claude Opus 5 (1M context) --- .../common/reindex/BulkProcessorListener.java | 54 ++++++++----------- .../java/com/dotmarketing/util/Logger.java | 30 ++++++++++- 2 files changed, 50 insertions(+), 34 deletions(-) diff --git a/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java b/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java index 9a743647d590..965470441633 100644 --- a/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java +++ b/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java @@ -22,7 +22,6 @@ import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicLong; /** * {@link IndexBulkListener} that handles the business logic before/after reindexing content. @@ -40,21 +39,15 @@ public class BulkProcessorListener implements IndexBulkListener { static final List RESERVED_IDS = List.of(Host.SYSTEM_HOST); - /** - * Config key: minutes between two systemic-rejection escalations - * ({@link #systemicFailureEscalation(IndexTag, int, Map)}). Default: 10. - */ - static final String SHADOW_ESCALATION_INTERVAL_KEY = - "REINDEX_SHADOW_FAILURE_ESCALATION_MINUTES"; - /** Stand-in used when the vendor reports a failed item without any message. */ static final String NO_FAILURE_MESSAGE = "(no failure message reported)"; /** - * Epoch millis of the last systemic-rejection escalation. Static on purpose: a fresh listener is - * created for every batch, so per-instance state could not rate-limit anything. + * Minimum gap between two systemic-rejection escalations for the same cause. Every batch is + * rejected while the cause lasts, so without a gap the escalation would flood the log at ERROR + * level — the very problem this reporting fixes. */ - private static final AtomicLong LAST_SYSTEMIC_ESCALATION = new AtomicLong(0L); + private static final int ESCALATION_EVERY_MILLIS = (int) TimeUnit.MINUTES.toMillis(10); private volatile long contentletsIndexed; private int lastBatchSize; @@ -210,11 +203,12 @@ private void logShadowBatchFailures(final List results) { + result.id() + ": " + result.failureMessage())); } - final Optional escalation = - systemicFailureEscalation(provider, results.size(), failuresByMessage); - if (escalation.isPresent() && shouldEscalateNow()) { - Logger.error(this.getClass(), escalation.get()); - } + // Throttled by cause, not by time alone: the same rejection is reported at most once per + // ESCALATION_EVERY_MILLIS, while a different rejection is reported immediately. + systemicFailureEscalation(provider, results.size(), failuresByMessage) + .ifPresent(escalation -> Logger.errorEvery(this.getClass(), + provider.name() + "|" + dominantFailure(failuresByMessage).orElse(""), + escalation, ESCALATION_EVERY_MILLIS)); } /** @@ -261,10 +255,9 @@ static Optional systemicFailureEscalation(final IndexTag provider, final if (batchSize <= 0 || failed < batchSize) { return Optional.empty(); } - return failuresByMessage.entrySet().stream() - .max(Map.Entry.comparingByValue()) + return dominantFailure(failuresByMessage) .flatMap(dominant -> IndexConfigHelper - .systemicFailureRemediation(dominant.getKey()) + .systemicFailureRemediation(dominant) .map(remediation -> "[" + provider.name() + "] EVERY document in this bulk" + " batch (" + batchSize + " item(s)) was rejected by " + provider.name() + " — likelyCause=" + remediation @@ -272,25 +265,20 @@ static Optional systemicFailureEscalation(final IndexTag provider, final + " the authoritative store, so it must not be promoted: advancing" + " the migration phase would serve reads from a stale index." + " Fix the cause above, then run a full reindex to resynchronise." - + " Rejection: " + dominant.getKey())); + + " Rejection: " + dominant)); } /** - * Rate-limits the systemic escalation to one entry per - * {@value #SHADOW_ESCALATION_INTERVAL_KEY} minutes. Static because a fresh listener is built - * for every batch, so per-instance state would not throttle anything. + * The failure message that accounts for most items of the batch. Doubles as the throttle + * identity of the escalation, so a change of cause is reported immediately. * - * @return {@code true} when this caller won the right to log the escalation now + * @param failuresByMessage output of {@link #summarizeFailures(List)} + * @return the dominant failure message, or empty when nothing failed */ - private static boolean shouldEscalateNow() { - final long intervalMs = TimeUnit.MINUTES.toMillis( - Config.getIntProperty(SHADOW_ESCALATION_INTERVAL_KEY, 10)); - final long now = System.currentTimeMillis(); - final long previous = LAST_SYSTEMIC_ESCALATION.get(); - if (previous != 0L && now - previous < intervalMs) { - return false; - } - return LAST_SYSTEMIC_ESCALATION.compareAndSet(previous, now); + static Optional dominantFailure(final Map failuresByMessage) { + return failuresByMessage.entrySet().stream() + .max(Map.Entry.comparingByValue()) + .map(Map.Entry::getKey); } static String getMatchingReservedIdIfAny(final String id) { diff --git a/dotCMS/src/main/java/com/dotmarketing/util/Logger.java b/dotCMS/src/main/java/com/dotmarketing/util/Logger.java index d2708664c261..80b8e9e53e98 100644 --- a/dotCMS/src/main/java/com/dotmarketing/util/Logger.java +++ b/dotCMS/src/main/java/com/dotmarketing/util/Logger.java @@ -244,10 +244,38 @@ public static void warnEvery(final Class cl, final String messageKey, final Stri } + /** + * this method will print the message at ERROR level at most once per millis set, keyed by + * {@code messageKey}: a repeating condition (e.g. every bulk batch being rejected for the same + * reason) is reported without flooding the log, while a different key is reported + * immediately. Mirrors {@link #warnEvery(Class, String, String, int)} at ERROR level. + * + * @param cl the class to attribute the log entry to + * @param messageKey identity of the condition; the same key is throttled together + * @param message the message to print + * @param errorEveryMillis minimum millis between two entries for this key + */ + public static void errorEvery(final Class cl, final String messageKey, final String message, + final int errorEveryMillis) { + + if (UtilMethods.isEmpty(messageKey)) { + return; + } + final org.apache.logging.log4j.Logger logger = loadLogger(cl); + + final Long hash = Long.valueOf(Objects.hashCode(messageKey.intern())); + final Long expireWhen = logMap.get().get(hash); + + if (expireWhen == null || expireWhen < System.currentTimeMillis()) { + logMap.get().put(hash, System.currentTimeMillis() + errorEveryMillis, true); + logger.error(message + " (log every " + errorEveryMillis + "ms)"); + } + } + /** * this method will print the message at WARN level every millis set and print the message plus * whole stack trace if at DEGUG level - * + * * @param cl * @param message * @param ex From 91db9ee5d437059ffdb6145c18fb4e8fc49b39c5 Mon Sep 17 00:00:00 2001 From: fabrizzio-dotCMS Date: Wed, 5 Aug 2026 17:01:44 -0600 Subject: [PATCH 3/3] fix(opensearch): contain any failure of the shadow-failure reporting (#36222) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reporting a shadow bulk failure must never become a failure itself: the OpenSearch adapter calls afterBulk inside its own try, so a throw here came back as afterBulk(Throwable) and was logged as "Bulk process failed entirely" — a defect in the reporting masquerading as the very condition the reporting exists to detect. Catches Throwable, not Exception: the escalation path initialises OSIndexAPIImpl, so the plausible failure is a LinkageError, which neither the adapter's catch(Exception) nor CompositeBulkProcessor.close()'s shadow-isolation branch would contain. Co-Authored-By: Claude Opus 5 (1M context) --- .../common/reindex/BulkProcessorListener.java | 34 ++++++++++++++++++- ...ulkProcessorListenerShadowFailureTest.java | 34 +++++++++++++++++++ 2 files changed, 67 insertions(+), 1 deletion(-) diff --git a/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java b/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java index 965470441633..fbedeed0bf46 100644 --- a/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java +++ b/dotCMS/src/main/java/com/dotmarketing/common/reindex/BulkProcessorListener.java @@ -128,7 +128,7 @@ public void afterBulk(final long executionId, final List re // line per failed item buried the log — a real reindex emits hundreds of thousands of // identical entries, hiding every other line including the actionable one below // (observed on issue #36222, TC-056: ~900 identical WARNs in one minute). - logShadowBatchFailures(results); + reportShadowBatchFailures(results); return; } Logger.debug(this.getClass(), "Bulk process completed"); @@ -179,6 +179,38 @@ public void afterBulk(final long executionId, final Throwable failure) { workingRecords.values().forEach(idx -> handleFailure(idx, msg)); } + /** + * Runs {@link #logShadowBatchFailures(List)} without ever letting it fail the callback. + * + *

Reporting a shadow failure must never become a failure itself. {@code flush()} of the + * OpenSearch adapter invokes this callback inside its own {@code try}, so a throw from here + * is re-entered as {@code afterBulk(executionId, throwable)} and logged as "Bulk process + * failed entirely" — a defect in the reporting would masquerade as the very condition + * this reporting exists to detect, on the signal an operator uses to decide whether the + * shadow store may be promoted.

+ * + *

{@link Throwable}, not {@link Exception}, on purpose: the escalation path reaches + * {@code OSIndexAPIImpl} through {@link IndexConfigHelper#systemicFailureRemediation(String)}, + * so a broken static initialiser in the OpenSearch adapter surfaces as an + * {@code ExceptionInInitializerError} / {@code NoClassDefFoundError}. Those are the failures + * worth containing here: nothing below catches them either + * ({@code CompositeBulkProcessor.close()} catches {@code Exception}, so its shadow-isolation + * branch is skipped) and they would reach {@code ReindexThread} as an opaque + * "ReindexThread Exception" naming neither OpenSearch nor the shadow store.

+ * + * @param results every item result of the completed batch, successful ones included + */ + private void reportShadowBatchFailures(final List results) { + try { + logShadowBatchFailures(results); + } catch (final Throwable t) { // NOSONAR - see javadoc: LinkageError is the case to contain + // Deliberately a plain warn: the recovery path must not re-enter the classifier that + // may have just failed, and a shadow-write reporting problem is not a reindex problem. + Logger.warn(this.getClass(), "[" + provider.name() + "] Unable to report shadow bulk" + + " failures — indexing itself is unaffected: " + t, t); + } + } + /** * Logs the failed items of a shadow bulk batch as one aggregated entry per distinct failure * message, keeps the per-item detail at {@code DEBUG}, and escalates once when the whole batch diff --git a/dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java b/dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java index 0cc44a759eb6..9d7165a3d3c9 100644 --- a/dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java +++ b/dotCMS/src/test/java/com/dotmarketing/common/reindex/BulkProcessorListenerShadowFailureTest.java @@ -182,4 +182,38 @@ public void emptyBatch_doesNotEscalate() { assertFalse(BulkProcessorListener.systemicFailureEscalation( IndexTag.OS, 0, Map.of()).isPresent()); } + + // ---- containment of the reporting itself ----------------------------------------------------- + + @Test + public void aFailureWhileReportingIsContained_andNeverReachesTheCaller() { + // Reporting a shadow failure must not become a failure. The OpenSearch adapter calls this + // callback inside its own try, so a throw here comes back as afterBulk(Throwable) and is + // logged as "Bulk process failed entirely" — a reporting defect would then look exactly like + // the shadow store having stopped accepting writes, the signal used to decide whether the + // store may be promoted. + // + // Error, not Exception, on purpose: the escalation path initialises OSIndexAPIImpl, so the + // realistic failure is a LinkageError, which neither the adapter's catch(Exception) nor + // CompositeBulkProcessor.close()'s shadow-isolation branch would contain. + final IndexBulkItemResult poisoned = new IndexBulkItemResult() { + @Override + public String id() { + return "boom"; + } + + @Override + public boolean failed() { + throw new NoClassDefFoundError("simulated broken static initialiser"); + } + + @Override + public String failureMessage() { + return null; + } + }; + + BulkProcessorListener.forShadowProvider(IndexTag.OS).afterBulk(1L, List.of(poisoned)); + // Reaching this line is the assertion: the callback returned instead of propagating. + } }