diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/RegisteredController.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/RegisteredController.java index e6aa6cbce6..ac5b7cd468 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/RegisteredController.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/RegisteredController.java @@ -18,7 +18,6 @@ import io.fabric8.kubernetes.api.model.HasMetadata; import io.javaoperatorsdk.operator.api.config.ControllerConfiguration; import io.javaoperatorsdk.operator.api.config.NamespaceChangeable; -import io.javaoperatorsdk.operator.api.event.EventRecorder; import io.javaoperatorsdk.operator.health.ControllerHealthInfo; public interface RegisteredController

extends NamespaceChangeable { @@ -26,17 +25,4 @@ public interface RegisteredController

extends NamespaceCh ControllerConfiguration

getConfiguration(); ControllerHealthInfo getControllerHealthInfo(); - - /** - * Returns the {@link EventRecorder} of this controller, to record Kubernetes events outside of a - * reconciliation, for example from a status listener or a background task. Within a - * reconciliation, use {@link io.javaoperatorsdk.operator.api.reconciler.Context#eventRecorder()} - * instead. - * - * @return the event recorder associated with this controller - */ - default EventRecorder eventRecorder() { - throw new UnsupportedOperationException( - "This implementation of RegisteredController does not provide an EventRecorder"); - } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java index 0b1d6b47cb..db1b9a5fa5 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java @@ -35,6 +35,7 @@ import io.fabric8.kubernetes.client.KubernetesClientBuilder; import io.fabric8.kubernetes.client.utils.KubernetesSerialization; import io.javaoperatorsdk.operator.api.event.DefaultEventRecorder; +import io.javaoperatorsdk.operator.api.event.EventRecorder; import io.javaoperatorsdk.operator.api.monitoring.Metrics; import io.javaoperatorsdk.operator.api.reconciler.Context; import io.javaoperatorsdk.operator.api.reconciler.Experimental; @@ -295,6 +296,24 @@ default String clusterScopedEventNamespace() { return DefaultEventRecorder.CLUSTER_SCOPED_EVENT_NAMESPACE; } + /** + * The {@link EventRecorder} the controllers of the operator record their Kubernetes events + * through, to plug in a custom implementation, for example one that assembles events differently + * by extending {@link DefaultEventRecorder}, or one that records them somewhere else entirely. + * + *

When empty, which is the default, every controller gets a {@link DefaultEventRecorder} of + * its own. A recorder configured here is shared by all controllers of the operator instead, which + * is why the reconciliation an event is recorded from is passed to it per call rather than + * configured on it: implementations have to be stateless and thread safe. + * + * @return the event recorder to use for the whole operator, or an empty optional to let each + * controller use its own default one + */ + @Experimental(Experimental.API_MIGHT_CHANGE) + default Optional eventRecorder() { + return Optional.empty(); + } + /** * if true, operator stops if there are some issues with informers {@link * io.javaoperatorsdk.operator.processing.event.source.informer.InformerEventSource} or {@link diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java index 9f0fd78356..caf40555fb 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java @@ -27,6 +27,7 @@ import io.fabric8.kubernetes.api.model.HasMetadata; import io.fabric8.kubernetes.client.KubernetesClient; import io.javaoperatorsdk.operator.Operator; +import io.javaoperatorsdk.operator.api.event.EventRecorder; import io.javaoperatorsdk.operator.api.monitoring.Metrics; import io.javaoperatorsdk.operator.api.reconciler.Experimental; import io.javaoperatorsdk.operator.api.reconciler.dependent.DependentResourceFactory; @@ -48,6 +49,7 @@ public class ConfigurationServiceOverrider { private ExecutorService workflowExecutorService; private LeaderElectionConfiguration leaderElectionConfiguration; private String clusterScopedEventNamespace; + private EventRecorder eventRecorder; private InformerStoppedHandler informerStoppedHandler; private Boolean stopOnInformerErrorDuringStartup; private Duration cacheSyncTimeout; @@ -148,6 +150,25 @@ public ConfigurationServiceOverrider withClusterScopedEventNamespace(String name return this; } + /** + * Replaces the {@link EventRecorder} the controllers of the operator record their Kubernetes + * events through by the specified one, which is then shared by all of them. Use this to record + * events differently, for example through a subclass of {@link + * io.javaoperatorsdk.operator.api.event.DefaultEventRecorder} that assembles them another way, or + * by delegating to another system (e.g. emitting events to an external store). + * + *

When not set, every controller records its events through a recorder of its own, which + * attributes them to that controller. + * + * @param eventRecorder the event recorder to use for the whole operator + * @return this {@link ConfigurationServiceOverrider} for chained customization + */ + @Experimental(Experimental.API_MIGHT_CHANGE) + public ConfigurationServiceOverrider withEventRecorder(EventRecorder eventRecorder) { + this.eventRecorder = eventRecorder; + return this; + } + public ConfigurationServiceOverrider withInformerStoppedHandler(InformerStoppedHandler handler) { this.informerStoppedHandler = handler; return this; @@ -297,6 +318,11 @@ public String clusterScopedEventNamespace() { : original.clusterScopedEventNamespace(); } + @Override + public Optional eventRecorder() { + return eventRecorder != null ? Optional.of(eventRecorder) : original.eventRecorder(); + } + @Override public Optional getInformerStoppedHandler() { return informerStoppedHandler != null diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorder.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorder.java index abf4ab025f..2e7023623d 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorder.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorder.java @@ -33,6 +33,8 @@ import io.fabric8.kubernetes.api.model.HasMetadata; import io.fabric8.kubernetes.api.model.ObjectReference; import io.fabric8.kubernetes.api.model.ObjectReferenceBuilder; +import io.javaoperatorsdk.operator.api.config.LeaderElectionConfiguration; +import io.javaoperatorsdk.operator.api.reconciler.Context; import static java.util.Objects.requireNonNullElse; @@ -71,50 +73,46 @@ public class DefaultEventRecorder implements EventRecorder { private static final int IDENTITY_HASH_LENGTH = 32; - private final String reportingController; - private final String reportingInstance; - private final String clusterScopedEventNamespace; private final EventSink sink; - public DefaultEventRecorder( - String reportingController, String reportingInstance, EventSink sink) { - this(reportingController, reportingInstance, CLUSTER_SCOPED_EVENT_NAMESPACE, sink); - } - - public DefaultEventRecorder( - String reportingController, - String reportingInstance, - String clusterScopedEventNamespace, - EventSink sink) { - this.reportingController = reportingController; - this.reportingInstance = reportingInstance; - this.clusterScopedEventNamespace = clusterScopedEventNamespace; + public DefaultEventRecorder(EventSink sink) { this.sink = sink; } /** * The instance name to report events under, when it is not otherwise configured. Uses the host * name, which for an operator running in a pod is the pod name. + * + *

Resolved once and cached: it cannot change over the life of the process, and looking the + * host name up can hit the name service, which is not something to do on every recorded event. */ public static String defaultReportingInstance() { - var fromEnv = System.getenv("HOSTNAME"); - if (fromEnv != null && !fromEnv.isBlank()) { - return fromEnv; - } - try { - return InetAddress.getLocalHost().getHostName(); - } catch (UnknownHostException e) { - log.debug("Could not determine host name to report events under", e); - return "unknown"; + return DefaultReportingInstance.VALUE; + } + + private static final class DefaultReportingInstance { + private static final String VALUE = resolve(); + + private static String resolve() { + var fromEnv = System.getenv("HOSTNAME"); + if (fromEnv != null && !fromEnv.isBlank()) { + return fromEnv; + } + try { + return InetAddress.getLocalHost().getHostName(); + } catch (UnknownHostException e) { + log.debug("Could not determine host name to report events under", e); + return "unknown"; + } } } @Override - public void record(HasMetadata regarding, EventRecord event) { - Objects.requireNonNull(regarding, "the object the event is about must not be null"); + public void record(EventRecord event, Context context) { + Objects.requireNonNull(context, "the context of the reconciliation must not be null"); Objects.requireNonNull(event, "event must not be null"); try { - sink.emit(toEvent(regarding, event)); + sink.emit(toEvent(context, event), context); } catch (Exception e) { // recording an event must never break the caller: a controller that fails to reconcile // because it could not write an event is strictly worse than one that records nothing @@ -122,26 +120,28 @@ public void record(HasMetadata regarding, EventRecord event) { "Could not record {} event with reason {} for resource {} in namespace {}", event.type(), event.reason(), - regarding.getMetadata().getName(), - regarding.getMetadata().getNamespace(), + context.getPrimaryResource().getMetadata().getName(), + context.getPrimaryResource().getMetadata().getNamespace(), e); } } @Override - public ResourceEventRecorder forResource(HasMetadata regarding) { - Objects.requireNonNull(regarding, "the object events will be about must not be null"); - return new BoundEventRecorder(this, regarding); + public ResourceEventRecorder forContext(Context context) { + Objects.requireNonNull(context, "the context events will be recorded from must not be null"); + return new BoundEventRecorder(this, context); } - protected Event toEvent(HasMetadata regarding, EventRecord record) { + protected Event toEvent(Context context, EventRecord record) { + var controllerName = context.getControllerConfiguration().getName(); + var regarding = context.getPrimaryResource(); var now = Instant.now().truncatedTo(ChronoUnit.SECONDS).toString(); var involvedObject = objectReferenceFor(regarding); var builder = new EventBuilder() .withNewMetadata() - .withName(eventName(regarding, record)) - .withNamespace(eventNamespace(regarding)) + .withName(eventName(regarding, record, controllerName)) + .withNamespace(eventNamespace(regarding, context)) .withLabels(record.labels()) .withAnnotations(record.annotations()) .endMetadata() @@ -152,19 +152,33 @@ protected Event toEvent(HasMetadata regarding, EventRecord record) { .withFirstTimestamp(now) .withLastTimestamp(now) .withCount(1) - .withReportingComponent(record.reportingComponent().orElse(reportingController)) - .withReportingInstance(reportingInstance) + .withReportingComponent(record.reportingComponent().orElse(controllerName)) + .withReportingInstance( + context + .getControllerConfiguration() + .getConfigurationService() + .getLeaderElectionConfiguration() + .flatMap(LeaderElectionConfiguration::getIdentity) + .orElseGet(DefaultEventRecorder::defaultReportingInstance)) // the deprecated source is still what kubectl renders in the "From" column .withNewSource() - .withComponent(record.reportingComponent().orElse(reportingController)) + .withComponent(record.reportingComponent().orElse(controllerName)) .endSource(); record.action().ifPresent(builder::withAction); return builder.build(); } - private String eventNamespace(HasMetadata regarding) { + private String eventNamespace(HasMetadata regarding, Context context) { var namespace = regarding.getMetadata().getNamespace(); - return namespace == null ? clusterScopedEventNamespace : namespace; + if (namespace != null) { + return namespace; + } + return requireNonNullElse( + context + .getControllerConfiguration() + .getConfigurationService() + .clusterScopedEventNamespace(), + CLUSTER_SCOPED_EVENT_NAMESPACE); } /** @@ -177,7 +191,7 @@ private String eventNamespace(HasMetadata regarding) { *

The object is identified by its uid, with the kind as a fallback for objects that do not * have one yet, such as a dependent resource that has only been built so far. */ - private String eventName(HasMetadata regarding, EventRecord record) { + private String eventName(HasMetadata regarding, EventRecord record, String reportingController) { var metadata = regarding.getMetadata(); var identity = String.join( @@ -234,22 +248,22 @@ private ObjectReference objectReferenceFor(HasMetadata resource) { .build(); } - private record BoundEventRecorder(EventRecorder delegate, HasMetadata regarding) + private record BoundEventRecorder(EventRecorder delegate, Context context) implements ResourceEventRecorder { @Override public void normal(String reason, String message) { - record(EventRecord.normal(reason, message)); + delegate.record(EventRecord.normal(reason, message), context); } @Override public void warn(String reason, String message) { - record(EventRecord.warning(reason, message)); + delegate.record(EventRecord.warning(reason, message), context); } @Override public void record(EventRecord event) { - delegate.record(regarding, event); + delegate.record(event, context); } } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventSink.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventSink.java index 827002a8e1..c2f5461ae8 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventSink.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/DefaultEventSink.java @@ -18,6 +18,7 @@ import io.fabric8.kubernetes.api.model.Event; import io.fabric8.kubernetes.api.model.EventBuilder; import io.fabric8.kubernetes.client.KubernetesClient; +import io.javaoperatorsdk.operator.api.reconciler.Context; import static java.util.Objects.requireNonNullElse; @@ -47,7 +48,7 @@ public DefaultEventSink(KubernetesClient client) { } @Override - public void emit(Event event) { + public void emit(Event event, Context context) { var events = client.v1().events().inNamespace(event.getMetadata().getNamespace()); var name = event.getMetadata().getName(); var existing = events.withName(name).get(); diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventRecorder.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventRecorder.java index a7d66d8e54..9c4984f187 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventRecorder.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventRecorder.java @@ -15,38 +15,46 @@ */ package io.javaoperatorsdk.operator.api.event; -import io.fabric8.kubernetes.api.model.HasMetadata; +import io.javaoperatorsdk.operator.api.reconciler.Context; +import io.javaoperatorsdk.operator.api.reconciler.Experimental; + +import static io.javaoperatorsdk.operator.api.reconciler.Experimental.API_MIGHT_CHANGE; /** * Records Kubernetes events on behalf of a controller. * - *

This is the unbound form of the API: it is scoped to a controller, not to a reconciliation, - * and can therefore be used outside of the reconciliation loop, for example from a status listener - * or a background task. Obtain it from {@link - * io.javaoperatorsdk.operator.RegisteredController#eventRecorder()}. Within a reconciliation, - * prefer {@link io.javaoperatorsdk.operator.api.reconciler.Context#eventRecorder()}, which is - * already bound to the primary resource. + *

This is the unbound form of the API: an instance is shared by all the controllers of the + * operator, and everything that varies between them - the primary resource an event is about, the + * controller the event is attributed to, and the configuration the event is assembled from - is + * passed per call, as the {@link Context} of the reconciliation recording the event. + * Implementations are therefore expected to be stateless and thread safe. To record events through + * an implementation of your own, see {@link + * io.javaoperatorsdk.operator.api.config.ConfigurationService#eventRecorder()}. Within a + * reconciliation, prefer {@link + * io.javaoperatorsdk.operator.api.reconciler.Context#eventRecorder()}, which is already bound to + * the context. * *

Recording an event is best effort: failures to write the event to the cluster are logged and * swallowed, and never fail the caller. */ +@Experimental(API_MIGHT_CHANGE) public interface EventRecorder { /** - * Records an event about the given object. + * Records an event about the primary resource of the given reconciliation. * - * @param regarding the object the event is about; it will be referenced as the involved object of - * the resulting event * @param event the event to record + * @param context the context of the reconciliation recording the event; the event is about its + * primary resource and is attributed to its controller */ - void record(HasMetadata regarding, EventRecord event); + void record(EventRecord event, Context context); /** - * Returns a view of this recorder bound to the given object, so that the object doesn't have to - * be passed for every event. + * Returns a view of this recorder bound to the given reconciliation, so that the context doesn't + * have to be passed for every event. * - * @param regarding the object subsequent events will be about - * @return a recorder bound to {@code regarding} + * @param context the context subsequent events will be recorded from + * @return a recorder bound to {@code context} */ - ResourceEventRecorder forResource(HasMetadata regarding); + ResourceEventRecorder forContext(Context context); } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventSink.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventSink.java index 54763902a8..93c58dcf54 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventSink.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventSink.java @@ -16,6 +16,10 @@ package io.javaoperatorsdk.operator.api.event; import io.fabric8.kubernetes.api.model.Event; +import io.javaoperatorsdk.operator.api.reconciler.Context; +import io.javaoperatorsdk.operator.api.reconciler.Experimental; + +import static io.javaoperatorsdk.operator.api.reconciler.Experimental.API_MIGHT_CHANGE; /** * Writes fully built events somewhere. Extracted from {@link EventRecorder} so that the assembly of @@ -29,7 +33,12 @@ public interface EventSink { /** * Delivers the event. * - * @param event the event to deliver + * @param event the event to deliver, fully assembled: everything the event says is already built + * into it + * @param context the context of the reconciliation the event was recorded from, for + * implementations that route the event based on it rather than on its contents. Ignored by + * {@link DefaultEventSink}. */ - void emit(Event event); + @Experimental(API_MIGHT_CHANGE) + void emit(Event event, Context context); } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/ResourceEventRecorder.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/ResourceEventRecorder.java index a1fcd6d272..0350688af3 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/ResourceEventRecorder.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/ResourceEventRecorder.java @@ -15,6 +15,10 @@ */ package io.javaoperatorsdk.operator.api.event; +import io.javaoperatorsdk.operator.api.reconciler.Experimental; + +import static io.javaoperatorsdk.operator.api.reconciler.Experimental.API_MIGHT_CHANGE; + /** * An {@link EventRecorder} bound to a single object, typically the primary resource of the current * reconciliation. @@ -22,6 +26,7 @@ *

Recording an event is best effort: failures to write the event to the cluster are logged and * swallowed, and never fail the caller. */ +@Experimental(API_MIGHT_CHANGE) public interface ResourceEventRecorder { /** Records a {@link EventType#NORMAL} event about the bound object. */ diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/Context.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/Context.java index df9c19b263..9743632404 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/Context.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/Context.java @@ -215,9 +215,9 @@ default Optional getSecondaryResource( ResourceOperations

resourceOperations(); /** - * Returns a {@link ResourceEventRecorder} bound to the primary resource, to record Kubernetes - * events about it. To record events outside of a reconciliation, or about another object, use - * {@link io.javaoperatorsdk.operator.RegisteredController#eventRecorder()}. + * Returns a {@link ResourceEventRecorder} bound to this context, to record Kubernetes events + * about the primary resource. To record them through an implementation of your own, see {@link + * io.javaoperatorsdk.operator.api.config.ConfigurationService#eventRecorder()}. * * @return an event recorder bound to the primary resource * @throws UnsupportedOperationException if the implementation does not provide an event recorder diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/DefaultContext.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/DefaultContext.java index 1c90c7535f..e399f7fdfd 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/DefaultContext.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/reconciler/DefaultContext.java @@ -211,7 +211,7 @@ public ResourceOperations

resourceOperations() { @Override public ResourceEventRecorder eventRecorder() { - return controller.eventRecorder().forResource(primaryResource); + return controller.eventRecorder().forContext(this); } @Override diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/Controller.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/Controller.java index 285eb3988c..89ba7e9d71 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/Controller.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/Controller.java @@ -38,7 +38,6 @@ import io.javaoperatorsdk.operator.RegisteredController; import io.javaoperatorsdk.operator.api.config.ControllerConfiguration; import io.javaoperatorsdk.operator.api.config.ExecutorServiceManager; -import io.javaoperatorsdk.operator.api.config.LeaderElectionConfiguration; import io.javaoperatorsdk.operator.api.config.workflow.WorkflowSpec; import io.javaoperatorsdk.operator.api.event.DefaultEventRecorder; import io.javaoperatorsdk.operator.api.event.DefaultEventSink; @@ -115,15 +114,9 @@ public Controller( this.kubernetesClient = kubernetesClient; this.metrics = Optional.ofNullable(configurationService.getMetrics()).orElse(Metrics.NOOP); this.eventRecorder = - new DefaultEventRecorder( - configuration.getName(), - configurationService - .getLeaderElectionConfiguration() - .flatMap(LeaderElectionConfiguration::getIdentity) - .orElseGet(DefaultEventRecorder::defaultReportingInstance), - Optional.ofNullable(configurationService.clusterScopedEventNamespace()) - .orElse(DefaultEventRecorder.CLUSTER_SCOPED_EVENT_NAMESPACE), - new DefaultEventSink(kubernetesClient)); + configurationService + .eventRecorder() + .orElseGet(() -> new DefaultEventRecorder(new DefaultEventSink(kubernetesClient))); contextInitializer = reconciler instanceof ContextInitializer; isCleaner = reconciler instanceof Cleaner; @@ -358,7 +351,11 @@ public ControllerHealthInfo getControllerHealthInfo() { return controllerHealthInfo; } - @Override + /** + * The {@link EventRecorder} this controller records its Kubernetes events through, either the one + * configured for the operator, see {@link ConfigurationService#eventRecorder()}, or a {@link + * DefaultEventRecorder} of its own. + */ public EventRecorder eventRecorder() { return eventRecorder; } diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java index 9df62bc03c..ab4d32e756 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java @@ -23,7 +23,12 @@ import org.junit.jupiter.api.Test; import io.fabric8.kubernetes.api.model.HasMetadata; +import io.javaoperatorsdk.operator.api.event.DefaultEventRecorder; +import io.javaoperatorsdk.operator.api.event.EventRecord; +import io.javaoperatorsdk.operator.api.event.EventRecorder; +import io.javaoperatorsdk.operator.api.event.ResourceEventRecorder; import io.javaoperatorsdk.operator.api.monitoring.Metrics; +import io.javaoperatorsdk.operator.api.reconciler.Context; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.api.Assertions.assertNotEquals; @@ -106,6 +111,28 @@ public R clone(R object) { config.reconciliationTerminationTimeout(), overridden.reconciliationTerminationTimeout()); } + @Test + void eventRecorderIsNotConfiguredByDefaultAndCanBeOverridden() { + final var eventRecorder = + new EventRecorder() { + @Override + public void record(EventRecord event, Context context) {} + + @Override + public ResourceEventRecorder forContext(Context context) { + return null; + } + }; + + assertThat(config.eventRecorder()).isEmpty(); + assertThat( + new ConfigurationServiceOverrider(config) + .withEventRecorder(eventRecorder) + .build() + .eventRecorder()) + .contains(eventRecorder); + } + @Test void threadCountConfiguredProperly() { final var overridden = @@ -118,4 +145,16 @@ void threadCountConfiguredProperly() { assertThat(((ThreadPoolExecutor) overridden.getWorkflowExecutorService()).getMaximumPoolSize()) .isEqualTo(14); } + + @Test + void clusterScopedEventNamespaceDefaultsToTheDefaultNamespaceAndCanBeOverridden() { + assertThat(config.clusterScopedEventNamespace()) + .isEqualTo(DefaultEventRecorder.CLUSTER_SCOPED_EVENT_NAMESPACE); + assertThat( + new ConfigurationServiceOverrider(config) + .withClusterScopedEventNamespace("operator-ns") + .build() + .clusterScopedEventNamespace()) + .isEqualTo("operator-ns"); + } } diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorderTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorderTest.java index 77ad3408d7..eb713e5d74 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorderTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/event/DefaultEventRecorderTest.java @@ -17,18 +17,26 @@ import java.util.ArrayList; import java.util.List; +import java.util.Optional; import org.junit.jupiter.api.Test; import io.fabric8.kubernetes.api.model.ConfigMap; import io.fabric8.kubernetes.api.model.ConfigMapBuilder; import io.fabric8.kubernetes.api.model.Event; +import io.fabric8.kubernetes.api.model.HasMetadata; import io.fabric8.kubernetes.api.model.Namespace; import io.fabric8.kubernetes.api.model.NamespaceBuilder; +import io.javaoperatorsdk.operator.api.config.ConfigurationService; +import io.javaoperatorsdk.operator.api.config.ControllerConfiguration; +import io.javaoperatorsdk.operator.api.reconciler.Context; +import static io.javaoperatorsdk.operator.api.config.LeaderElectionConfigurationBuilder.aLeaderElectionConfiguration; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatCode; import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; class DefaultEventRecorderTest { @@ -37,11 +45,11 @@ class DefaultEventRecorderTest { private final List emitted = new ArrayList<>(); private final DefaultEventRecorder recorder = - new DefaultEventRecorder(CONTROLLER, INSTANCE, emitted::add); + new DefaultEventRecorder((event, context) -> emitted.add(event)); @Test void fillsInEverythingDerivableFromTheControllerAndTheInvolvedObject() { - recorder.record(configMap(), EventRecord.warning("Failed", "could not do the thing")); + recorder.record(EventRecord.warning("Failed", "could not do the thing"), context(configMap())); assertThat(emitted).hasSize(1); var event = emitted.get(0); @@ -63,9 +71,35 @@ void fillsInEverythingDerivableFromTheControllerAndTheInvolvedObject() { assertThat(involved.getResourceVersion()).isEqualTo("42"); } + @Test + void takesTheReportingControllerFromTheControllerConfigurationOfTheContext() { + var context = context(configMap()); + when(context.getControllerConfiguration().getName()).thenReturn("othercontroller"); + + recorder.record(EventRecord.normal("Created", "created"), context); + + assertThat(emitted.get(0).getReportingComponent()).isEqualTo("othercontroller"); + assertThat(emitted.get(0).getSource().getComponent()).isEqualTo("othercontroller"); + } + + @Test + void fallsBackToTheHostNameWhenLeaderElectionConfiguresNoIdentity() { + var context = context(configMap()); + when(context + .getControllerConfiguration() + .getConfigurationService() + .getLeaderElectionConfiguration()) + .thenReturn(Optional.empty()); + + recorder.record(EventRecord.normal("Created", "created"), context); + + assertThat(emitted.get(0).getReportingInstance()) + .isEqualTo(DefaultEventRecorder.defaultReportingInstance()); + } + @Test void createsTheEventInTheNamespaceOfTheInvolvedObject() { - recorder.record(configMap(), EventRecord.normal("Created", "created")); + recorder.record(EventRecord.normal("Created", "created"), context(configMap())); assertThat(emitted.get(0).getMetadata().getNamespace()).isEqualTo("ns1"); assertThat(emitted.get(0).getMetadata().getName()).startsWith("test1."); @@ -73,7 +107,7 @@ void createsTheEventInTheNamespaceOfTheInvolvedObject() { @Test void recordsEventsForClusterScopedObjectsInTheDefaultNamespace() { - recorder.record(clusterScoped(), EventRecord.normal("Created", "created")); + recorder.record(EventRecord.normal("Created", "created"), context(clusterScoped())); assertThat(emitted.get(0).getMetadata().getNamespace()) .isEqualTo(DefaultEventRecorder.CLUSTER_SCOPED_EVENT_NAMESPACE); @@ -82,18 +116,15 @@ void recordsEventsForClusterScopedObjectsInTheDefaultNamespace() { @Test void clusterScopedEventNamespaceCanBeOverridden() { - var configured = new DefaultEventRecorder(CONTROLLER, INSTANCE, "operator-ns", emitted::add); - - configured.record(clusterScoped(), EventRecord.normal("Created", "created")); + recorder.record( + EventRecord.normal("Created", "created"), context(clusterScoped(), "operator-ns")); assertThat(emitted.get(0).getMetadata().getNamespace()).isEqualTo("operator-ns"); } @Test void anOverriddenClusterScopedNamespaceDoesNotAffectNamespacedResources() { - var configured = new DefaultEventRecorder(CONTROLLER, INSTANCE, "operator-ns", emitted::add); - - configured.record(configMap(), EventRecord.normal("Created", "created")); + recorder.record(EventRecord.normal("Created", "created"), context(configMap(), "operator-ns")); assertThat(emitted.get(0).getMetadata().getNamespace()).isEqualTo("ns1"); } @@ -101,13 +132,13 @@ void anOverriddenClusterScopedNamespaceDoesNotAffectNamespacedResources() { @Test void perEventReportingComponentOverridesTheControllerName() { recorder.record( - configMap(), EventRecord.builder() .reason("Submitted") .message("submitted") .reportingComponent("JobManagerDeployment") .action("Submit") - .build()); + .build(), + context(configMap())); assertThat(emitted.get(0).getReportingComponent()).isEqualTo("JobManagerDeployment"); assertThat(emitted.get(0).getSource().getComponent()).isEqualTo("JobManagerDeployment"); @@ -119,13 +150,13 @@ void perEventReportingComponentOverridesTheControllerName() { @Test void passesLabelsAndAnnotationsThrough() { recorder.record( - configMap(), EventRecord.builder() .reason("Scaling") .message("scaling up") .label("group", "autoscaler") .annotation("recommendation", "4") - .build()); + .build(), + context(configMap())); assertThat(emitted.get(0).getMetadata().getLabels()).containsEntry("group", "autoscaler"); assertThat(emitted.get(0).getMetadata().getAnnotations()).containsEntry("recommendation", "4"); @@ -135,19 +166,29 @@ void passesLabelsAndAnnotationsThrough() { void aFailingSinkNeverFailsTheCaller() { var failing = new DefaultEventRecorder( - CONTROLLER, - INSTANCE, - event -> { + (event, context) -> { throw new RuntimeException("API server said no"); }); - assertThatCode(() -> failing.record(configMap(), EventRecord.normal("Created", "created"))) + assertThatCode( + () -> failing.record(EventRecord.normal("Created", "created"), context(configMap()))) .doesNotThrowAnyException(); } @Test - void boundRecorderRecordsAboutTheBoundObject() { - var bound = recorder.forResource(configMap()); + void passesTheContextOnToTheSink() { + var contexts = new ArrayList>(); + var recording = new DefaultEventRecorder((event, context) -> contexts.add(context)); + var context = context(configMap()); + + recording.record(EventRecord.normal("Created", "created"), context); + + assertThat(contexts).containsExactly(context); + } + + @Test + void boundRecorderRecordsAboutThePrimaryResourceOfTheBoundContext() { + var bound = recorder.forContext(context(configMap())); bound.normal("Created", "created"); bound.warn("Failed", "failed"); @@ -170,7 +211,7 @@ void truncatesTheNameOfTheInvolvedObjectToStayWithinTheKubernetesNameLimit() { .endMetadata() .build(); - recorder.record(configMap, EventRecord.normal("Created", "created")); + recorder.record(EventRecord.normal("Created", "created"), context(configMap)); var name = emitted.get(0).getMetadata().getName(); assertThat(name).hasSizeLessThanOrEqualTo(253); @@ -187,7 +228,7 @@ void reasonIsRequired() { @Test void namesEventsWithADnsSafeHashSuffix() { - recorder.record(configMap(), EventRecord.normal("Created", "created")); + recorder.record(EventRecord.normal("Created", "created"), context(configMap())); assertThat(emitted.get(0).getMetadata().getName()).matches("test1\\.[0-9a-f]{32}"); } @@ -200,14 +241,43 @@ void givesEventsWhoseMessagesCollideUnderStringHashCodeDistinctNames() { // name and the sink would take the second for a repeat of the first and drop it. assertThat("Aa".hashCode()).isEqualTo("BB".hashCode()); - recorder.record(configMap(), EventRecord.warning("Failed", "Aa")); - recorder.record(configMap(), EventRecord.warning("Failed", "BB")); + recorder.record(EventRecord.warning("Failed", "Aa"), context(configMap())); + recorder.record(EventRecord.warning("Failed", "BB"), context(configMap())); assertThat(emitted).hasSize(2); assertThat(emitted.get(0).getMetadata().getName()) .isNotEqualTo(emitted.get(1).getMetadata().getName()); } + Context context(HasMetadata primaryResource) { + return context(primaryResource, DefaultEventRecorder.CLUSTER_SCOPED_EVENT_NAMESPACE); + } + + /** + * The recorder reads everything it does not get from the {@link EventRecord} off the context: the + * primary resource the event is about, the name of the controller the event is attributed to, and + * the configuration service the reporting instance and the cluster scoped event namespace come + * from. + */ + @SuppressWarnings({"unchecked", "rawtypes"}) + Context context(HasMetadata primaryResource, String clusterScopedEventNamespace) { + var configurationService = mock(ConfigurationService.class); + when(configurationService.getLeaderElectionConfiguration()) + .thenReturn( + Optional.of(aLeaderElectionConfiguration("lease").withIdentity(INSTANCE).build())); + when(configurationService.clusterScopedEventNamespace()) + .thenReturn(clusterScopedEventNamespace); + + ControllerConfiguration controllerConfiguration = mock(ControllerConfiguration.class); + when(controllerConfiguration.getName()).thenReturn(CONTROLLER); + when(controllerConfiguration.getConfigurationService()).thenReturn(configurationService); + + Context context = mock(Context.class); + when(context.getPrimaryResource()).thenReturn(primaryResource); + when(context.getControllerConfiguration()).thenReturn(controllerConfiguration); + return context; + } + ConfigMap configMap() { return new ConfigMapBuilder() .withNewMetadata() diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/ControllerTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/ControllerTest.java index b725f49132..de582f593d 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/ControllerTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/ControllerTest.java @@ -28,6 +28,8 @@ import io.javaoperatorsdk.operator.api.config.ConfigurationService; import io.javaoperatorsdk.operator.api.config.MockControllerConfiguration; import io.javaoperatorsdk.operator.api.config.workflow.WorkflowSpec; +import io.javaoperatorsdk.operator.api.event.DefaultEventRecorder; +import io.javaoperatorsdk.operator.api.event.EventRecorder; import io.javaoperatorsdk.operator.api.monitoring.Metrics; import io.javaoperatorsdk.operator.api.reconciler.Cleaner; import io.javaoperatorsdk.operator.api.reconciler.DefaultContext; @@ -110,6 +112,36 @@ void doesNotNotifyMetricsWhenEventProcessorNotStarted() { verify(metrics, never()).eventProcessingStarted(controller); } + @Test + void recordsEventsThroughTheEventRecorderConfiguredForTheOperator() { + final var client = MockKubernetesClient.client(Secret.class); + final var eventRecorder = mock(EventRecorder.class); + final var configurationService = + ConfigurationService.newOverriddenConfigurationService( + new BaseConfigurationService(), + o -> o.withEventRecorder(eventRecorder).withKubernetesClient(client)); + final var configuration = + MockControllerConfiguration.forResource(Secret.class, configurationService); + + final var controller = new Controller(reconciler, configuration, client); + + assertThat(controller.eventRecorder()).isSameAs(eventRecorder); + } + + @Test + void recordsEventsThroughAnEventRecorderOfItsOwnWhenNoneIsConfigured() { + final var client = MockKubernetesClient.client(Secret.class); + final var configuration = + MockControllerConfiguration.forResource( + Secret.class, + ConfigurationService.newOverriddenConfigurationService( + new BaseConfigurationService(), o -> o.withKubernetesClient(client))); + + final var controller = new Controller(reconciler, configuration, client); + + assertThat(controller.eventRecorder()).isInstanceOf(DefaultEventRecorder.class); + } + @Test void crdShouldNotBeCheckedForCustomResourcesIfDisabled() { final var client = MockKubernetesClient.client(TestCustomResource.class); diff --git a/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java b/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java index fedaf81eb6..44fac32b7d 100644 --- a/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java +++ b/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java @@ -122,6 +122,19 @@ void applyConfigsAppliesDurations() { assertThat(result.reconciliationTerminationTimeout()).isEqualTo(Duration.ofSeconds(5)); } + @Test + void applyConfigsAppliesStrings() { + var loader = + new ConfigLoader( + mapProvider(Map.of("josdk.events.cluster-scoped-namespace", "operator-ns"))); + + var base = new BaseConfigurationService(null); + var result = + ConfigurationService.newOverriddenConfigurationService(base, loader.applyConfigs()); + + assertThat(result.clusterScopedEventNamespace()).isEqualTo("operator-ns"); + } + @Test void applyConfigsOnlyAppliesPresentKeys() { // Only one key present — other defaults must be unchanged.