From 3c6ce5f0ae053651555394f78a438bf57f9a202a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Attila=20M=C3=A9sz=C3=A1ros?= Date: Mon, 31 Aug 2026 09:32:51 +0200 Subject: [PATCH 1/5] refactor: configure the event recorder on the operator instead of the controller Removes EventRecorder from RegisteredController and makes the instance the controllers record their Kubernetes events through configurable for the whole operator, via ConfigurationService.eventRecorder() / ConfigurationServiceOverrider.withEventRecorder(). When none is configured, each controller keeps recording through a DefaultEventRecorder of its own, which attributes the events to that controller, as before. Recording events outside of a reconciliation is now done through the configured instance, which the caller owns, rather than through one handed out by the registered controller. --- .../operator/RegisteredController.java | 14 -------- .../api/config/ConfigurationService.java | 20 ++++++++++++ .../config/ConfigurationServiceOverrider.java | 26 +++++++++++++++ .../operator/api/event/EventRecorder.java | 13 +++++--- .../operator/api/reconciler/Context.java | 5 +-- .../operator/processing/Controller.java | 29 +++++++++++------ .../ConfigurationServiceOverriderTest.java | 25 +++++++++++++++ .../operator/processing/ControllerTest.java | 32 +++++++++++++++++++ 8 files changed, 134 insertions(+), 30 deletions(-) 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..d42cf3c6fc 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,25 @@ 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, which attributes the events it records to that controller. A recorder configured here + * is shared by all controllers of the operator, so it decides on its own what the events it + * records are attributed to, and it is up to the caller to hold on to the instance if it also + * records events outside of a reconciliation. + * + * @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..944200c742 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 + * to hold on to the recorder in order to also record events outside of a reconciliation. + * + *

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/EventRecorder.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/event/EventRecorder.java index a7d66d8e54..c4d4919508 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 @@ -16,20 +16,25 @@ package io.javaoperatorsdk.operator.api.event; import io.fabric8.kubernetes.api.model.HasMetadata; +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. + * or a background task. To use it that way, configure the instance the operator records its events + * through, see {@link io.javaoperatorsdk.operator.api.config.ConfigurationService#eventRecorder()}, + * and keep a reference to it. Within a reconciliation, prefer {@link + * io.javaoperatorsdk.operator.api.reconciler.Context#eventRecorder()}, which is already bound to + * the primary resource. * *

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 { /** 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..7c10c0963e 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 @@ -216,8 +216,9 @@ default Optional getSecondaryResource( /** * 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()}. + * events about it. To record events outside of a reconciliation, or about another object, use an + * {@link io.javaoperatorsdk.operator.api.event.EventRecorder} configured for the operator, 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/processing/Controller.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/Controller.java index 285eb3988c..43cb0a8a83 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 @@ -36,6 +36,7 @@ import io.javaoperatorsdk.operator.MissingCRDException; import io.javaoperatorsdk.operator.OperatorException; import io.javaoperatorsdk.operator.RegisteredController; +import io.javaoperatorsdk.operator.api.config.ConfigurationService; import io.javaoperatorsdk.operator.api.config.ControllerConfiguration; import io.javaoperatorsdk.operator.api.config.ExecutorServiceManager; import io.javaoperatorsdk.operator.api.config.LeaderElectionConfiguration; @@ -115,15 +116,19 @@ 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( + configuration.getName(), + configurationService + .getLeaderElectionConfiguration() + .flatMap(LeaderElectionConfiguration::getIdentity) + .orElseGet(DefaultEventRecorder::defaultReportingInstance), + Optional.ofNullable(configurationService.clusterScopedEventNamespace()) + .orElse(DefaultEventRecorder.CLUSTER_SCOPED_EVENT_NAMESPACE), + new DefaultEventSink(kubernetesClient))); contextInitializer = reconciler instanceof ContextInitializer; isCleaner = reconciler instanceof Cleaner; @@ -358,7 +363,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..9d5db71c6f 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,6 +23,9 @@ import org.junit.jupiter.api.Test; import io.fabric8.kubernetes.api.model.HasMetadata; +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 static org.assertj.core.api.Assertions.assertThat; @@ -106,6 +109,28 @@ public R clone(R object) { config.reconciliationTerminationTimeout(), overridden.reconciliationTerminationTimeout()); } + @Test + void eventRecorderIsNotConfiguredByDefaultAndCanBeOverridden() { + final var eventRecorder = + new EventRecorder() { + @Override + public void record(HasMetadata regarding, EventRecord event) {} + + @Override + public ResourceEventRecorder forResource(HasMetadata regarding) { + return null; + } + }; + + assertThat(config.eventRecorder()).isEmpty(); + assertThat( + new ConfigurationServiceOverrider(config) + .withEventRecorder(eventRecorder) + .build() + .eventRecorder()) + .contains(eventRecorder); + } + @Test void threadCountConfiguredProperly() { final var overridden = 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); From ce3f75d833411be40a3a04d0e0ac6a327b7cb561 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Attila=20M=C3=A9sz=C3=A1ros?= Date: Mon, 31 Aug 2026 13:44:07 +0200 Subject: [PATCH 2/5] Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- .../java/io/javaoperatorsdk/operator/processing/Controller.java | 1 - 1 file changed, 1 deletion(-) 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 43cb0a8a83..bdd858e813 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 @@ -36,7 +36,6 @@ import io.javaoperatorsdk.operator.MissingCRDException; import io.javaoperatorsdk.operator.OperatorException; import io.javaoperatorsdk.operator.RegisteredController; -import io.javaoperatorsdk.operator.api.config.ConfigurationService; import io.javaoperatorsdk.operator.api.config.ControllerConfiguration; import io.javaoperatorsdk.operator.api.config.ExecutorServiceManager; import io.javaoperatorsdk.operator.api.config.LeaderElectionConfiguration; From c5ebd7e42710762d404d6a6e2c312e931617f5f4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Attila=20M=C3=A9sz=C3=A1ros?= Date: Wed, 2 Sep 2026 14:24:11 +0200 Subject: [PATCH 3/5] wip MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Attila Mészáros --- .../api/config/ConfigurationService.java | 7 +- .../api/event/DefaultEventRecorder.java | 106 +++++++++------- .../operator/api/event/DefaultEventSink.java | 3 +- .../operator/api/event/EventRecorder.java | 35 +++--- .../operator/api/event/EventSink.java | 9 +- .../api/event/ResourceEventRecorder.java | 5 + .../operator/api/reconciler/Context.java | 7 +- .../api/reconciler/DefaultContext.java | 2 +- .../operator/processing/Controller.java | 13 +- .../ConfigurationServiceOverriderTest.java | 5 +- .../api/event/DefaultEventRecorderTest.java | 118 ++++++++++++++---- 11 files changed, 198 insertions(+), 112 deletions(-) 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 d42cf3c6fc..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 @@ -302,10 +302,9 @@ default String clusterScopedEventNamespace() { * 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, which attributes the events it records to that controller. A recorder configured here - * is shared by all controllers of the operator, so it decides on its own what the events it - * records are attributed to, and it is up to the caller to hold on to the instance if it also - * records events outside of a reconciliation. + * 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 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 c4d4919508..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,7 +15,7 @@ */ 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; @@ -23,13 +23,16 @@ /** * 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. To use it that way, configure the instance the operator records its events - * through, see {@link io.javaoperatorsdk.operator.api.config.ConfigurationService#eventRecorder()}, - * and keep a reference to it. Within a reconciliation, prefer {@link + *

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 primary resource. + * 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. @@ -38,20 +41,20 @@ 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..6db5a983c6 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,7 @@ package io.javaoperatorsdk.operator.api.event; import io.fabric8.kubernetes.api.model.Event; +import io.javaoperatorsdk.operator.api.reconciler.Context; /** * Writes fully built events somewhere. Extracted from {@link EventRecorder} so that the assembly of @@ -29,7 +30,11 @@ 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); + 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 7c10c0963e..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,10 +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 an - * {@link io.javaoperatorsdk.operator.api.event.EventRecorder} configured for the operator, see - * {@link io.javaoperatorsdk.operator.api.config.ConfigurationService#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 bdd858e813..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; @@ -117,17 +116,7 @@ public Controller( this.eventRecorder = configurationService .eventRecorder() - .orElseGet( - () -> - 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))); + .orElseGet(() -> new DefaultEventRecorder(new DefaultEventSink(kubernetesClient))); contextInitializer = reconciler instanceof ContextInitializer; isCleaner = reconciler instanceof Cleaner; 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 9d5db71c6f..fe92e63275 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 @@ -27,6 +27,7 @@ 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; @@ -114,10 +115,10 @@ void eventRecorderIsNotConfiguredByDefaultAndCanBeOverridden() { final var eventRecorder = new EventRecorder() { @Override - public void record(HasMetadata regarding, EventRecord event) {} + public void record(EventRecord event, Context context) {} @Override - public ResourceEventRecorder forResource(HasMetadata regarding) { + public ResourceEventRecorder forContext(Context context) { return null; } }; 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() From 28dd54b9a2f88f90f9231af576af794fbe4aadd2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Attila=20M=C3=A9sz=C3=A1ros?= Date: Wed, 2 Sep 2026 16:29:18 +0200 Subject: [PATCH 4/5] wip MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Attila Mészáros --- .../java/io/javaoperatorsdk/operator/api/event/EventSink.java | 4 ++++ 1 file changed, 4 insertions(+) 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 6db5a983c6..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 @@ -17,6 +17,9 @@ 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 @@ -36,5 +39,6 @@ public interface EventSink { * implementations that route the event based on it rather than on its contents. Ignored by * {@link DefaultEventSink}. */ + @Experimental(API_MIGHT_CHANGE) void emit(Event event, Context context); } From 5f33a74cb22ed862cf3d3aed435be6218a2ef198 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Attila=20M=C3=A9sz=C3=A1ros?= Date: Wed, 2 Sep 2026 16:32:35 +0200 Subject: [PATCH 5/5] wip MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Attila Mészáros --- .../operator/api/config/ConfigurationServiceOverrider.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 944200c742..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 @@ -155,7 +155,7 @@ public ConfigurationServiceOverrider withClusterScopedEventNamespace(String name * 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 - * to hold on to the recorder in order to also record events outside of a reconciliation. + * 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.