diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventProcessor.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventProcessor.java index ddc0f73a27..9292d96673 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventProcessor.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventProcessor.java @@ -64,7 +64,7 @@ public class EventProcessor

implements EventHandler, Life private final Cache

cache; private final EventSourceManager

eventSourceManager; private final RateLimiter rateLimiter; - private final ResourceStateManager resourceStateManager = new ResourceStateManager(); + private final ResourceStateManager resourceStateManager; private final Map metricsMetadata; private ExecutorService executor; @@ -107,6 +107,8 @@ private EventProcessor( this.metrics = metrics != null ? metrics : Metrics.NOOP; this.eventSourceManager = eventSourceManager; this.rateLimiter = controllerConfiguration.getRateLimiter(); + this.resourceStateManager = + new ResourceStateManager(controllerConfiguration.triggerReconcilerOnAllEvents()); metricsMetadata = Optional.ofNullable(eventSourceManager.getController()) @@ -194,7 +196,7 @@ private void submitReconciliationExecution(ResourceState state) { state.getRetry(), state.deleteEventPresent(), state.isDeleteFinalStateUnknown()); - state.unMarkEventReceived(triggerOnAllEvents()); + state.unMarkEventReceived(); metrics.reconciliationSubmitted(latest, state.getRetry(), metricsMetadata); log.debug("Executing events for custom resource. Scope: {}", executionScope); executor.execute(new ReconcilerExecutor(resourceID, executionScope)); @@ -249,10 +251,10 @@ private void handleEventMarking(Event event, ResourceState state) { // removed, but also the informers websocket is disconnected and later reconnected. So // meanwhile the resource could be deleted and recreated. In this case we just mark a new // event as below. - state.markEventReceived(triggerOnAllEvents()); + state.markEventReceived(); } } else if (!state.deleteEventPresent() && !state.processedMarkForDeletionPresent()) { - state.markEventReceived(triggerOnAllEvents()); + state.markEventReceived(); } else if (isTriggerOnAllEventAndDeleteEventPresent(state)) { state.markAdditionalEventAfterDeleteEvent(); } else if (log.isDebugEnabled()) { @@ -381,7 +383,7 @@ private void handleRetryOnException( boolean eventPresent = state.eventPresent() || (triggerOnAllEvents() && state.isAdditionalEventPresentAfterDeleteEvent()); - state.markEventReceived(triggerOnAllEvents()); + state.markEventReceived(); retryAwareErrorLogging( state.getRetry(), eventPresent, errorHandledByReconciler, exception, executionScope); metrics.reconciliationFailed( diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceState.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceState.java index dac24e7941..89ae8396fa 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceState.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceState.java @@ -47,6 +47,7 @@ private enum EventingState { } private final ResourceID id; + private final boolean triggerOnAllEvents; private boolean underProcessing; private RetryExecution retry; @@ -55,8 +56,9 @@ private enum EventingState { private HasMetadata lastKnownResource; private boolean isDeleteFinalStateUnknown = false; - public ResourceState(ResourceID id) { + public ResourceState(ResourceID id, boolean triggerOnAllEvents) { this.id = id; + this.triggerOnAllEvents = triggerOnAllEvents; eventing = EventingState.NO_EVENT_PRESENT; } @@ -108,8 +110,8 @@ public boolean processedMarkForDeletionPresent() { return eventing == EventingState.PROCESSED_MARK_FOR_DELETION; } - public void markEventReceived(boolean isAllEventMode) { - if (!isAllEventMode && deleteEventPresent()) { + public void markEventReceived() { + if (!triggerOnAllEvents && deleteEventPresent()) { throw new IllegalStateException("Cannot receive event after a delete event received"); } log.debug("Marking event received for: {}", getId()); @@ -151,7 +153,7 @@ public HasMetadata getLastKnownResource() { return lastKnownResource; } - public void unMarkEventReceived(boolean isAllEventReconcileMode) { + public void unMarkEventReceived() { switch (eventing) { case EVENT_PRESENT: eventing = EventingState.NO_EVENT_PRESENT; @@ -159,12 +161,12 @@ public void unMarkEventReceived(boolean isAllEventReconcileMode) { case PROCESSED_MARK_FOR_DELETION: throw new IllegalStateException("Cannot unmark processed marked for deletion."); case DELETE_EVENT_PRESENT: - if (!isAllEventReconcileMode) { + if (!triggerOnAllEvents) { throw new IllegalStateException("Cannot unmark delete event."); } break; case ADDITIONAL_EVENT_PRESENT_AFTER_DELETE_EVENT: - if (!isAllEventReconcileMode) { + if (!triggerOnAllEvents) { throw new IllegalStateException( "This state should not happen in non all-event-reconciliation mode"); } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManager.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManager.java index 39a94b7735..9b25c7ae0c 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManager.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManager.java @@ -28,6 +28,11 @@ class ResourceStateManager { // will process to avoid under- or over-sizing the state maps and avoid too many resizing that // take time and memory? private final Map states = new ConcurrentHashMap<>(100); + private final boolean triggerOnAllEvents; + + public ResourceStateManager(boolean triggerOnAllEvents) { + this.triggerOnAllEvents = triggerOnAllEvents; + } public Optional getOrCreateOnResourceEvent(Event event) { var resourceId = event.getRelatedCustomResourceID(); @@ -36,7 +41,7 @@ public Optional getOrCreateOnResourceEvent(Event event) { return Optional.of(state); } if (event instanceof ResourceEvent) { - state = new ResourceState(resourceId); + state = new ResourceState(resourceId, triggerOnAllEvents); states.put(resourceId, state); return Optional.of(state); } else { @@ -45,7 +50,7 @@ public Optional getOrCreateOnResourceEvent(Event event) { } public ResourceState getOrCreate(ResourceID resourceID) { - return states.computeIfAbsent(resourceID, ResourceState::new); + return states.computeIfAbsent(resourceID, id -> new ResourceState(id, triggerOnAllEvents)); } public Optional get(ResourceID resourceID) { diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManagerTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManagerTest.java index d480dd06f8..8ac3be8c35 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManagerTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/ResourceStateManagerTest.java @@ -27,7 +27,7 @@ class ResourceStateManagerTest { - private final ResourceStateManager manager = new ResourceStateManager(); + private final ResourceStateManager manager = new ResourceStateManager(false); private final ResourceID sampleResourceID = new ResourceID("test-name"); private final ResourceID sampleResourceID2 = new ResourceID("test-name2"); private ResourceState state; @@ -49,7 +49,7 @@ public void returnsNoEventPresentIfNotMarkedYet() { @Test public void marksEvent() { - state.markEventReceived(false); + state.markEventReceived(); assertThat(state.eventPresent()).isTrue(); assertThat(state.deleteEventPresent()).isFalse(); @@ -65,7 +65,7 @@ public void marksDeleteEvent() { @Test public void afterDeleteEventMarkEventIsNotRelevant() { - state.markEventReceived(false); + state.markEventReceived(); state.markDeleteEventReceived(TestUtils.testCustomResource(), true); @@ -75,7 +75,7 @@ public void afterDeleteEventMarkEventIsNotRelevant() { @Test public void cleansUp() { - state.markEventReceived(false); + state.markEventReceived(); state.markDeleteEventReceived(TestUtils.testCustomResource(), true); manager.remove(sampleResourceID); @@ -91,15 +91,15 @@ public void cannotMarkEventAfterDeleteEventReceived() { IllegalStateException.class, () -> { state.markDeleteEventReceived(TestUtils.testCustomResource(), true); - state.markEventReceived(false); + state.markEventReceived(); }); } @Test public void listsResourceIDSWithEventsPresent() { - state.markEventReceived(false); - state2.markEventReceived(false); - state.unMarkEventReceived(false); + state.markEventReceived(); + state2.markEventReceived(); + state.unMarkEventReceived(); var res = manager.resourcesWithEventPresent();