diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Informable.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Informable.java index 12c6b4fe06..5175efb898 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Informable.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Informable.java @@ -15,7 +15,10 @@ */ package io.javaoperatorsdk.operator.api.config; +import java.util.Optional; + import io.fabric8.kubernetes.api.model.HasMetadata; +import io.fabric8.kubernetes.client.KubernetesClient; import io.javaoperatorsdk.operator.api.config.informer.InformerConfiguration; public interface Informable { @@ -29,4 +32,12 @@ default String getResourceTypeName() { default Class getResourceClass() { return getInformerConfig().getResourceClass(); } + + /** + * Optional, specific kubernetes client, typically to connect to a different cluster than the rest + * of the operator. Note that this is solely for multi cluster support. + */ + default Optional getKubernetesClient() { + return Optional.empty(); + } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/informer/InformerEventSourceConfiguration.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/informer/InformerEventSourceConfiguration.java index b6f7939728..ae2b12fe16 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/informer/InformerEventSourceConfiguration.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/informer/InformerEventSourceConfiguration.java @@ -83,14 +83,6 @@ default String name() { return getInformerConfig().getName(); } - /** - * Optional, specific kubernetes client, typically to connect to a different cluster than the rest - * of the operator. Note that this is solely for multi cluster support. - */ - default Optional getKubernetesClient() { - return Optional.empty(); - } - class DefaultInformerEventSourceConfiguration implements InformerEventSourceConfiguration { private final PrimaryToSecondaryMapper primaryToSecondaryMapper; diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/PrimaryUpdateAndCacheUtils.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/PrimaryUpdateAndCacheUtils.java index f74cd49ee7..1be8e3f09e 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/PrimaryUpdateAndCacheUtils.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/PrimaryUpdateAndCacheUtils.java @@ -25,12 +25,12 @@ import org.slf4j.LoggerFactory; import io.fabric8.kubernetes.api.model.HasMetadata; -import io.fabric8.kubernetes.api.model.ObjectMeta; import io.fabric8.kubernetes.client.KubernetesClient; import io.fabric8.kubernetes.client.KubernetesClientException; import io.fabric8.kubernetes.client.dsl.base.PatchContext; import io.fabric8.kubernetes.client.dsl.base.PatchType; import io.javaoperatorsdk.operator.OperatorException; +import io.javaoperatorsdk.operator.ReconcilerUtilsInternal; import io.javaoperatorsdk.operator.processing.event.ResourceID; import static io.javaoperatorsdk.operator.processing.KubernetesResourceUtils.getUID; @@ -430,10 +430,7 @@ public static

P addFinalizerWithSSA( } try { P resource = (P) originalResource.getClass().getConstructor().newInstance(); - ObjectMeta objectMeta = new ObjectMeta(); - objectMeta.setName(originalResource.getMetadata().getName()); - objectMeta.setNamespace(originalResource.getMetadata().getNamespace()); - resource.setMetadata(objectMeta); + resource.initNameAndNamespaceFrom(originalResource); resource.addFinalizer(finalizerName); return client .resource(resource) @@ -456,43 +453,10 @@ public static

P addFinalizerWithSSA( } public static int compareResourceVersions(HasMetadata h1, HasMetadata h2) { - return compareResourceVersions( - h1.getMetadata().getResourceVersion(), h2.getMetadata().getResourceVersion()); + return ReconcilerUtilsInternal.validateAndCompareResourceVersions(h1, h2); } public static int compareResourceVersions(String v1, String v2) { - int v1Length = validateResourceVersion(v1); - int v2Length = validateResourceVersion(v2); - int comparison = v1Length - v2Length; - if (comparison != 0) { - return comparison; - } - for (int i = 0; i < v2Length; i++) { - int comp = v1.charAt(i) - v2.charAt(i); - if (comp != 0) { - return comp; - } - } - return 0; - } - - private static int validateResourceVersion(String v1) { - int v1Length = v1.length(); - if (v1Length == 0) { - throw new NonComparableResourceVersionException("Resource version is empty"); - } - for (int i = 0; i < v1Length; i++) { - char char1 = v1.charAt(i); - if (char1 == '0') { - if (i == 0) { - throw new NonComparableResourceVersionException( - "Resource version cannot begin with 0: " + v1); - } - } else if (char1 < '0' || char1 > '9') { - throw new NonComparableResourceVersionException( - "Non numeric characters in resource version: " + v1); - } - } - return v1Length; + return ReconcilerUtilsInternal.validateAndCompareResourceVersions(v1, v2); } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/ResourceOperations.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/ResourceOperations.java index 603548e33d..4c8ff8de9c 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/ResourceOperations.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/ResourceOperations.java @@ -556,7 +556,7 @@ public R jsonPatch(R actualResource, UnaryOperator un */ public R jsonPatch( R actualResource, UnaryOperator unaryOperator, Options options) { - R desired = desiredForJsonPatch(actualResource, unaryOperator, options); + R desired = desiredForJsonPatch(actualResource, unaryOperator); return resourcePatch( desired, actualResource, @@ -580,7 +580,7 @@ public R jsonPatch( UnaryOperator unaryOperator, InformerEventSource informerEventSource, Options options) { - R desired = desiredForJsonPatch(actualResource, unaryOperator, options); + R desired = desiredForJsonPatch(actualResource, unaryOperator); return resourcePatch( desired, actualResource, @@ -620,7 +620,7 @@ public R jsonPatchStatus( */ public R jsonPatchStatus( R actualResource, UnaryOperator unaryOperator, Options options) { - R desired = desiredForJsonPatch(actualResource, unaryOperator, options); + R desired = desiredForJsonPatch(actualResource, unaryOperator); return resourcePatch( desired, actualResource, @@ -645,7 +645,7 @@ public R jsonPatchStatus( UnaryOperator unaryOperator, InformerEventSource informerEventSource, Options options) { - R desired = desiredForJsonPatch(actualResource, unaryOperator, options); + R desired = desiredForJsonPatch(actualResource, unaryOperator); return resourcePatch( desired, actualResource, @@ -680,7 +680,7 @@ public P jsonPatchPrimary(P actualResource, UnaryOperator

unaryOperator) { * @return the patched resource as returned by the API server */ public P jsonPatchPrimary(P actualResource, UnaryOperator

unaryOperator, Options options) { - P desired = desiredForJsonPatch(actualResource, unaryOperator, options); + P desired = desiredForJsonPatch(actualResource, unaryOperator); return resourcePatch( desired, actualResource, @@ -717,7 +717,7 @@ public P jsonPatchPrimaryStatus(P actualResource, UnaryOperator

unaryOperator */ public P jsonPatchPrimaryStatus( P actualResource, UnaryOperator

unaryOperator, Options options) { - P desired = desiredForJsonPatch(actualResource, unaryOperator, options); + P desired = desiredForJsonPatch(actualResource, unaryOperator); return resourcePatch( desired, actualResource, @@ -1441,7 +1441,7 @@ public enum Mode { } private T desiredForJsonPatch( - T actualResource, UnaryOperator unaryOperator, Options options) { + T actualResource, UnaryOperator unaryOperator) { var cloned = context .getControllerConfiguration() diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/GenericKubernetesResourceMatcher.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/GenericKubernetesResourceMatcher.java index b5a0728e16..23fb29151f 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/GenericKubernetesResourceMatcher.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/GenericKubernetesResourceMatcher.java @@ -38,6 +38,12 @@ public class GenericKubernetesResourceMatcher SPEC_PREFIX = List.of(SPEC); + private static final List STATUS_PREFIX = List.of(STATUS); + private static final List METADATA_PREFIX = List.of(METADATA); + private static final List LABELS_AND_ANNOTATIONS_PREFIX = + List.of(METADATA_LABELS, METADATA_ANNOTATIONS); + private static final String PATH = "path"; private static final String[] EMPTY_ARRAY = {}; @@ -182,11 +188,11 @@ public static Matcher.Result m boolean matched = true; for (int i = 0; i < wholeDiffJsonPatch.size() && matched; i++) { var node = wholeDiffJsonPatch.get(i); - if (nodeIsChildOf(node, List.of(SPEC))) { + if (nodeIsChildOf(node, SPEC_PREFIX)) { matched = match(valuesEquality, node, ignoreList); - } else if (nodeIsChildOf(node, List.of(METADATA))) { + } else if (nodeIsChildOf(node, METADATA_PREFIX)) { // conditionally consider labels and annotations - if (nodeIsChildOf(node, List.of(METADATA_LABELS, METADATA_ANNOTATIONS))) { + if (nodeIsChildOf(node, LABELS_AND_ANNOTATIONS_PREFIX)) { matched = match(labelsAndAnnotationsEquality, node, Collections.emptyList()); } } else if (!nodeIsChildOf(node, IGNORED_FIELDS)) { @@ -241,7 +247,7 @@ public static Matcher.Result m boolean matched = true; for (int i = 0; i < wholeDiffJsonPatch.size() && matched; i++) { var node = wholeDiffJsonPatch.get(i); - if (nodeIsChildOf(node, List.of(STATUS))) { + if (nodeIsChildOf(node, STATUS_PREFIX)) { matched = match(valuesEquality, node, Collections.emptyList()); } } @@ -261,7 +267,12 @@ private static boolean match(boolean equality, JsonNode diff, final List static boolean nodeIsChildOf(JsonNode n, List prefixes) { var path = getPath(n); - return prefixes.stream().anyMatch(path::startsWith); + for (int i = 0; i < prefixes.size(); i++) { + if (path.startsWith(prefixes.get(i))) { + return true; + } + } + return false; } static String getPath(JsonNode n) { diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/SSABasedGenericKubernetesResourceMatcher.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/SSABasedGenericKubernetesResourceMatcher.java index d3e5b6dbc5..abec13290d 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/SSABasedGenericKubernetesResourceMatcher.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/kubernetes/SSABasedGenericKubernetesResourceMatcher.java @@ -38,6 +38,7 @@ import io.fabric8.kubernetes.api.model.apps.Deployment; import io.fabric8.kubernetes.api.model.apps.ReplicaSet; import io.fabric8.kubernetes.api.model.apps.StatefulSet; +import io.fabric8.kubernetes.api.model.apps.StatefulSetSpec; import io.fabric8.kubernetes.client.utils.KubernetesSerialization; import io.javaoperatorsdk.operator.OperatorException; import io.javaoperatorsdk.operator.api.reconciler.Context; @@ -199,25 +200,7 @@ protected void sanitizeState(R actual, R desired, Map actualMap) && desired instanceof StatefulSet desiredStatefulSet) { var actualSpec = actualStatefulSet.getSpec(); var desiredSpec = desiredStatefulSet.getSpec(); - int claims = desiredSpec.getVolumeClaimTemplates().size(); - if (claims == actualSpec.getVolumeClaimTemplates().size()) { - for (int i = 0; i < claims; i++) { - var claim = desiredSpec.getVolumeClaimTemplates().get(i); - if (claim.getSpec().getVolumeMode() == null) { - Optional.ofNullable( - GenericKubernetesResource.get( - actualMap, "spec", "volumeClaimTemplates", i, "spec")) - .map(Map.class::cast) - .ifPresent(m -> m.remove("volumeMode")); - } - if (claim.getStatus() == null) { - Optional.ofNullable( - GenericKubernetesResource.get(actualMap, "spec", "volumeClaimTemplates", i)) - .map(Map.class::cast) - .ifPresent(m -> m.remove("status")); - } - } - } + sanitizeVolumeClaimTemplates(actualMap, actualSpec, desiredSpec); sanitizePodTemplateSpec(actualMap, actualSpec.getTemplate(), desiredSpec.getTemplate()); } else if (actual instanceof Deployment actualDeployment && desired instanceof Deployment desiredDeployment) { @@ -240,6 +223,29 @@ protected void sanitizeState(R actual, R desired, Map actualMap) } } + private static void sanitizeVolumeClaimTemplates( + Map actualMap, StatefulSetSpec actualSpec, StatefulSetSpec desiredSpec) { + int claims = desiredSpec.getVolumeClaimTemplates().size(); + if (claims != actualSpec.getVolumeClaimTemplates().size()) { + return; + } + for (int i = 0; i < claims; i++) { + var claim = desiredSpec.getVolumeClaimTemplates().get(i); + if (claim.getSpec().getVolumeMode() == null) { + Optional.ofNullable( + GenericKubernetesResource.get(actualMap, "spec", "volumeClaimTemplates", i, "spec")) + .map(Map.class::cast) + .ifPresent(m -> m.remove("volumeMode")); + } + if (claim.getStatus() == null) { + Optional.ofNullable( + GenericKubernetesResource.get(actualMap, "spec", "volumeClaimTemplates", i)) + .map(Map.class::cast) + .ifPresent(m -> m.remove("status")); + } + } + } + @SuppressWarnings("unchecked") static void keepOnlyManagedFields( Map result, diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/workflow/AbstractWorkflowExecutor.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/workflow/AbstractWorkflowExecutor.java index d3907b657a..665d80063b 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/workflow/AbstractWorkflowExecutor.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/workflow/AbstractWorkflowExecutor.java @@ -52,7 +52,7 @@ protected AbstractWorkflowExecutor(DefaultWorkflow

workflow, P primary, Conte this.context = context; this.primaryID = ResourceID.fromResource(primary); executorService = context.getWorkflowExecutorService(); - results = new HashMap<>(workflow.getDependentResourcesByName().size()); + results = new HashMap<>(workflow.size()); } protected abstract Logger logger(); 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/main/java/io/javaoperatorsdk/operator/processing/event/source/ExternalResourceCachingEventSource.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/ExternalResourceCachingEventSource.java index 8a4c476443..109e83b413 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/ExternalResourceCachingEventSource.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/ExternalResourceCachingEventSource.java @@ -242,7 +242,7 @@ public Set getSecondaryResources(ResourceID primaryID) { if (cachedValues == null) { return Collections.emptySet(); } else { - return new HashSet<>(cache.get(primaryID).values()); + return new HashSet<>(cachedValues.values()); } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/EventFilterWindow.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/EventFilterWindow.java index 826551656e..c63261c0b1 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/EventFilterWindow.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/EventFilterWindow.java @@ -239,8 +239,7 @@ public synchronized void addRelatedEvent(ExtendedResourceEvent event) { event.setPartOfReList(true); } - relatedEvents.put( - Long.valueOf(event.getResource().orElseThrow().getMetadata().getResourceVersion()), event); + relatedEvents.put(event.getResourceVersion(), event); } public synchronized void setReListStarted() { diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerEventSource.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerEventSource.java index cb0fdaa8dd..a8b1a96bec 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerEventSource.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerEventSource.java @@ -178,7 +178,9 @@ public synchronized void start() { super.start(); // this makes sure that on first reconciliation all resources are // present on the index - manager().list().forEach(r -> primaryToSecondaryIndex.onAddOrUpdate(r, null)); + if (useSecondaryToPrimaryIndex()) { + manager().list().forEach(r -> primaryToSecondaryIndex.onAddOrUpdate(r, null)); + } } @SuppressWarnings("unchecked") diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerManager.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerManager.java index 3908bbcf09..6caf39ccd9 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerManager.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/InformerManager.java @@ -33,7 +33,6 @@ import io.javaoperatorsdk.operator.api.config.ControllerConfiguration; import io.javaoperatorsdk.operator.api.config.Informable; import io.javaoperatorsdk.operator.api.config.informer.InformerConfiguration; -import io.javaoperatorsdk.operator.api.config.informer.InformerEventSourceConfiguration; import io.javaoperatorsdk.operator.health.InformerHealthIndicator; import io.javaoperatorsdk.operator.processing.event.ResourceID; import io.javaoperatorsdk.operator.processing.event.source.Cache; @@ -179,13 +178,11 @@ private KubernetesClient getTargetClient() { // to see the very same instance. ConfigurationService#getKubernetesClient is expected to return // a stable instance, but its default implementation does create a new client on every call. if (targetClient == null) { - targetClient = controllerConfiguration.getConfigurationService().getKubernetesClient(); - if (configuration instanceof InformerEventSourceConfiguration iesc) { - var remoteClient = iesc.getKubernetesClient().orElse(null); - if (remoteClient != null) { - targetClient = remoteClient; - } - } + targetClient = + configuration + .getKubernetesClient() + .orElseGet( + () -> controllerConfiguration.getConfigurationService().getKubernetesClient()); } return targetClient; } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/Mappers.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/Mappers.java index efc6a981c3..5636fc3893 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/Mappers.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/Mappers.java @@ -124,6 +124,8 @@ private static SecondaryToPrimaryMapper fromMetadata( String typeKey, Class primaryResourceType, boolean isLabel) { + final var expectedGvk = GroupVersionKind.gvkFor(primaryResourceType); + final var expectedGvkString = expectedGvk.toGVKString(); return resource -> { final var metadata = resource.getMetadata(); if (metadata == null) { @@ -143,8 +145,8 @@ private static SecondaryToPrimaryMapper fromMetadata( String gvkSimple = map.get(typeKey); if (gvkSimple != null - && !GroupVersionKind.fromString(gvkSimple) - .equals(GroupVersionKind.gvkFor(primaryResourceType))) { + && !expectedGvkString.equals(gvkSimple) + && !GroupVersionKind.fromString(gvkSimple).equals(expectedGvk)) { return Set.of(); } @@ -183,7 +185,7 @@ SecondaryToPrimaryMapper fromOwnerType(Class clazz) { } return owners.stream() .filter(it -> kind.equals(it.getKind())) - .map(it -> new ResourceID(it.getName(), resource.getMetadata().getNamespace())) + .map(it -> ResourceID.fromOwnerReference(resource, it, false)) .collect(Collectors.toSet()); }; } @@ -191,16 +193,16 @@ SecondaryToPrimaryMapper fromOwnerType(Class clazz) { public static class SecondaryToPrimaryFromDefaultAnnotation implements SecondaryToPrimaryMapper { - private final Class primaryResourceType; + private final SecondaryToPrimaryMapper delegate; public SecondaryToPrimaryFromDefaultAnnotation( Class primaryResourceType) { - this.primaryResourceType = primaryResourceType; + this.delegate = Mappers.fromDefaultAnnotations(primaryResourceType); } @Override public Set toPrimaryResourceIDs(HasMetadata resource) { - return Mappers.fromDefaultAnnotations(primaryResourceType).toPrimaryResourceIDs(resource); + return delegate.toPrimaryResourceIDs(resource); } } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/pool/AbstractInformerPool.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/pool/AbstractInformerPool.java index 81f8f979c1..2c960a7b6f 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/pool/AbstractInformerPool.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/pool/AbstractInformerPool.java @@ -128,15 +128,11 @@ protected SharedIndexInformer createInformer(InformerClassifier classifier) { return null; }); } else { - final var apiTypeClass = informer.getApiTypeClass(); - final var fullResourceName = HasMetadata.getFullResourceName(apiTypeClass); - final var version = HasMetadata.getVersion(apiTypeClass); throw new IllegalStateException( "Cannot retrieve 'stopped' callback to listen to informer stopping for" + " informer for " - + fullResourceName - + "/" - + version); + + ReconcilerUtilsInternal.getResourceTypeNameWithVersion( + informer.getApiTypeClass())); } }); if (!configurationService.stopOnInformerErrorDuringStartup()) { diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java index 0f0eb78a69..69c7267dbf 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java @@ -76,8 +76,9 @@ public PerResourcePollingEventSource( private Set getAndCacheResource(P primary, boolean fromGetter) { var values = resourceFetcher.fetchResources(primary); - handleResources(ResourceID.fromResource(primary), values, !fromGetter); - fetchedForPrimaries.add(ResourceID.fromResource(primary)); + var primaryID = ResourceID.fromResource(primary); + handleResources(primaryID, values, !fromGetter); + fetchedForPrimaries.add(primaryID); return values; } 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(); diff --git a/operator-framework-junit/src/main/java/io/javaoperatorsdk/operator/junit/LocallyRunOperatorExtension.java b/operator-framework-junit/src/main/java/io/javaoperatorsdk/operator/junit/LocallyRunOperatorExtension.java index 2b2c3bec48..5e7602b094 100644 --- a/operator-framework-junit/src/main/java/io/javaoperatorsdk/operator/junit/LocallyRunOperatorExtension.java +++ b/operator-framework-junit/src/main/java/io/javaoperatorsdk/operator/junit/LocallyRunOperatorExtension.java @@ -51,6 +51,7 @@ import io.javaoperatorsdk.operator.RegisteredController; import io.javaoperatorsdk.operator.api.config.ConfigurationServiceOverrider; import io.javaoperatorsdk.operator.api.config.ControllerConfigurationOverrider; +import io.javaoperatorsdk.operator.api.config.Utils; import io.javaoperatorsdk.operator.api.reconciler.Reconciler; import io.javaoperatorsdk.operator.processing.retry.Retry; @@ -558,11 +559,7 @@ public Builder withReconciler(Reconciler value, Retry retry) { @SuppressWarnings("rawtypes") public Builder withReconciler(Class value) { - try { - reconcilers.add(new ReconcilerSpec(value.getConstructor().newInstance(), null)); - } catch (Exception e) { - throw new RuntimeException(e); - } + reconcilers.add(new ReconcilerSpec(Utils.instantiate(value, Reconciler.class, null), null)); return this; }