From d919b94c88ae4d60d47dfa6c03d481c1d12f1cab Mon Sep 17 00:00:00 2001 From: Kai Huang Date: Fri, 1 May 2026 12:09:38 -0700 Subject: [PATCH] Add AnalyticsFrontEndExtension SPI + producer-side wiring in AnalyticsPlugin MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds the SPI surface in analytics-framework that lets frontend plugins (e.g., opensearch-sql) consume analytics-engine services without taking a hard install-time dependency on analytics-engine. Adds the producer-side wiring in AnalyticsPlugin: discover consumers via ExtensiblePlugin#loadExtensions, push the executor + schemaProvider bundle once Guice constructs DefaultPlanExecutor. Mirrors the JobSchedulerExtension pattern from opensearch-job-scheduler. analytics-framework (sandbox/libs/analytics-framework): - AnalyticsFrontEndExtension — single-method SPI interface, setAnalyticsServices(AnalyticsServices). Implementers declared via standard Java SPI (META-INF/services). - AnalyticsServices — record bundling QueryPlanExecutor and SchemaProvider. Bundled rather than separate setters so future services can be added without changing the SPI signature. analytics-engine (sandbox/plugins/analytics-engine): - AnalyticsPlugin#loadExtensions: collect AnalyticsFrontEndExtension consumers alongside the existing AnalyticsSearchBackendPlugin collection. - AnalyticsPlugin#createGuiceModules: bind a Guice TypeListener that fires when DefaultPlanExecutor is instantiated. The InjectionListener wraps OpenSearchSchemaBuilder::buildSchema as a SchemaProvider, packs both into AnalyticsServices, and pushes to every registered consumer (per-consumer try/catch so one bad consumer doesn't crash the others). - The push is idempotent — guarded by a servicesPushed flag so Guice re-instantiation (no Singleton scope today) doesn't re-fire the callback. Verified end-to-end against opensearch-sql: - WITH analytics-engine: SQL routes parquet_* queries through the pushed executor (confirmed via /_plugins/_ppl/_explain showing Calcite logical plan). - WITHOUT analytics-engine: SQL plugin loads fine; parquet_* queries fall through to the legacy path with a clean IndexNotFoundException — no NoClassDefFoundError, no startup crash. Pairs with the SQL-side consumer wiring in opensearch-project/sql. Signed-off-by: Kai Huang --- .../spi/AnalyticsFrontEndExtension.java | 42 +++++++++++++++++++ .../analytics/spi/AnalyticsServices.java | 29 +++++++++++++ .../opensearch/analytics/AnalyticsPlugin.java | 41 ++++++++++++++++++ 3 files changed, 112 insertions(+) create mode 100644 sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsFrontEndExtension.java create mode 100644 sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsServices.java diff --git a/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsFrontEndExtension.java b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsFrontEndExtension.java new file mode 100644 index 0000000000000..0dc2166b8642d --- /dev/null +++ b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsFrontEndExtension.java @@ -0,0 +1,42 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.analytics.spi; + +/** + * SPI for frontend plugins (e.g., opensearch-sql) that integrate with analytics-engine. + * + *

Implementers are discovered by {@code AnalyticsPlugin} via + * {@link org.opensearch.plugins.ExtensiblePlugin#loadExtensions}; analytics-engine pushes its + * services to each consumer once Guice has constructed them. Mirrors the + * {@code JobSchedulerExtension} pattern from opensearch-job-scheduler — the consumer plugin + * declares its capability via this interface; the publishing plugin (analytics-engine) handles + * discovery and lifecycle. + * + *

This SPI lets a frontend declare analytics-engine as an OPTIONAL extended plugin + * ({@code extendedPlugins = ['analytics-engine;optional=true']}). When analytics-engine is not + * installed, no consumer ever receives a callback; analytics-routing code paths stay inert and + * the frontend plugin boots normally. + * + *

Lifecycle. {@link #setAnalyticsServices} is invoked exactly once per consumer per node + * lifecycle, AFTER the node Guice injector is built (i.e., after every plugin's + * {@code createComponents} returns) and BEFORE the first analytics query is dispatched. + * Implementations should stash the bundle for later use; do not assume the services are available + * during {@code createComponents}. + * + * @opensearch.internal + */ +public interface AnalyticsFrontEndExtension { + + /** + * Receives the bundle of analytics-engine services. Called exactly once after the services are + * constructed and before any analytics query is dispatched. Each service inside the bundle is + * safe to invoke from any thread once received. + */ + void setAnalyticsServices(AnalyticsServices services); +} diff --git a/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsServices.java b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsServices.java new file mode 100644 index 0000000000000..c6e9d5a0cec3f --- /dev/null +++ b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsServices.java @@ -0,0 +1,29 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.analytics.spi; + +import org.apache.calcite.rel.RelNode; +import org.opensearch.analytics.exec.QueryPlanExecutor; +import org.opensearch.analytics.schema.SchemaProvider; + +/** + * Bundle of services that {@code AnalyticsPlugin} pushes to each {@link AnalyticsFrontEndExtension} + * consumer once Guice has constructed them. + * + *

Bundled rather than passed through individual setters so future analytics-engine services can + * be added without changing the {@link AnalyticsFrontEndExtension} signature — frontends that do + * not consume the new service simply ignore the new accessor. + * + * @param queryPlanExecutor coordinator-level query plan executor + * @param schemaProvider builds a Calcite {@code SchemaPlus} from the current cluster state + * + * @opensearch.internal + */ +public record AnalyticsServices(QueryPlanExecutor> queryPlanExecutor, SchemaProvider schemaProvider) { +} diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java index 6c5786847a4d7..a20ec4daa0507 100644 --- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java +++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java @@ -24,11 +24,19 @@ import org.opensearch.analytics.planner.CapabilityRegistry; import org.opensearch.analytics.planner.FieldStorageResolver; import org.opensearch.analytics.schema.OpenSearchSchemaBuilder; +import org.opensearch.analytics.schema.SchemaProvider; +import org.opensearch.analytics.spi.AnalyticsFrontEndExtension; import org.opensearch.analytics.spi.AnalyticsSearchBackendPlugin; +import org.opensearch.analytics.spi.AnalyticsServices; +import org.opensearch.cluster.ClusterState; import org.opensearch.cluster.metadata.IndexNameExpressionResolver; import org.opensearch.cluster.service.ClusterService; import org.opensearch.common.inject.Module; import org.opensearch.common.inject.TypeLiteral; +import org.opensearch.common.inject.matcher.Matchers; +import org.opensearch.common.inject.spi.InjectionListener; +import org.opensearch.common.inject.spi.TypeEncounter; +import org.opensearch.common.inject.spi.TypeListener; import org.opensearch.core.action.ActionResponse; import org.opensearch.core.common.io.stream.NamedWriteableRegistry; import org.opensearch.core.xcontent.NamedXContentRegistry; @@ -66,12 +74,14 @@ public class AnalyticsPlugin extends Plugin implements ExtensiblePlugin, ActionP public AnalyticsPlugin() {} private final List backEnds = new ArrayList<>(); + private final List frontEnds = new ArrayList<>(); private SqlOperatorTable operatorTable; @SuppressWarnings("rawtypes") @Override public void loadExtensions(ExtensionLoader loader) { backEnds.addAll(loader.loadExtensions(AnalyticsSearchBackendPlugin.class)); + frontEnds.addAll(loader.loadExtensions(AnalyticsFrontEndExtension.class)); } @Override @@ -112,9 +122,40 @@ public Collection createGuiceModules() { }).to(DefaultPlanExecutor.class); b.bind(EngineContext.class).to(DefaultEngineContext.class); b.bind(Scheduler.class).to(QueryScheduler.class); + // Push the executor + schemaProvider bundle to every registered AnalyticsFrontEndExtension + // once Guice constructs DefaultPlanExecutor. The InjectionListener fires on the singleton + // construction; pushAnalyticsServices guards against re-firing if Guice ever instantiates + // more than once. + b.bindListener(Matchers.any(), new TypeListener() { + @Override + public void hear(TypeLiteral type, TypeEncounter encounter) { + if (!DefaultPlanExecutor.class.isAssignableFrom(type.getRawType())) { + return; + } + encounter.register((InjectionListener) instance -> pushAnalyticsServices((DefaultPlanExecutor) instance)); + } + }); }); } + private boolean servicesPushed = false; + + private synchronized void pushAnalyticsServices(DefaultPlanExecutor executor) { + if (servicesPushed) { + return; + } + servicesPushed = true; + SchemaProvider schemaProvider = clusterState -> OpenSearchSchemaBuilder.buildSchema((ClusterState) clusterState); + AnalyticsServices services = new AnalyticsServices(executor, schemaProvider); + for (AnalyticsFrontEndExtension consumer : frontEnds) { + try { + consumer.setAnalyticsServices(services); + } catch (Exception e) { + logger.warn("AnalyticsFrontEndExtension {} threw on setAnalyticsServices", consumer.getClass().getName(), e); + } + } + } + @Override public List> getActions() { return List.of(new ActionHandler<>(AnalyticsQueryAction.INSTANCE, DefaultPlanExecutor.class));