plannerRules() {
+ return List.of();
+ }
+
+ /**
+ * 创建保留 ANSI 双引号输入、继承目标方言大小写语义的解析配置。
+ *
+ * 数据库若对表名与列名采用不同的大小写规则,可以在 Adapter 中覆盖。
+ *
+ * @param dialect 目标数据库方言
+ * @return Calcite 解析配置
+ */
+ default SqlParser.Config parserConfig(SqlDialect dialect) {
+ return SqlParser.config()
+ .withQuotedCasing(dialect.getQuotedCasing())
+ .withUnquotedCasing(dialect.getUnquotedCasing())
+ .withCaseSensitive(dialect.isCaseSensitive());
+ }
+
+ /**
+ * 将 JDBC 参数类型映射为 Calcite 类型声明,供动态参数参与校验和类型推导。
+ *
+ * 默认实现补齐 JDBC 4.2 时区类型,并将 {@link Types#OTHER} 解释为 UUID。
+ * 厂商 Adapter 可以覆盖此方法,直接返回带精度、长度或专有类型名的
+ * Calcite 类型声明。
+ *
+ * @param jdbcType {@link java.sql.Types} 类型值
+ * @param parserPosition 动态参数的解析位置
+ * @return Calcite 类型声明;无法映射时返回 null
+ */
+ default SqlDataTypeSpec parameterTypeSpec(int jdbcType, SqlParserPos parserPosition) {
+ SqlTypeName typeName = switch (jdbcType) {
+ case Types.TIME_WITH_TIMEZONE -> SqlTypeName.TIME_TZ;
+ case Types.TIMESTAMP_WITH_TIMEZONE -> SqlTypeName.TIMESTAMP_TZ;
+ case Types.OTHER -> SqlTypeName.UUID;
+ default -> SqlTypeName.getNameForJdbcType(jdbcType);
+ };
+ if (typeName == null || typeName.isSpecial() || !typeName.allowsNoPrecNoScale()) {
+ return null;
+ }
+ return new SqlDataTypeSpec(
+ new SqlBasicTypeNameSpec(typeName, parserPosition),
+ parserPosition
+ );
+ }
+
+ /**
+ * 返回目标数据库 Fragment 执行器。
+ *
+ * @return Fragment 执行器
+ */
+ FederationFragmentExecutor fragmentExecutor();
+
+ /**
+ * 返回可选的数据库物理 Explain 实现。
+ *
+ * @return 物理 Explain SPI
+ */
+ default Optional fragmentExplainer() {
+ return Optional.empty();
+ }
+
+ /**
+ * 返回可选的数据库目录统计采集器。
+ *
+ * 统计采集由引擎管理缓存、并发合并、失效和失败降级,Adapter 只负责
+ * 当前数据库的目录语义。
+ *
+ * @return 统计采集 SPI;未适配时为空
+ */
+ default Optional statisticsCollector() {
+ return Optional.empty();
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationSqlAdapterRegistry.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationSqlAdapterRegistry.java
new file mode 100644
index 0000000..06375b2
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationSqlAdapterRegistry.java
@@ -0,0 +1,94 @@
+package com.easyagents.federation.sql.adapter;
+
+import com.easyagents.federation.sql.api.FederationSqlErrorCode;
+import com.easyagents.federation.sql.api.FederationSqlException;
+import java.util.Collection;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.ServiceLoader;
+import java.util.Set;
+
+/**
+ * 支持显式注册与 ServiceLoader 的 Adapter 注册表。
+ */
+public final class FederationSqlAdapterRegistry {
+
+ private final Map providers;
+
+ /**
+ * 创建注册表;显式 Provider 优先于 ServiceLoader Provider。
+ *
+ * @param explicitProviders 显式 Provider
+ * @param classLoader ServiceLoader 使用的类加载器
+ */
+ public FederationSqlAdapterRegistry(
+ Collection explicitProviders,
+ ClassLoader classLoader
+ ) {
+ Map loaded = new LinkedHashMap<>();
+ ServiceLoader.load(FederationSqlAdapterProvider.class, classLoader)
+ .forEach(provider -> putUnique(loaded, provider));
+ if (explicitProviders != null) {
+ Set explicitIds = new HashSet<>();
+ for (FederationSqlAdapterProvider provider : explicitProviders) {
+ Objects.requireNonNull(provider, "adapter provider must not be null");
+ if (!explicitIds.add(provider.adapterId())) {
+ throw new FederationSqlException(
+ FederationSqlErrorCode.SOURCE_DEFINITION_CONFLICT,
+ "duplicate explicitly registered adapter id: " + provider.adapterId()
+ );
+ }
+ loaded.put(provider.adapterId(), provider);
+ }
+ }
+ this.providers = Map.copyOf(loaded);
+ }
+
+ private static void putUnique(
+ Map providers,
+ FederationSqlAdapterProvider provider
+ ) {
+ FederationSqlAdapterProvider previous = providers.putIfAbsent(provider.adapterId(), provider);
+ if (previous != null && !previous.getClass().equals(provider.getClass())) {
+ throw new FederationSqlException(
+ FederationSqlErrorCode.SOURCE_DEFINITION_CONFLICT,
+ "duplicate adapter id from ServiceLoader: " + provider.adapterId()
+ );
+ }
+ }
+
+ /**
+ * 查找 Adapter Provider。
+ *
+ * @param adapterId Adapter 标识
+ * @return 可选 Provider
+ */
+ public Optional find(String adapterId) {
+ return Optional.ofNullable(providers.get(adapterId));
+ }
+
+ /**
+ * 返回 Adapter Provider,缺失时抛出稳定错误。
+ *
+ * @param adapterId Adapter 标识
+ * @return Provider
+ */
+ public FederationSqlAdapterProvider require(String adapterId) {
+ return find(adapterId).orElseThrow(() -> new FederationSqlException(
+ FederationSqlErrorCode.ADAPTER_NOT_FOUND,
+ "adapter is not registered: " + adapterId
+ ));
+ }
+
+ /**
+ * 返回不可变 Provider 视图。
+ *
+ * @return Provider 映射
+ */
+ public Map providers() {
+ return providers;
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationStatisticsCollectionContext.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationStatisticsCollectionContext.java
new file mode 100644
index 0000000..ab7f808
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationStatisticsCollectionContext.java
@@ -0,0 +1,40 @@
+package com.easyagents.federation.sql.adapter;
+
+import com.easyagents.federation.sql.source.FederationSourceDefinition;
+import java.sql.Connection;
+import java.time.Instant;
+import java.util.Objects;
+
+/**
+ * Adapter 采集数据库目录统计时使用的只读上下文。
+ *
+ * @param sourceDefinition 当前物理数据源定义
+ * @param connection 已从运行时连接池借出的 JDBC 连接
+ * @param collectedAt 本轮统计采集时间
+ * @param expiresAt 本轮统计默认失效时间
+ * @param queryTimeoutSeconds 单条目录查询超时秒数
+ */
+public record FederationStatisticsCollectionContext(
+ FederationSourceDefinition sourceDefinition,
+ Connection connection,
+ Instant collectedAt,
+ Instant expiresAt,
+ int queryTimeoutSeconds
+) {
+
+ /**
+ * 校验统计采集上下文。
+ */
+ public FederationStatisticsCollectionContext {
+ sourceDefinition = Objects.requireNonNull(sourceDefinition, "sourceDefinition");
+ connection = Objects.requireNonNull(connection, "connection");
+ collectedAt = Objects.requireNonNull(collectedAt, "collectedAt");
+ expiresAt = Objects.requireNonNull(expiresAt, "expiresAt");
+ if (!expiresAt.isAfter(collectedAt)) {
+ throw new IllegalArgumentException("expiresAt must be after collectedAt");
+ }
+ if (queryTimeoutSeconds <= 0) {
+ throw new IllegalArgumentException("queryTimeoutSeconds must be positive");
+ }
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationStatisticsCollector.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationStatisticsCollector.java
new file mode 100644
index 0000000..68a7b35
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/adapter/FederationStatisticsCollector.java
@@ -0,0 +1,27 @@
+package com.easyagents.federation.sql.adapter;
+
+import com.easyagents.federation.sql.federation.FederationStatisticsSnapshot;
+import com.easyagents.federation.sql.federation.FederationTableStatistics;
+import java.sql.SQLException;
+import java.util.Map;
+
+/**
+ * 数据库 Adapter 提供的轻量目录统计采集 SPI。
+ *
+ * 实现应使用数据库系统目录或 JDBC 元数据批量采集,禁止执行逐表
+ * {@code COUNT(*)}。采集异常由引擎统一降级,不应在实现中伪造成功结果。
+ */
+@FunctionalInterface
+public interface FederationStatisticsCollector {
+
+ /**
+ * 采集一个物理数据源当前 revision 的表统计。
+ *
+ * @param context 统计采集上下文
+ * @return 按逻辑 Schema 和物理表索引的不可变统计;不支持时返回空映射
+ * @throws SQLException 数据库目录或 JDBC 元数据读取失败
+ */
+ Map collect(
+ FederationStatisticsCollectionContext context
+ ) throws SQLException;
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationCleanupMetrics.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationCleanupMetrics.java
new file mode 100644
index 0000000..de2ab07
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationCleanupMetrics.java
@@ -0,0 +1,35 @@
+package com.easyagents.federation.sql.api;
+
+import java.io.Serializable;
+
+/**
+ * 节点本地 JDBC 终止与游标清理通道的累计观测指标。
+ *
+ * @param overflowFallbacks 主清理队列拒绝后转入隔离通道的次数
+ * @param deferredRetries 隔离通道拒绝后进入有界延期重试队列的次数
+ * @param unresolvedCleanups 延期队列溢出或 Engine 有界关闭后仍未完成的清理数
+ * @param deferredQueueDepth 当前等待重试的清理数
+ */
+public record FederationCleanupMetrics(
+ long overflowFallbacks,
+ long deferredRetries,
+ long unresolvedCleanups,
+ int deferredQueueDepth
+) implements Serializable {
+
+ private static final FederationCleanupMetrics EMPTY = new FederationCleanupMetrics(
+ 0L,
+ 0L,
+ 0L,
+ 0
+ );
+
+ /**
+ * 返回无清理压力的空指标。
+ *
+ * @return 空指标
+ */
+ public static FederationCleanupMetrics empty() {
+ return EMPTY;
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlEngine.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlEngine.java
new file mode 100644
index 0000000..229a74e
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlEngine.java
@@ -0,0 +1,88 @@
+package com.easyagents.federation.sql.api;
+
+import com.easyagents.federation.sql.compile.FederationSqlPlan;
+import com.easyagents.federation.sql.compile.SqlCompileRequest;
+import com.easyagents.federation.sql.compile.SqlExplainRequest;
+import com.easyagents.federation.sql.compile.SqlExplainResult;
+import com.easyagents.federation.sql.execute.FederationResultCursor;
+import com.easyagents.federation.sql.execute.QueryId;
+import com.easyagents.federation.sql.source.FederationSourceManager;
+
+/**
+ * SQL 编译、补全、查询、Explain、取消和数据源管理的统一公共入口。
+ */
+public interface FederationSqlEngine extends AutoCloseable {
+
+ /**
+ * 返回数据源管理入口。
+ *
+ * @return 数据源管理器
+ */
+ FederationSourceManager sources();
+
+ /**
+ * 编译节点本地计划。
+ *
+ * @param request 编译请求
+ * @return 节点本地计划
+ */
+ FederationSqlPlan compile(SqlCompileRequest request);
+
+ /**
+ * 执行节点本地计划。
+ *
+ * @param plan 编译计划
+ * @param context 执行上下文
+ * @return 流式游标
+ */
+ FederationResultCursor execute(FederationSqlPlan plan, SqlExecutionContext context);
+
+ /**
+ * 在当前节点完成编译或缓存命中并立即执行。
+ *
+ * @param command 可跨节点查询命令
+ * @return 流式游标
+ */
+ FederationResultCursor query(SqlQueryCommand command);
+
+ /**
+ * 返回不含运行对象的 Explain 结果。
+ *
+ * @param request Explain 请求
+ * @return Explain 结果
+ */
+ SqlExplainResult explain(SqlExplainRequest request);
+
+ /**
+ * 根据当前查询范围返回 Calcite SQL 上下文补全候选。
+ *
+ * @param request 补全请求
+ * @return 替换区间与候选列表
+ */
+ SqlCompletionResult complete(SqlCompletionRequest request);
+
+ /**
+ * 尝试取消当前节点正在执行的查询。
+ *
+ * @param queryId 查询标识
+ * @return 是否找到并发起取消
+ */
+ boolean cancel(QueryId queryId);
+
+ /**
+ * 返回节点本地 JDBC 终止与游标清理通道的累计指标。
+ *
+ * 自定义 Engine 未提供资源治理指标时返回空快照。
+ *
+ * @return 清理通道指标
+ */
+ default FederationCleanupMetrics cleanupMetrics() {
+ return FederationCleanupMetrics.empty();
+ }
+
+ /**
+ * 关闭 Engine、订阅、Runtime 和独占句柄。
+ */
+ @Override
+ void close();
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlEngines.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlEngines.java
new file mode 100644
index 0000000..570fb87
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlEngines.java
@@ -0,0 +1,277 @@
+package com.easyagents.federation.sql.api;
+
+import com.easyagents.federation.sql.adapter.FederationSqlAdapterProvider;
+import com.easyagents.federation.sql.adapter.FederationSqlAdapterRegistry;
+import com.easyagents.federation.sql.compile.FederationSqlPolicy;
+import com.easyagents.federation.sql.execute.FederationQueryAdmissionController;
+import com.easyagents.federation.sql.execute.LocalFederationQueryAdmissionController;
+import com.easyagents.federation.sql.federation.FederationExecutionPolicy;
+import com.easyagents.federation.sql.federation.FederationTableStatisticsProvider;
+import com.easyagents.federation.sql.runtime.DefaultFederationSqlEngine;
+import com.easyagents.federation.sql.runtime.DefaultFederationSourceManager;
+import com.easyagents.federation.sql.source.FederationDataSourceResolver;
+import com.easyagents.federation.sql.source.FederationSourceStateProvider;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Objects;
+import java.time.Duration;
+
+/**
+ * 使用显式依赖构建独立 FederationSqlEngine 的入口。
+ */
+public final class FederationSqlEngines {
+
+ private FederationSqlEngines() {
+ }
+
+ /**
+ * 创建 Engine Builder。
+ *
+ * @return Builder
+ */
+ public static Builder builder() {
+ return new Builder();
+ }
+
+ /**
+ * FederationSqlEngine 的轻量配置 Builder。
+ */
+ public static final class Builder {
+
+ private final List adapters = new ArrayList<>();
+ private final List policies = new ArrayList<>();
+ private FederationDataSourceResolver resolver;
+ private FederationSourceStateProvider stateProvider = FederationSourceStateProvider.none();
+ private FederationQueryAdmissionController admissionController =
+ new LocalFederationQueryAdmissionController(64);
+ private int maximumPlanCacheEntries = 1024;
+ private long maximumPlanCacheWeightBytes = 64L * 1024L * 1024L;
+ private Duration planCacheTimeToLive = Duration.ofMinutes(30);
+ private int maximumConcurrentCompilations = Math.max(
+ 1,
+ Math.min(8, Runtime.getRuntime().availableProcessors())
+ );
+ private boolean crossSourceEnabled = true;
+ private FederationExecutionPolicy executionPolicy = FederationExecutionPolicy.basic();
+ private long maximumNodeIntermediateBytes = 512L * 1024L * 1024L;
+ private FederationTableStatisticsProvider statisticsProvider;
+ private ClassLoader classLoader = Thread.currentThread().getContextClassLoader();
+
+ private Builder() {
+ }
+
+ /**
+ * 设置调用方 DataSource Resolver。
+ *
+ * @param resolver Resolver
+ * @return 当前 Builder
+ */
+ public Builder dataSourceResolver(FederationDataSourceResolver resolver) {
+ this.resolver = Objects.requireNonNull(resolver, "resolver must not be null");
+ return this;
+ }
+
+ /**
+ * 显式注册 Adapter;同 id 时覆盖 ServiceLoader 实现。
+ *
+ * @param adapter Adapter Provider
+ * @return 当前 Builder
+ */
+ public Builder adapter(FederationSqlAdapterProvider adapter) {
+ this.adapters.add(Objects.requireNonNull(adapter, "adapter must not be null"));
+ return this;
+ }
+
+ /**
+ * 增加 SQL 策略。
+ *
+ * @param policy 策略
+ * @return 当前 Builder
+ */
+ public Builder policy(FederationSqlPolicy policy) {
+ this.policies.add(Objects.requireNonNull(policy, "policy must not be null"));
+ return this;
+ }
+
+ /**
+ * 设置共享状态 Provider。
+ *
+ * @param stateProvider 状态 Provider
+ * @return 当前 Builder
+ */
+ public Builder stateProvider(FederationSourceStateProvider stateProvider) {
+ this.stateProvider = Objects.requireNonNull(stateProvider, "stateProvider must not be null");
+ return this;
+ }
+
+ /**
+ * 设置查询准入控制器。
+ *
+ * @param admissionController 准入控制器
+ * @return 当前 Builder
+ */
+ public Builder admissionController(FederationQueryAdmissionController admissionController) {
+ this.admissionController = Objects.requireNonNull(
+ admissionController,
+ "admissionController must not be null"
+ );
+ return this;
+ }
+
+ /**
+ * 设置计划缓存最大条目数。
+ *
+ * @param maximumPlanCacheEntries 最大条目数
+ * @return 当前 Builder
+ */
+ public Builder maximumPlanCacheEntries(int maximumPlanCacheEntries) {
+ if (maximumPlanCacheEntries <= 0) {
+ throw new IllegalArgumentException("maximumPlanCacheEntries must be positive");
+ }
+ this.maximumPlanCacheEntries = maximumPlanCacheEntries;
+ return this;
+ }
+
+ /**
+ * 设置计划缓存最大估算权重。
+ *
+ * @param maximumPlanCacheWeightBytes 最大估算字节数
+ * @return 当前 Builder
+ */
+ public Builder maximumPlanCacheWeightBytes(long maximumPlanCacheWeightBytes) {
+ if (maximumPlanCacheWeightBytes <= 0) {
+ throw new IllegalArgumentException("maximumPlanCacheWeightBytes must be positive");
+ }
+ this.maximumPlanCacheWeightBytes = maximumPlanCacheWeightBytes;
+ return this;
+ }
+
+ /**
+ * 设置计划缓存条目存活时间。
+ *
+ * @param planCacheTimeToLive 存活时间
+ * @return 当前 Builder
+ */
+ public Builder planCacheTimeToLive(Duration planCacheTimeToLive) {
+ if (planCacheTimeToLive == null
+ || planCacheTimeToLive.isZero()
+ || planCacheTimeToLive.isNegative()) {
+ throw new IllegalArgumentException("planCacheTimeToLive must be positive");
+ }
+ this.planCacheTimeToLive = planCacheTimeToLive;
+ return this;
+ }
+
+ /**
+ * 设置 Calcite 冷编译最大并发数。
+ *
+ * @param maximumConcurrentCompilations 最大并发冷编译数
+ * @return 当前 Builder
+ */
+ public Builder maximumConcurrentCompilations(int maximumConcurrentCompilations) {
+ if (maximumConcurrentCompilations <= 0) {
+ throw new IllegalArgumentException("maximumConcurrentCompilations must be positive");
+ }
+ this.maximumConcurrentCompilations = maximumConcurrentCompilations;
+ return this;
+ }
+
+ /**
+ * 设置联邦执行开关。
+ *
+ * @param enabled 是否开启
+ * @return 当前 Builder
+ */
+ public Builder crossSourceEnabled(boolean enabled) {
+ this.crossSourceEnabled = enabled;
+ return this;
+ }
+
+ /**
+ * 设置 Engine 级联邦资源硬上限。
+ *
+ * @param executionPolicy 资源策略
+ * @return 当前 Builder
+ */
+ public Builder federationExecutionPolicy(FederationExecutionPolicy executionPolicy) {
+ this.executionPolicy = Objects.requireNonNull(
+ executionPolicy,
+ "executionPolicy must not be null"
+ );
+ return this;
+ }
+
+ /**
+ * 设置节点同时预留的联邦中间结果内存总上限。
+ *
+ * @param maximumNodeIntermediateBytes 节点内存上限
+ * @return 当前 Builder
+ */
+ public Builder maximumNodeIntermediateBytes(long maximumNodeIntermediateBytes) {
+ if (maximumNodeIntermediateBytes <= 0L) {
+ throw new IllegalArgumentException(
+ "maximumNodeIntermediateBytes must be positive"
+ );
+ }
+ this.maximumNodeIntermediateBytes = maximumNodeIntermediateBytes;
+ return this;
+ }
+
+ /**
+ * 设置联邦表统计 Provider,覆盖引擎内建的 Adapter 自动采集能力。
+ *
+ * @param statisticsProvider 调用方完全托管的只读统计快照 Provider
+ * @return 当前 Builder
+ */
+ public Builder tableStatisticsProvider(
+ FederationTableStatisticsProvider statisticsProvider
+ ) {
+ this.statisticsProvider = Objects.requireNonNull(
+ statisticsProvider,
+ "statisticsProvider must not be null"
+ );
+ return this;
+ }
+
+ /**
+ * 设置 ServiceLoader 类加载器。
+ *
+ * @param classLoader 类加载器
+ * @return 当前 Builder
+ */
+ public Builder classLoader(ClassLoader classLoader) {
+ this.classLoader = Objects.requireNonNull(classLoader, "classLoader must not be null");
+ return this;
+ }
+
+ /**
+ * 构建独立 Engine。
+ *
+ * @return Engine
+ */
+ public FederationSqlEngine build() {
+ if (resolver == null) {
+ throw new IllegalStateException("dataSourceResolver must be configured");
+ }
+ FederationSqlAdapterRegistry registry = new FederationSqlAdapterRegistry(adapters, classLoader);
+ DefaultFederationSourceManager sourceManager = new DefaultFederationSourceManager(
+ resolver,
+ registry,
+ stateProvider
+ );
+ return new DefaultFederationSqlEngine(
+ sourceManager,
+ admissionController,
+ policies,
+ maximumPlanCacheEntries,
+ maximumConcurrentCompilations,
+ crossSourceEnabled,
+ executionPolicy,
+ maximumPlanCacheWeightBytes,
+ planCacheTimeToLive,
+ statisticsProvider,
+ maximumNodeIntermediateBytes
+ );
+ }
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlErrorCode.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlErrorCode.java
new file mode 100644
index 0000000..821bf80
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlErrorCode.java
@@ -0,0 +1,71 @@
+package com.easyagents.federation.sql.api;
+
+/**
+ * SQL 联邦查询稳定错误码。
+ */
+public enum FederationSqlErrorCode {
+ /** 公共参数不合法。 */
+ INVALID_ARGUMENT,
+ /** Engine 或 SourceManager 已关闭。 */
+ ENGINE_CLOSED,
+ /** SQL 解析失败。 */
+ SQL_PARSE_FAILED,
+ /** SQL 校验失败。 */
+ SQL_VALIDATION_FAILED,
+ /** SQL 超出只读查询基线。 */
+ SQL_NOT_READ_ONLY,
+ /** SQL 编译或关系转换失败。 */
+ SQL_COMPILE_FAILED,
+ /** SQL 冷编译等待或编译过程超过统一时限。 */
+ SQL_COMPILE_TIMEOUT,
+ /** SQL 编辑器补全失败。 */
+ SQL_COMPLETION_FAILED,
+ /** 单源计划仍含不可执行的本地残余算子。 */
+ SQL_NOT_FULLY_PUSHDOWN,
+ /** 跨数据源能力未开启。 */
+ CROSS_SOURCE_DISABLED,
+ /** 跨数据源执行在当前阶段未实现。 */
+ CROSS_SOURCE_EXECUTION_UNSUPPORTED,
+ /** 查询范围或 Binding 声明不合法。 */
+ INVALID_QUERY_SCOPE,
+ /** 联邦本地执行暂不支持当前关系算子。 */
+ FEDERATION_OPERATOR_UNSUPPORTED,
+ /** 联邦中间结果行数、字节数或执行时间超过限制。 */
+ FEDERATION_RESOURCE_LIMIT_EXCEEDED,
+ /** 节点本地计划绑定的 Runtime 身份已经失效。 */
+ PLAN_STALE,
+ /** 数据源未登记且无法从共享状态恢复。 */
+ SOURCE_NOT_FOUND,
+ /** 数据源已被墓碑删除。 */
+ SOURCE_REMOVED,
+ /** 节点本地数据源版本不满足请求。 */
+ SOURCE_REVISION_NOT_READY,
+ /** 同 revision 出现不同 Definition 校验和。 */
+ SOURCE_DEFINITION_CONFLICT,
+ /** 数据源 Runtime 初始化失败。 */
+ SOURCE_INITIALIZATION_FAILED,
+ /** Adapter 未注册。 */
+ ADAPTER_NOT_FOUND,
+ /** Adapter 不支持当前数据库。 */
+ ADAPTER_UNSUPPORTED,
+ /** SQL 动态参数数量不匹配。 */
+ PARAMETER_COUNT_MISMATCH,
+ /** 查询准入等待超时或被中断。 */
+ QUERY_ADMISSION_TIMEOUT,
+ /** 节点本地联邦中间结果内存准入超时。 */
+ NODE_MEMORY_ADMISSION_TIMEOUT,
+ /** JDBC 连接池获取连接达到超时。 */
+ CONNECTION_ACQUISITION_TIMEOUT,
+ /** JDBC 连接获取因网络、认证或连接池关闭等原因失败。 */
+ CONNECTION_ACQUISITION_FAILED,
+ /** 查询被主动取消。 */
+ QUERY_CANCELLED,
+ /** JDBC 查询或结果读取达到驱动超时。 */
+ QUERY_TIMEOUT,
+ /** JDBC 查询执行失败。 */
+ EXECUTION_FAILED,
+ /** 物理数据库 Explain 执行失败。 */
+ EXPLAIN_FAILED,
+ /** JDBC 或 Runtime 资源关闭失败。 */
+ RESOURCE_CLOSE_FAILED
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlException.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlException.java
new file mode 100644
index 0000000..d4e7af7
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/FederationSqlException.java
@@ -0,0 +1,44 @@
+package com.easyagents.federation.sql.api;
+
+import java.util.Objects;
+
+/**
+ * SQL 联邦查询异常,携带稳定错误码供调用方分类处理。
+ */
+public class FederationSqlException extends RuntimeException {
+
+ /** 稳定错误码。 */
+ private final FederationSqlErrorCode errorCode;
+
+ /**
+ * 创建异常。
+ *
+ * @param errorCode 稳定错误码
+ * @param message 可安全返回的错误说明
+ */
+ public FederationSqlException(FederationSqlErrorCode errorCode, String message) {
+ super(message);
+ this.errorCode = Objects.requireNonNull(errorCode, "errorCode must not be null");
+ }
+
+ /**
+ * 创建带原始原因的异常。
+ *
+ * @param errorCode 稳定错误码
+ * @param message 可安全返回的错误说明
+ * @param cause 原始异常
+ */
+ public FederationSqlException(FederationSqlErrorCode errorCode, String message, Throwable cause) {
+ super(message, cause);
+ this.errorCode = Objects.requireNonNull(errorCode, "errorCode must not be null");
+ }
+
+ /**
+ * 返回稳定错误码。
+ *
+ * @return 错误码
+ */
+ public FederationSqlErrorCode errorCode() {
+ return errorCode;
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionItem.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionItem.java
new file mode 100644
index 0000000..271588d
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionItem.java
@@ -0,0 +1,31 @@
+package com.easyagents.federation.sql.api;
+
+import java.util.List;
+
+/**
+ * 一个可插入 SQL 编辑器的补全候选。
+ *
+ * @param label 面向用户展示的短名称
+ * @param insertText Calcite 生成的替换文本
+ * @param kind 候选类型
+ * @param qualifiedName 候选的完整限定名称
+ */
+public record SqlCompletionItem(
+ String label,
+ String insertText,
+ SqlCompletionKind kind,
+ List qualifiedName
+) {
+
+ /**
+ * 校验并防御性复制候选信息。
+ */
+ public SqlCompletionItem {
+ if (label == null || label.isBlank() || insertText == null || kind == null) {
+ throw new IllegalArgumentException(
+ "completion label, insertText and kind must be provided"
+ );
+ }
+ qualifiedName = List.copyOf(qualifiedName == null ? List.of() : qualifiedName);
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionKind.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionKind.java
new file mode 100644
index 0000000..4640a24
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionKind.java
@@ -0,0 +1,23 @@
+package com.easyagents.federation.sql.api;
+
+/**
+ * SQL 补全候选类型。
+ */
+public enum SqlCompletionKind {
+ /** SQL 关键字。 */
+ KEYWORD,
+ /** SQL 函数。 */
+ FUNCTION,
+ /** 逻辑表。 */
+ TABLE,
+ /** 逻辑视图。 */
+ VIEW,
+ /** Schema。 */
+ SCHEMA,
+ /** Catalog。 */
+ CATALOG,
+ /** 字段。 */
+ COLUMN,
+ /** 无法进一步分类的候选。 */
+ OTHER
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionRequest.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionRequest.java
new file mode 100644
index 0000000..d83a150
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionRequest.java
@@ -0,0 +1,29 @@
+package com.easyagents.federation.sql.api;
+
+import com.easyagents.federation.sql.federation.FederationQueryScopeDefinition;
+
+/**
+ * SQL 编辑器补全请求。
+ *
+ * @param queryScope 当前编辑器可见的查询范围
+ * @param sql 允许不完整的 SQL 文本
+ * @param cursorOffset 光标 UTF-16 字符偏移
+ */
+public record SqlCompletionRequest(
+ FederationQueryScopeDefinition queryScope,
+ String sql,
+ int cursorOffset
+) {
+
+ /**
+ * 校验补全请求。
+ */
+ public SqlCompletionRequest {
+ if (queryScope == null || sql == null) {
+ throw new IllegalArgumentException("queryScope and sql must be provided");
+ }
+ if (cursorOffset < 0 || cursorOffset > sql.length()) {
+ throw new IllegalArgumentException("cursorOffset is outside the SQL text");
+ }
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionResult.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionResult.java
new file mode 100644
index 0000000..5adbc35
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlCompletionResult.java
@@ -0,0 +1,27 @@
+package com.easyagents.federation.sql.api;
+
+import java.util.List;
+
+/**
+ * SQL 补全结果。
+ *
+ * @param replaceStart 建议替换区间起点,使用 UTF-16 字符偏移
+ * @param replaceEnd 建议替换区间终点,使用 UTF-16 字符偏移
+ * @param items 补全候选
+ */
+public record SqlCompletionResult(
+ int replaceStart,
+ int replaceEnd,
+ List items
+) {
+
+ /**
+ * 校验并防御性复制补全结果。
+ */
+ public SqlCompletionResult {
+ if (replaceStart < 0 || replaceEnd < replaceStart) {
+ throw new IllegalArgumentException("completion replacement range is invalid");
+ }
+ items = List.copyOf(items == null ? List.of() : items);
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlExecutionContext.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlExecutionContext.java
new file mode 100644
index 0000000..022370c
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlExecutionContext.java
@@ -0,0 +1,46 @@
+package com.easyagents.federation.sql.api;
+
+import com.easyagents.federation.sql.execute.QueryId;
+import com.easyagents.federation.sql.execute.SqlExecutionOptions;
+import com.easyagents.federation.sql.execute.SqlParameter;
+import java.time.Duration;
+import java.util.List;
+
+/**
+ * 执行节点本地编译计划的上下文。
+ *
+ * @param queryId 查询标识
+ * @param parameters 参数值
+ * @param options JDBC 执行限制
+ * @param admissionTimeout 查询准入等待上限
+ */
+public record SqlExecutionContext(
+ QueryId queryId,
+ List parameters,
+ SqlExecutionOptions options,
+ Duration admissionTimeout
+) {
+
+ /**
+ * 校验并防御性复制执行上下文。
+ */
+ public SqlExecutionContext {
+ queryId = queryId == null ? QueryId.create() : queryId;
+ parameters = List.copyOf(parameters == null ? List.of() : parameters);
+ options = options == null ? SqlExecutionOptions.defaults() : options;
+ admissionTimeout = admissionTimeout == null ? Duration.ofSeconds(5) : admissionTimeout;
+ if (admissionTimeout.isNegative()) {
+ throw new IllegalArgumentException("admissionTimeout must not be negative");
+ }
+ }
+
+ /**
+ * 创建默认执行上下文。
+ *
+ * @param parameters 参数值
+ * @return 执行上下文
+ */
+ public static SqlExecutionContext of(List parameters) {
+ return new SqlExecutionContext(null, parameters, null, null);
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlQueryCommand.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlQueryCommand.java
new file mode 100644
index 0000000..e7049fb
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/api/SqlQueryCommand.java
@@ -0,0 +1,155 @@
+package com.easyagents.federation.sql.api;
+
+import com.easyagents.federation.sql.execute.QueryId;
+import com.easyagents.federation.sql.execute.SqlExecutionOptions;
+import com.easyagents.federation.sql.execute.SqlParameter;
+import com.easyagents.federation.sql.federation.FederationQueryScopeDefinition;
+import com.easyagents.federation.sql.federation.FederationSourceBindingDefinition;
+import com.easyagents.federation.sql.source.SourceId;
+import java.io.Serializable;
+import java.time.Duration;
+import java.util.List;
+
+/**
+ * 可持久化或跨节点传递的一体化查询命令。
+ *
+ * @param queryId 查询标识,可为空并在构造命令时生成
+ * @param sql 单条只读 SQL
+ * @param queryScope 查询可见的数据源范围
+ * @param parameters 参数值
+ * @param options JDBC 执行限制
+ * @param admissionTimeoutMillis 查询准入等待毫秒数
+ * @param policyVersion 策略版本
+ */
+public record SqlQueryCommand(
+ QueryId queryId,
+ String sql,
+ FederationQueryScopeDefinition queryScope,
+ List parameters,
+ SqlExecutionOptions options,
+ long admissionTimeoutMillis,
+ String policyVersion
+) implements Serializable {
+
+ /**
+ * 校验并防御性复制查询命令。
+ */
+ public SqlQueryCommand {
+ if (sql == null || sql.isBlank() || queryScope == null) {
+ throw new IllegalArgumentException("sql and queryScope must be provided");
+ }
+ if (admissionTimeoutMillis < 0) {
+ throw new IllegalArgumentException("timeout must not be negative");
+ }
+ queryId = queryId == null ? QueryId.create() : queryId;
+ parameters = List.copyOf(parameters == null ? List.of() : parameters);
+ options = options == null ? SqlExecutionOptions.defaults() : options;
+ policyVersion = policyVersion == null || policyVersion.isBlank() ? "default" : policyVersion;
+ }
+
+ /**
+ * 使用单物理数据源创建兼容查询命令。
+ *
+ * @param queryId 查询标识
+ * @param sql 单条只读 SQL
+ * @param sourceId 默认数据源
+ * @param minimumRevision 最低数据源版本
+ * @param parameters 参数值
+ * @param options JDBC 执行限制
+ * @param admissionTimeoutMillis 准入等待毫秒数
+ * @param policyVersion 策略版本
+ */
+ public SqlQueryCommand(
+ QueryId queryId,
+ String sql,
+ SourceId sourceId,
+ long minimumRevision,
+ List parameters,
+ SqlExecutionOptions options,
+ long admissionTimeoutMillis,
+ String policyVersion
+ ) {
+ this(
+ queryId,
+ sql,
+ FederationQueryScopeDefinition.single(sourceId, minimumRevision),
+ parameters,
+ options,
+ admissionTimeoutMillis,
+ policyVersion
+ );
+ }
+
+ /**
+ * 返回默认 Binding 的物理数据源,供单源调用方兼容读取。
+ *
+ * @return 默认物理数据源
+ */
+ public SourceId sourceId() {
+ return defaultBinding().sourceId();
+ }
+
+ /**
+ * 返回默认 Binding 的最低物理 Definition 版本。
+ *
+ * @return 最低版本
+ */
+ public long minimumRevision() {
+ return defaultBinding().minimumRevision();
+ }
+
+ private FederationSourceBindingDefinition defaultBinding() {
+ return queryScope.defaultBindingDefinition();
+ }
+
+ /**
+ * 创建使用默认执行限制的查询命令。
+ *
+ * @param sql 单条只读 SQL
+ * @param sourceId 数据源标识
+ * @param minimumRevision 最低数据源版本
+ * @param parameters 参数
+ * @return 查询命令
+ */
+ public static SqlQueryCommand of(
+ String sql,
+ SourceId sourceId,
+ long minimumRevision,
+ List parameters
+ ) {
+ return new SqlQueryCommand(
+ null,
+ sql,
+ sourceId,
+ minimumRevision,
+ parameters,
+ SqlExecutionOptions.defaults(),
+ Duration.ofSeconds(5).toMillis(),
+ "default"
+ );
+ }
+
+ /**
+ * 创建使用默认执行限制的查询范围命令。
+ *
+ * @param sql 单条只读 SQL
+ * @param queryScope 查询范围
+ * @param parameters 参数
+ * @return 查询命令
+ */
+ public static SqlQueryCommand of(
+ String sql,
+ FederationQueryScopeDefinition queryScope,
+ List parameters
+ ) {
+ return new SqlQueryCommand(
+ null,
+ sql,
+ queryScope,
+ parameters,
+ SqlExecutionOptions.defaults(),
+ Duration.ofSeconds(5).toMillis(),
+ "default"
+ );
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationFragmentExplain.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationFragmentExplain.java
new file mode 100644
index 0000000..50f4726
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationFragmentExplain.java
@@ -0,0 +1,93 @@
+package com.easyagents.federation.sql.compile;
+
+import com.easyagents.federation.sql.execute.FederationColumn;
+import com.easyagents.federation.sql.execute.FederationPhysicalExplain;
+import com.easyagents.federation.sql.federation.FederationCostEstimate;
+import com.easyagents.federation.sql.source.SourceId;
+import java.io.Serializable;
+import java.util.List;
+
+/**
+ * Explain 中一个目标数据库分片的纯数据视图。
+ *
+ * @param fragmentId 分片标识
+ * @param bindingName 查询范围 Binding 名称
+ * @param sourceId 物理数据源
+ * @param adapterId Adapter 标识
+ * @param executableSql 目标方言参数化 SQL
+ * @param parameterMapping 分片参数到原查询参数的映射
+ * @param columns 分片输出列
+ * @param costEstimate 分片搬运成本估算
+ * @param pushedDownOperators 已下推算子
+ * @param physicalExplain 显式物理 Explain;逻辑级别时为空
+ */
+public record FederationFragmentExplain(
+ String fragmentId,
+ String bindingName,
+ SourceId sourceId,
+ String adapterId,
+ String executableSql,
+ List parameterMapping,
+ List columns,
+ FederationCostEstimate costEstimate,
+ List pushedDownOperators,
+ FederationPhysicalExplain physicalExplain
+) implements Serializable {
+
+ /**
+ * 防御性复制集合字段。
+ */
+ public FederationFragmentExplain {
+ parameterMapping = List.copyOf(parameterMapping == null ? List.of() : parameterMapping);
+ columns = List.copyOf(columns == null ? List.of() : columns);
+ pushedDownOperators = List.copyOf(
+ pushedDownOperators == null ? List.of() : pushedDownOperators
+ );
+ }
+
+ /**
+ * 创建旧字段集合的兼容 Explain 分片。
+ *
+ * @param fragmentId 分片标识
+ * @param bindingName Binding 名称
+ * @param sourceId 物理源
+ * @param adapterId Adapter 标识
+ * @param executableSql 目标 SQL
+ * @param parameterMapping 参数映射
+ * @param columns 输出列
+ * @param physicalExplain 物理 Explain
+ */
+ public FederationFragmentExplain(
+ String fragmentId,
+ String bindingName,
+ SourceId sourceId,
+ String adapterId,
+ String executableSql,
+ List parameterMapping,
+ List columns,
+ FederationPhysicalExplain physicalExplain
+ ) {
+ this(
+ fragmentId,
+ bindingName,
+ sourceId,
+ adapterId,
+ executableSql,
+ parameterMapping,
+ columns,
+ new FederationCostEstimate(
+ 0,
+ 0,
+ 0,
+ "calcite-default",
+ "none",
+ java.time.Instant.EPOCH,
+ true,
+ com.easyagents.federation.sql.federation.FederationStatisticsStatus.MISSING,
+ false
+ ),
+ List.of(),
+ physicalExplain
+ );
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationSqlPlan.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationSqlPlan.java
new file mode 100644
index 0000000..7d2f6f9
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationSqlPlan.java
@@ -0,0 +1,191 @@
+package com.easyagents.federation.sql.compile;
+
+import com.easyagents.federation.sql.adapter.AdapterCompatibility;
+import com.easyagents.federation.sql.execute.FederationColumn;
+import com.easyagents.federation.sql.federation.FederationFragmentPlan;
+import com.easyagents.federation.sql.federation.FederationJoinOptimization;
+import com.easyagents.federation.sql.federation.FederationQueryMode;
+import com.easyagents.federation.sql.federation.FederationQueryScopeDefinition;
+import com.easyagents.federation.sql.federation.FederationSourceRuntimeIdentity;
+import com.easyagents.federation.sql.source.SourceId;
+import java.time.Instant;
+import java.util.List;
+import java.util.Set;
+import org.apache.calcite.rel.RelRoot;
+import org.apache.calcite.sql.SqlNode;
+
+/**
+ * Engine 签发的节点本地 SQL 编译计划。
+ *
+ * 该接口仅用于读取编译事实。调用方不能自行创建可执行计划,且计划只能交回
+ * 签发它的 Engine 实例执行。
+ */
+public interface FederationSqlPlan {
+
+ /**
+ * 返回主数据源。
+ *
+ * @return 主数据源
+ */
+ SourceId sourceId();
+
+ /**
+ * 返回编译时使用的不可变查询范围。
+ *
+ * @return 查询范围
+ */
+ FederationQueryScopeDefinition queryScope();
+
+ /**
+ * 返回根据实际引用源确定的查询模式。
+ *
+ * @return 查询模式
+ */
+ FederationQueryMode queryMode();
+
+ /**
+ * 返回数据源版本。
+ *
+ * @return 数据源版本
+ */
+ long sourceRevision();
+
+ /**
+ * 返回 Calcite 规范化 SQL。
+ *
+ * @return Calcite 规范化 SQL
+ */
+ String normalizedSql();
+
+ /**
+ * 返回单源目标数据库参数化 SQL。
+ *
+ * 联邦计划应读取 {@link #fragments()};本兼容视图不代表任一数据库可执行 SQL。
+ *
+ * @return 单源目标 SQL,或联邦调用方原始 SQL 兼容视图
+ */
+ String executableSql();
+
+ /**
+ * 返回 Calcite 已校验 SQL 节点。
+ *
+ * @return Calcite 已校验 SQL 节点
+ */
+ SqlNode sqlNode();
+
+ /**
+ * 返回 Calcite 关系计划。
+ *
+ * @return Calcite 关系计划
+ */
+ RelRoot relRoot();
+
+ /**
+ * 返回动态参数数量。
+ *
+ * @return 动态参数数量
+ */
+ int parameterCount();
+
+ /**
+ * 返回编译时声明的原始 JDBC 参数类型。
+ *
+ * @return JDBC 参数类型;未显式声明时为空
+ */
+ List parameterJdbcTypes();
+
+ /**
+ * 返回目标 SQL 占位符到原始参数的零基索引映射。
+ *
+ * @return 参数映射
+ */
+ List parameterMapping();
+
+ /**
+ * 返回物理数据源分片;单源计划也包含一个分片。
+ *
+ * @return 分片列表
+ */
+ List fragments();
+
+ /**
+ * 返回跨源 Join 的优化选择。
+ *
+ * @return 不可变 Join 优化列表
+ */
+ default List joinOptimizations() {
+ return List.of();
+ }
+
+ /**
+ * 返回实际引用 Binding 对应的节点本地运行身份。
+ *
+ * @return 运行身份列表
+ */
+ List sourceRuntimeIdentities();
+
+ /**
+ * 返回查询范围的稳定校验和。
+ *
+ * @return Scope 校验和
+ */
+ String scopeChecksum();
+
+ /**
+ * 返回结果列。
+ *
+ * @return 结果列
+ */
+ List columns();
+
+ /**
+ * 返回引用的数据源集合。
+ *
+ * @return 引用的数据源集合
+ */
+ Set referencedSources();
+
+ /**
+ * 返回 Adapter 兼容性。
+ *
+ * @return Adapter 兼容性
+ */
+ AdapterCompatibility compatibility();
+
+ /**
+ * 返回是否允许直接执行。
+ *
+ * @return 是否允许直接执行
+ */
+ boolean executable();
+
+ /**
+ * 返回编译时的数据源 Definition 校验和。
+ *
+ * @return Definition 校验和
+ */
+ String sourceChecksum();
+
+ /**
+ * 返回编译时的 Adapter 标识。
+ *
+ * @return Adapter 标识
+ */
+ String adapterId();
+
+ /**
+ * 返回编译时的数据库与驱动指纹。
+ *
+ * @return 运行指纹摘要
+ */
+ String runtimeFingerprint();
+
+ /**
+ * 返回该计划所依赖统计快照的最早失效时间。
+ *
+ * @return 最早失效时间;未使用有期限统计时为 {@link Instant#MAX}
+ */
+ default Instant statisticsValidUntil() {
+ return Instant.MAX;
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationSqlPolicy.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationSqlPolicy.java
new file mode 100644
index 0000000..f1845bd
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/FederationSqlPolicy.java
@@ -0,0 +1,27 @@
+package com.easyagents.federation.sql.compile;
+
+/**
+ * 调用方在 SQL 已校验并转换为 RelRoot 后执行的策略 SPI。
+ */
+@FunctionalInterface
+public interface FederationSqlPolicy {
+
+ /**
+ * 返回策略实现的稳定版本,用于隔离计划缓存。
+ *
+ * 策略规则发生变化时应同步更新版本。默认版本适用于 Engine 生命周期内
+ * 逻辑不变的无状态策略。
+ *
+ * @return 稳定策略版本
+ */
+ default String version() {
+ return "1";
+ }
+
+ /**
+ * 校验已编译 SQL;拒绝时应抛出 FederationSqlException。
+ *
+ * @param context 策略上下文
+ */
+ void validate(SqlPolicyContext context);
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlCompileRequest.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlCompileRequest.java
new file mode 100644
index 0000000..fdee17e
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlCompileRequest.java
@@ -0,0 +1,108 @@
+package com.easyagents.federation.sql.compile;
+
+import com.easyagents.federation.sql.federation.FederationQueryScopeDefinition;
+import com.easyagents.federation.sql.federation.FederationSourceBindingDefinition;
+import com.easyagents.federation.sql.source.SourceId;
+import java.util.List;
+
+/**
+ * 节点本地 SQL 编译请求。
+ *
+ * @param sql 单条只读 SQL
+ * @param queryScope 查询可见的数据源范围
+ * @param parameterJdbcTypes 参数 JDBC 类型列表
+ * @param policyVersion 调用方策略版本,用于隔离计划缓存
+ */
+public record SqlCompileRequest(
+ String sql,
+ FederationQueryScopeDefinition queryScope,
+ List parameterJdbcTypes,
+ String policyVersion
+) {
+
+ /**
+ * 校验并防御性复制编译请求。
+ */
+ public SqlCompileRequest {
+ if (sql == null || sql.isBlank()) {
+ throw new IllegalArgumentException("sql must not be blank");
+ }
+ if (queryScope == null) {
+ throw new IllegalArgumentException("queryScope must not be null");
+ }
+ parameterJdbcTypes = List.copyOf(parameterJdbcTypes == null ? List.of() : parameterJdbcTypes);
+ policyVersion = policyVersion == null || policyVersion.isBlank() ? "default" : policyVersion;
+ }
+
+ /**
+ * 使用单物理数据源创建兼容编译请求。
+ *
+ * @param sql 单条只读 SQL
+ * @param sourceId 默认数据源
+ * @param minimumRevision 最低数据源版本
+ * @param parameterJdbcTypes 参数 JDBC 类型
+ * @param policyVersion 调用方策略版本
+ */
+ public SqlCompileRequest(
+ String sql,
+ SourceId sourceId,
+ long minimumRevision,
+ List parameterJdbcTypes,
+ String policyVersion
+ ) {
+ this(
+ sql,
+ FederationQueryScopeDefinition.single(sourceId, minimumRevision),
+ parameterJdbcTypes,
+ policyVersion
+ );
+ }
+
+ /**
+ * 返回默认 Binding 的物理数据源,供单源调用方兼容读取。
+ *
+ * @return 默认物理数据源
+ */
+ public SourceId sourceId() {
+ return defaultBinding().sourceId();
+ }
+
+ /**
+ * 返回默认 Binding 的最低物理 Definition 版本。
+ *
+ * @return 最低版本
+ */
+ public long minimumRevision() {
+ return defaultBinding().minimumRevision();
+ }
+
+ private FederationSourceBindingDefinition defaultBinding() {
+ return queryScope.defaultBindingDefinition();
+ }
+
+ /**
+ * 创建无参数的默认编译请求。
+ *
+ * @param sql 单条只读 SQL
+ * @param sourceId 默认数据源
+ * @param minimumRevision 最低数据源版本
+ * @return 编译请求
+ */
+ public static SqlCompileRequest of(String sql, SourceId sourceId, long minimumRevision) {
+ return new SqlCompileRequest(sql, sourceId, minimumRevision, List.of(), "default");
+ }
+
+ /**
+ * 创建无参数的查询范围编译请求。
+ *
+ * @param sql 单条只读 SQL
+ * @param queryScope 查询范围
+ * @return 编译请求
+ */
+ public static SqlCompileRequest of(
+ String sql,
+ FederationQueryScopeDefinition queryScope
+ ) {
+ return new SqlCompileRequest(sql, queryScope, List.of(), "default");
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainLevel.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainLevel.java
new file mode 100644
index 0000000..4d331b4
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainLevel.java
@@ -0,0 +1,13 @@
+package com.easyagents.federation.sql.compile;
+
+/**
+ * Explain 深度。
+ */
+public enum SqlExplainLevel {
+
+ /** 只生成 Calcite 逻辑计划和物理分片 SQL,不访问数据库 Optimizer。 */
+ LOGICAL,
+
+ /** 在逻辑计划基础上显式请求各物理数据库的非 ANALYZE Explain。 */
+ PHYSICAL
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainRequest.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainRequest.java
new file mode 100644
index 0000000..93fe1fd
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainRequest.java
@@ -0,0 +1,54 @@
+package com.easyagents.federation.sql.compile;
+
+import com.easyagents.federation.sql.execute.SqlParameter;
+import java.util.List;
+
+/**
+ * SQL Explain 请求。
+ *
+ * @param compileRequest 编译请求
+ * @param level Explain 深度
+ * @param parameters 保留的兼容字段;为避免数据库计划回显敏感值,只允许为空
+ */
+public record SqlExplainRequest(
+ SqlCompileRequest compileRequest,
+ SqlExplainLevel level,
+ List parameters
+) {
+
+ /**
+ * 校验 Explain 请求。
+ */
+ public SqlExplainRequest {
+ if (compileRequest == null) {
+ throw new IllegalArgumentException("compileRequest must not be null");
+ }
+ level = level == null ? SqlExplainLevel.PHYSICAL : level;
+ parameters = List.copyOf(parameters == null ? List.of() : parameters);
+ if (!parameters.isEmpty()) {
+ throw new IllegalArgumentException(
+ "physical Explain does not accept parameter values; "
+ + "declare JDBC types in compileRequest"
+ );
+ }
+ }
+
+ /**
+ * 创建默认物理 Explain 请求,不提供实际参数值。
+ *
+ * @param compileRequest 编译请求
+ */
+ public SqlExplainRequest(SqlCompileRequest compileRequest) {
+ this(compileRequest, SqlExplainLevel.PHYSICAL, List.of());
+ }
+
+ /**
+ * 创建指定深度的 Explain 请求。
+ *
+ * @param compileRequest 编译请求
+ * @param level Explain 深度
+ */
+ public SqlExplainRequest(SqlCompileRequest compileRequest, SqlExplainLevel level) {
+ this(compileRequest, level, List.of());
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainResult.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainResult.java
new file mode 100644
index 0000000..07b052b
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlExplainResult.java
@@ -0,0 +1,165 @@
+package com.easyagents.federation.sql.compile;
+
+import com.easyagents.federation.sql.adapter.AdapterCompatibility;
+import com.easyagents.federation.sql.federation.FederationQueryMode;
+import com.easyagents.federation.sql.federation.FederationJoinOptimization;
+import com.easyagents.federation.sql.federation.FederationStatisticsStatus;
+import com.easyagents.federation.sql.source.SourceId;
+import java.io.Serializable;
+import java.util.List;
+import java.util.Set;
+
+/**
+ * 不含节点本地 Calcite/JDBC 对象的 Explain 结果。
+ *
+ * @param level Explain 深度
+ * @param queryMode 实际查询模式
+ * @param statisticsStatus 成本统计完整性与时效状态
+ * @param estimateAvailable 聚合成本数值是否为有效估算
+ * @param estimatedTransferBytes 预计从物理源搬运的总字节数
+ * @param estimatedLocalMemoryBytes 本地 Join 构建侧估算内存字节数
+ * @param joinOptimizations 跨源 Join 优化选择
+ * @param normalizedSql Calcite 规范化 SQL
+ * @param executableSql 单源目标方言 SQL;联邦计划仅为兼容视图,应读取 fragments
+ * @param relationalPlan 关系计划文本
+ * @param executionPlan 实际单源下推或联邦本地执行计划文本
+ * @param fragments 物理分片与可选数据库计划
+ * @param referencedSources 引用的数据源
+ * @param compatibility Adapter 兼容性
+ * @param executable 是否允许执行
+ * @param planCacheHit 是否命中节点本地计划缓存
+ * @param diagnostic 诊断说明
+ */
+public record SqlExplainResult(
+ SqlExplainLevel level,
+ FederationQueryMode queryMode,
+ FederationStatisticsStatus statisticsStatus,
+ boolean estimateAvailable,
+ double estimatedTransferBytes,
+ long estimatedLocalMemoryBytes,
+ List joinOptimizations,
+ String normalizedSql,
+ String executableSql,
+ String relationalPlan,
+ String executionPlan,
+ List fragments,
+ Set referencedSources,
+ AdapterCompatibility compatibility,
+ boolean executable,
+ boolean planCacheHit,
+ String diagnostic
+) implements Serializable {
+
+ /**
+ * 防御性复制引用集合。
+ */
+ public SqlExplainResult {
+ level = level == null ? SqlExplainLevel.LOGICAL : level;
+ queryMode = queryMode == null ? FederationQueryMode.SINGLE_SOURCE : queryMode;
+ statisticsStatus = statisticsStatus == null
+ ? FederationStatisticsStatus.MISSING
+ : statisticsStatus;
+ if (!Double.isFinite(estimatedTransferBytes) || estimatedTransferBytes < 0
+ || estimatedLocalMemoryBytes < 0) {
+ throw new IllegalArgumentException("Explain cost values must be non-negative");
+ }
+ joinOptimizations = List.copyOf(
+ joinOptimizations == null ? List.of() : joinOptimizations
+ );
+ fragments = List.copyOf(fragments == null ? List.of() : fragments);
+ referencedSources = Set.copyOf(referencedSources);
+ diagnostic = diagnostic == null ? "" : diagnostic;
+ }
+
+ /**
+ * 创建未包含聚合成本字段的兼容 Explain 结果。
+ *
+ * @param level Explain 深度
+ * @param queryMode 查询模式
+ * @param normalizedSql 规范化 SQL
+ * @param executableSql 可执行 SQL
+ * @param relationalPlan 关系计划
+ * @param executionPlan 执行计划
+ * @param fragments 分片计划
+ * @param referencedSources 引用源
+ * @param compatibility Adapter 兼容性
+ * @param executable 是否可执行
+ * @param planCacheHit 是否命中缓存
+ * @param diagnostic 诊断信息
+ */
+ public SqlExplainResult(
+ SqlExplainLevel level,
+ FederationQueryMode queryMode,
+ String normalizedSql,
+ String executableSql,
+ String relationalPlan,
+ String executionPlan,
+ List fragments,
+ Set referencedSources,
+ AdapterCompatibility compatibility,
+ boolean executable,
+ boolean planCacheHit,
+ String diagnostic
+ ) {
+ this(
+ level,
+ queryMode,
+ FederationStatisticsStatus.MISSING,
+ false,
+ 0D,
+ 0L,
+ List.of(),
+ normalizedSql,
+ executableSql,
+ relationalPlan,
+ executionPlan,
+ fragments,
+ referencedSources,
+ compatibility,
+ executable,
+ planCacheHit,
+ diagnostic
+ );
+ }
+
+ /**
+ * 创建旧单源字段视图的兼容 Explain 结果。
+ *
+ * @param normalizedSql Calcite 规范化 SQL
+ * @param executableSql 目标方言 SQL
+ * @param relationalPlan 关系计划文本
+ * @param referencedSources 引用的数据源
+ * @param compatibility Adapter 兼容性
+ * @param executable 是否允许执行
+ * @param diagnostic 诊断说明
+ */
+ public SqlExplainResult(
+ String normalizedSql,
+ String executableSql,
+ String relationalPlan,
+ Set referencedSources,
+ AdapterCompatibility compatibility,
+ boolean executable,
+ String diagnostic
+ ) {
+ this(
+ SqlExplainLevel.LOGICAL,
+ FederationQueryMode.SINGLE_SOURCE,
+ FederationStatisticsStatus.MISSING,
+ false,
+ 0D,
+ 0L,
+ List.of(),
+ normalizedSql,
+ executableSql,
+ relationalPlan,
+ relationalPlan,
+ List.of(),
+ referencedSources,
+ compatibility,
+ executable,
+ false,
+ diagnostic
+ );
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlPolicyContext.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlPolicyContext.java
new file mode 100644
index 0000000..d63a162
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/compile/SqlPolicyContext.java
@@ -0,0 +1,22 @@
+package com.easyagents.federation.sql.compile;
+
+import com.easyagents.federation.sql.source.SourceId;
+import java.util.Set;
+import org.apache.calcite.rel.RelRoot;
+import org.apache.calcite.sql.SqlNode;
+
+/**
+ * SQL 策略直接读取 Calcite 事实对象的上下文。
+ *
+ * @param request 原始编译请求
+ * @param validatedSql 已校验 SqlNode
+ * @param relRoot 关系计划
+ * @param referencedSources 引用的数据源
+ */
+public record SqlPolicyContext(
+ SqlCompileRequest request,
+ SqlNode validatedSql,
+ RelRoot relRoot,
+ Set referencedSources
+) {
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationColumn.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationColumn.java
new file mode 100644
index 0000000..1f878a6
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationColumn.java
@@ -0,0 +1,21 @@
+package com.easyagents.federation.sql.execute;
+
+import java.io.Serializable;
+
+/**
+ * 查询结果列元数据。
+ *
+ * @param index 从 1 开始的列序号
+ * @param label 列标签
+ * @param jdbcType JDBC 类型
+ * @param typeName 数据库类型名
+ * @param nullable 是否允许空值
+ */
+public record FederationColumn(
+ int index,
+ String label,
+ int jdbcType,
+ String typeName,
+ boolean nullable
+) implements Serializable {
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationExecutionGuard.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationExecutionGuard.java
new file mode 100644
index 0000000..8289dd5
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationExecutionGuard.java
@@ -0,0 +1,67 @@
+package com.easyagents.federation.sql.execute;
+
+import java.util.concurrent.TimeUnit;
+
+/**
+ * 贯穿连接获取、Statement 执行和结果读取的查询终态检查器。
+ */
+public interface FederationExecutionGuard {
+
+ /**
+ * 检查查询是否仍允许继续执行。
+ *
+ * @throws RuntimeException 查询取消或超时时抛出稳定异常
+ */
+ void ensureAllowed();
+
+ /**
+ * 返回查询剩余时限。
+ *
+ * @return 剩余纳秒数;无限制时返回 {@link Long#MAX_VALUE}
+ */
+ long remainingNanos();
+
+ /**
+ * 将调用方 JDBC 秒级超时收敛到统一剩余时限。
+ *
+ * @param requestedSeconds 调用方超时,0 表示未指定
+ * @return 至少 1 秒的 JDBC 超时;无限制且未指定时返回 0
+ */
+ default int boundedQueryTimeoutSeconds(int requestedSeconds) {
+ if (remainingNanos() == Long.MAX_VALUE) {
+ return requestedSeconds;
+ }
+ long remainingSeconds = Math.max(
+ 1L,
+ TimeUnit.NANOSECONDS.toSeconds(Math.max(1L, remainingNanos()))
+ );
+ int bounded = (int) Math.min(Integer.MAX_VALUE, remainingSeconds);
+ return requestedSeconds == 0 ? bounded : Math.min(requestedSeconds, bounded);
+ }
+
+ /**
+ * 返回无限制检查器,供旧 Adapter 调用兼容使用。
+ *
+ * @return 无限制检查器
+ */
+ static FederationExecutionGuard none() {
+ return NoopHolder.INSTANCE;
+ }
+
+ /** 无状态实例持有者。 */
+ final class NoopHolder {
+ private static final FederationExecutionGuard INSTANCE = new FederationExecutionGuard() {
+ @Override
+ public void ensureAllowed() {
+ }
+
+ @Override
+ public long remainingNanos() {
+ return Long.MAX_VALUE;
+ }
+ };
+
+ private NoopHolder() {
+ }
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationExecutionObserver.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationExecutionObserver.java
new file mode 100644
index 0000000..9f37b74
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationExecutionObserver.java
@@ -0,0 +1,41 @@
+package com.easyagents.federation.sql.execute;
+
+/**
+ * Adapter 向 Core 回传 Fragment 执行阶段耗时的轻量观察器。
+ */
+public interface FederationExecutionObserver {
+
+ /**
+ * 记录获取物理连接的耗时。
+ *
+ * @param elapsedNanos 获取连接耗时
+ */
+ default void connectionAcquired(long elapsedNanos) {
+ }
+
+ /**
+ * 记录数据库完成 Statement 执行并返回 ResultSet 的耗时。
+ *
+ * @param elapsedNanos 数据库执行耗时
+ */
+ default void databaseExecutionCompleted(long elapsedNanos) {
+ }
+
+ /**
+ * 记录 ResultSet 返回首行的耗时。
+ *
+ * @param elapsedNanos 从 ResultSet 创建到首行可用的耗时
+ */
+ default void firstRowAvailable(long elapsedNanos) {
+ }
+
+ /**
+ * 返回不采集指标的观察器。
+ *
+ * @return 空观察器
+ */
+ static FederationExecutionObserver none() {
+ return new FederationExecutionObserver() {
+ };
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExecutionContext.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExecutionContext.java
new file mode 100644
index 0000000..b0e456c
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExecutionContext.java
@@ -0,0 +1,122 @@
+package com.easyagents.federation.sql.execute;
+
+import com.easyagents.federation.sql.adapter.AdapterCompatibility;
+import java.util.List;
+import java.util.Map;
+import javax.sql.DataSource;
+
+/**
+ * 单数据源 SQL Fragment 的执行上下文。
+ *
+ * @param queryId 查询标识
+ * @param sql 已按目标方言生成的参数化 SQL
+ * @param parameters JDBC 参数
+ * @param options 强制执行限制
+ * @param dataSource 调用方提供的 DataSource
+ * @param compatibility 当前数据库与驱动兼容性信息
+ * @param adapterOptions 不含凭据的 Adapter 执行选项
+ * @param statementLifecycle Statement 取消登记回调
+ * @param observer Fragment 执行阶段观察器
+ * @param executionGuard 查询取消与统一截止时间检查器
+ */
+public record FederationFragmentExecutionContext(
+ QueryId queryId,
+ String sql,
+ List parameters,
+ SqlExecutionOptions options,
+ DataSource dataSource,
+ AdapterCompatibility compatibility,
+ Map adapterOptions,
+ StatementLifecycle statementLifecycle,
+ FederationExecutionObserver observer,
+ FederationExecutionGuard executionGuard
+) {
+
+ /**
+ * 防御性复制参数并校验必需字段。
+ */
+ public FederationFragmentExecutionContext {
+ if (queryId == null || sql == null || sql.isBlank() || options == null
+ || dataSource == null || compatibility == null || statementLifecycle == null) {
+ throw new IllegalArgumentException("fragment execution context contains null or blank values");
+ }
+ parameters = List.copyOf(parameters == null ? List.of() : parameters);
+ adapterOptions = Map.copyOf(adapterOptions == null ? Map.of() : adapterOptions);
+ observer = observer == null ? FederationExecutionObserver.none() : observer;
+ executionGuard = executionGuard == null ? FederationExecutionGuard.none() : executionGuard;
+ }
+
+ /**
+ * 创建不采集 Adapter 阶段指标的兼容执行上下文。
+ *
+ * @param queryId 查询标识
+ * @param sql 参数化 SQL
+ * @param parameters JDBC 参数
+ * @param options 执行限制
+ * @param dataSource 数据源
+ * @param compatibility 数据库兼容信息
+ * @param adapterOptions Adapter 选项
+ * @param statementLifecycle Statement 生命周期
+ */
+ public FederationFragmentExecutionContext(
+ QueryId queryId,
+ String sql,
+ List parameters,
+ SqlExecutionOptions options,
+ DataSource dataSource,
+ AdapterCompatibility compatibility,
+ Map adapterOptions,
+ StatementLifecycle statementLifecycle
+ ) {
+ this(
+ queryId,
+ sql,
+ parameters,
+ options,
+ dataSource,
+ compatibility,
+ adapterOptions,
+ statementLifecycle,
+ FederationExecutionObserver.none(),
+ FederationExecutionGuard.none()
+ );
+ }
+
+ /**
+ * 创建带阶段观察器的旧调用兼容上下文。
+ *
+ * @param queryId 查询标识
+ * @param sql 参数化 SQL
+ * @param parameters JDBC 参数
+ * @param options 执行限制
+ * @param dataSource 数据源
+ * @param compatibility 数据库兼容信息
+ * @param adapterOptions Adapter 选项
+ * @param statementLifecycle Statement 生命周期
+ * @param observer Fragment 观察器
+ */
+ public FederationFragmentExecutionContext(
+ QueryId queryId,
+ String sql,
+ List parameters,
+ SqlExecutionOptions options,
+ DataSource dataSource,
+ AdapterCompatibility compatibility,
+ Map adapterOptions,
+ StatementLifecycle statementLifecycle,
+ FederationExecutionObserver observer
+ ) {
+ this(
+ queryId,
+ sql,
+ parameters,
+ options,
+ dataSource,
+ compatibility,
+ adapterOptions,
+ statementLifecycle,
+ observer,
+ FederationExecutionGuard.none()
+ );
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExecutor.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExecutor.java
new file mode 100644
index 0000000..5d0fb8a
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExecutor.java
@@ -0,0 +1,16 @@
+package com.easyagents.federation.sql.execute;
+
+/**
+ * Adapter 提供的单数据源 SQL Fragment 执行器。
+ */
+@FunctionalInterface
+public interface FederationFragmentExecutor {
+
+ /**
+ * 执行参数化 SQL 并返回流式游标。
+ *
+ * @param context 执行上下文
+ * @return 流式游标
+ */
+ FederationResultCursor execute(FederationFragmentExecutionContext context);
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExplainContext.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExplainContext.java
new file mode 100644
index 0000000..50b0492
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExplainContext.java
@@ -0,0 +1,74 @@
+package com.easyagents.federation.sql.execute;
+
+import com.easyagents.federation.sql.adapter.AdapterCompatibility;
+import java.util.List;
+import java.util.Map;
+import javax.sql.DataSource;
+
+/**
+ * Adapter 执行单个物理分片 Explain 的上下文。
+ *
+ * @param sql 目标数据库方言参数化 SQL
+ * @param parameters 已按分片参数映射排序的参数
+ * @param dataSource 调用方提供的 DataSource
+ * @param compatibility 数据库与驱动兼容性
+ * @param adapterOptions 不含凭据的 Adapter 选项
+ * @param queryTimeoutSeconds Explain 超时秒数
+ * @param executionGuard 统一查询终态与截止时间检查器
+ */
+public record FederationFragmentExplainContext(
+ String sql,
+ List parameters,
+ DataSource dataSource,
+ AdapterCompatibility compatibility,
+ Map adapterOptions,
+ int queryTimeoutSeconds,
+ FederationExecutionGuard executionGuard
+) {
+
+ /**
+ * 校验并创建不可变 Explain 上下文。
+ */
+ public FederationFragmentExplainContext {
+ if (sql == null || sql.isBlank() || dataSource == null || compatibility == null) {
+ throw new IllegalArgumentException("fragment Explain context is incomplete");
+ }
+ if (queryTimeoutSeconds < 0) {
+ throw new IllegalArgumentException("queryTimeoutSeconds must not be negative");
+ }
+ parameters = List.copyOf(parameters == null ? List.of() : parameters);
+ adapterOptions = Map.copyOf(adapterOptions == null ? Map.of() : adapterOptions);
+ executionGuard = executionGuard == null
+ ? FederationExecutionGuard.none()
+ : executionGuard;
+ }
+
+ /**
+ * 保留旧 Adapter 与调用方的兼容构造器。
+ *
+ * @param sql 目标数据库 SQL
+ * @param parameters 参数
+ * @param dataSource 数据源
+ * @param compatibility 兼容性信息
+ * @param adapterOptions Adapter 选项
+ * @param queryTimeoutSeconds Explain 超时秒数
+ */
+ public FederationFragmentExplainContext(
+ String sql,
+ List parameters,
+ DataSource dataSource,
+ AdapterCompatibility compatibility,
+ Map adapterOptions,
+ int queryTimeoutSeconds
+ ) {
+ this(
+ sql,
+ parameters,
+ dataSource,
+ compatibility,
+ adapterOptions,
+ queryTimeoutSeconds,
+ FederationExecutionGuard.none()
+ );
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExplainer.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExplainer.java
new file mode 100644
index 0000000..2c8123d
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentExplainer.java
@@ -0,0 +1,16 @@
+package com.easyagents.federation.sql.execute;
+
+/**
+ * Adapter 可选的物理数据库 Explain SPI。
+ */
+@FunctionalInterface
+public interface FederationFragmentExplainer {
+
+ /**
+ * 执行不会运行真实数据查询的物理 Explain。
+ *
+ * @param context 分片 Explain 上下文
+ * @return 物理计划
+ */
+ FederationPhysicalExplain explain(FederationFragmentExplainContext context);
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentMetrics.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentMetrics.java
new file mode 100644
index 0000000..a56da5f
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationFragmentMetrics.java
@@ -0,0 +1,51 @@
+package com.easyagents.federation.sql.execute;
+
+import com.easyagents.federation.sql.source.SourceId;
+import java.io.Serializable;
+
+/**
+ * 单个物理分片的查询消耗快照。
+ *
+ * @param fragmentId 分片标识
+ * @param sourceId 物理数据源标识
+ * @param rowsRead 已读取行数
+ * @param bytesRead 已读取估算字节数;Adapter 未安全提供时为 -1
+ * @param elapsedNanos 当前或最终耗时
+ * @param connectionAcquireNanos 获取连接耗时
+ * @param databaseExecutionNanos Statement 返回 ResultSet 的耗时
+ * @param firstRowNanos ResultSet 创建到首行可用的耗时;不可用时为 -1
+ * @param complete 是否已完成或关闭
+ */
+public record FederationFragmentMetrics(
+ String fragmentId,
+ SourceId sourceId,
+ long rowsRead,
+ long bytesRead,
+ long elapsedNanos,
+ long connectionAcquireNanos,
+ long databaseExecutionNanos,
+ long firstRowNanos,
+ boolean complete
+) implements Serializable {
+
+ /**
+ * 创建只包含旧基础字段的兼容分片指标。
+ *
+ * @param fragmentId 分片标识
+ * @param sourceId 数据源
+ * @param rowsRead 读取行数
+ * @param bytesRead 读取字节数
+ * @param elapsedNanos 耗时
+ * @param complete 是否完成
+ */
+ public FederationFragmentMetrics(
+ String fragmentId,
+ SourceId sourceId,
+ long rowsRead,
+ long bytesRead,
+ long elapsedNanos,
+ boolean complete
+ ) {
+ this(fragmentId, sourceId, rowsRead, bytesRead, elapsedNanos, 0, 0, -1, complete);
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationLocalOperatorMetrics.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationLocalOperatorMetrics.java
new file mode 100644
index 0000000..6421478
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationLocalOperatorMetrics.java
@@ -0,0 +1,19 @@
+package com.easyagents.federation.sql.execute;
+
+import java.io.Serializable;
+
+/**
+ * Calcite 本地联邦算子的累计执行指标。
+ *
+ * @param operatorName 算子名称
+ * @param outputRows 交给下游的输出行数
+ * @param outputBytes 输出估算字节数
+ * @param executionNanos 算子产生输出的累计耗时
+ */
+public record FederationLocalOperatorMetrics(
+ String operatorName,
+ long outputRows,
+ long outputBytes,
+ long executionNanos
+) implements Serializable {
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationPhysicalExplain.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationPhysicalExplain.java
new file mode 100644
index 0000000..1fc9da7
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationPhysicalExplain.java
@@ -0,0 +1,59 @@
+package com.easyagents.federation.sql.execute;
+
+import java.io.Serializable;
+import java.util.List;
+
+/**
+ * 物理数据库 Optimizer 的非 ANALYZE Explain 结果。
+ *
+ * @param available 是否获得物理计划
+ * @param nativePlan 数据库原生计划文本
+ * @param nodeType 首个主要计划节点类型
+ * @param scanType 扫描或访问方式
+ * @param candidateIndexes 数据库返回的候选索引
+ * @param chosenIndex 数据库选择的索引
+ * @param estimatedRows 数据库估算行数
+ * @param extraCondition 额外过滤或索引条件
+ * @param diagnostic 不含凭据和参数值的诊断
+ */
+public record FederationPhysicalExplain(
+ boolean available,
+ String nativePlan,
+ String nodeType,
+ String scanType,
+ List candidateIndexes,
+ String chosenIndex,
+ Long estimatedRows,
+ String extraCondition,
+ String diagnostic
+) implements Serializable {
+
+ /**
+ * 防御性复制候选索引。
+ */
+ public FederationPhysicalExplain {
+ candidateIndexes = List.copyOf(candidateIndexes == null ? List.of() : candidateIndexes);
+ nativePlan = nativePlan == null ? "" : nativePlan;
+ diagnostic = diagnostic == null ? "" : diagnostic;
+ }
+
+ /**
+ * 创建数据库不支持或未能提供物理计划的结果。
+ *
+ * @param diagnostic 诊断说明
+ * @return 不可用结果
+ */
+ public static FederationPhysicalExplain unavailable(String diagnostic) {
+ return new FederationPhysicalExplain(
+ false,
+ "",
+ null,
+ null,
+ List.of(),
+ null,
+ null,
+ null,
+ diagnostic
+ );
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryAdmissionController.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryAdmissionController.java
new file mode 100644
index 0000000..90d3139
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryAdmissionController.java
@@ -0,0 +1,98 @@
+package com.easyagents.federation.sql.execute;
+
+import com.easyagents.federation.sql.api.FederationSqlErrorCode;
+import com.easyagents.federation.sql.api.FederationSqlException;
+import com.easyagents.federation.sql.source.SourceId;
+import java.time.Duration;
+import java.util.function.BooleanSupplier;
+
+/**
+ * 可替换的查询并发准入控制器。
+ */
+@FunctionalInterface
+public interface FederationQueryAdmissionController extends AutoCloseable {
+
+ /**
+ * 获取查询许可。
+ *
+ * @param sourceId 数据源标识
+ * @param queryId 查询标识
+ * @param timeout 最大等待时间
+ * @return 查询许可
+ */
+ FederationQueryPermit acquire(SourceId sourceId, QueryId queryId, Duration timeout);
+
+ /**
+ * 获取支持查询级取消的许可。
+ *
+ * 自定义实现可以覆盖此方法及时中断分布式或远程准入等待。
+ *
+ * @param sourceId 数据源标识
+ * @param queryId 查询标识
+ * @param timeout 最大等待时间
+ * @param cancellationRequested 取消状态
+ * @return 查询许可
+ */
+ default FederationQueryPermit acquire(
+ SourceId sourceId,
+ QueryId queryId,
+ Duration timeout,
+ BooleanSupplier cancellationRequested
+ ) {
+ return acquire(sourceId, queryId, timeout);
+ }
+
+ /**
+ * 为一次单源或联邦查询获取一份查询级许可。
+ *
+ * 兼容实现只接受单源请求。联邦查询必须由实现方明确覆盖本方法,避免其余
+ * 物理源静默绕过源级配额。
+ *
+ * @param request 查询级准入请求
+ * @return 查询许可
+ */
+ default FederationQueryPermit acquire(QueryAdmissionRequest request) {
+ if (request.sourceIds().size() != 1) {
+ throw new FederationSqlException(
+ FederationSqlErrorCode.INVALID_QUERY_SCOPE,
+ "admission controller does not declare multi-source query support"
+ );
+ }
+ return acquire(
+ request.primarySourceId(),
+ request.queryId(),
+ request.timeout(),
+ request.cancellationRequested()
+ );
+ }
+
+ /**
+ * 关闭控制器;默认无额外资源。
+ */
+ @Override
+ default void close() {
+ }
+
+ /**
+ * 返回无并发限制的控制器。
+ *
+ * @return 无限制控制器
+ */
+ static FederationQueryAdmissionController unlimited() {
+ return new FederationQueryAdmissionController() {
+ @Override
+ public FederationQueryPermit acquire(
+ SourceId sourceId,
+ QueryId queryId,
+ Duration timeout
+ ) {
+ return () -> { };
+ }
+
+ @Override
+ public FederationQueryPermit acquire(QueryAdmissionRequest request) {
+ return () -> { };
+ }
+ };
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryMetricsSnapshot.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryMetricsSnapshot.java
new file mode 100644
index 0000000..6a19ecd
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryMetricsSnapshot.java
@@ -0,0 +1,152 @@
+package com.easyagents.federation.sql.execute;
+
+import com.easyagents.federation.sql.federation.FederationQueryMode;
+import java.io.Serializable;
+import java.util.List;
+
+/**
+ * 查询游标当前或关闭后的不可变消耗指标快照。
+ *
+ * @param queryId 查询标识;不可用快照可为空
+ * @param queryMode 查询模式
+ * @param planCacheHit 是否命中计划缓存
+ * @param planningNanos 编译阶段耗时
+ * @param admissionWaitNanos 准入等待耗时
+ * @param connectionAcquireNanos 全部分片获取连接累计耗时
+ * @param databaseExecutionNanos 全部分片执行 Statement 累计耗时
+ * @param localExecutionNanos 本地算子累计耗时
+ * @param executionNanos 从执行开始到当前或结束的耗时
+ * @param firstRowNanos 从执行开始到首行的耗时;尚未返回首行时为 -1
+ * @param returnedRows 调用方已消费的最终结果行数
+ * @param returnedBytes 最终结果估算字节数;Adapter 未安全提供时为 -1
+ * @param intermediateRows 全部分片读取的中间结果行数
+ * @param intermediateBytes 全部分片读取的中间结果估算字节数;Adapter 未安全提供时为 -1
+ * @param complete 查询是否正常消费完成
+ * @param cancelled 查询是否因取消结束
+ * @param timedOut 查询是否因统一执行时限结束
+ * @param truncated 是否因最终行数上限停止继续消费
+ * @param terminalErrorCode 失败终态错误码;成功或尚未失败时为空
+ * @param fragments 分片指标
+ * @param localOperators 本地算子指标
+ */
+public record FederationQueryMetricsSnapshot(
+ QueryId queryId,
+ FederationQueryMode queryMode,
+ boolean planCacheHit,
+ long planningNanos,
+ long admissionWaitNanos,
+ long connectionAcquireNanos,
+ long databaseExecutionNanos,
+ long localExecutionNanos,
+ long executionNanos,
+ long firstRowNanos,
+ long returnedRows,
+ long returnedBytes,
+ long intermediateRows,
+ long intermediateBytes,
+ boolean complete,
+ boolean cancelled,
+ boolean timedOut,
+ boolean truncated,
+ String terminalErrorCode,
+ List fragments,
+ List localOperators
+) implements Serializable {
+
+ /**
+ * 防御性复制分片指标。
+ */
+ public FederationQueryMetricsSnapshot {
+ fragments = List.copyOf(fragments == null ? List.of() : fragments);
+ localOperators = List.copyOf(localOperators == null ? List.of() : localOperators);
+ terminalErrorCode = terminalErrorCode == null ? "" : terminalErrorCode;
+ }
+
+ /**
+ * 创建旧基础字段视图的兼容指标快照。
+ *
+ * @param queryId 查询标识
+ * @param queryMode 查询模式
+ * @param planCacheHit 是否命中计划缓存
+ * @param planningNanos 编译耗时
+ * @param executionNanos 执行耗时
+ * @param firstRowNanos 首行耗时
+ * @param returnedRows 返回行数
+ * @param returnedBytes 返回字节数
+ * @param intermediateRows 中间行数
+ * @param intermediateBytes 中间字节数
+ * @param complete 是否完成
+ * @param cancelled 是否取消
+ * @param fragments 分片指标
+ */
+ public FederationQueryMetricsSnapshot(
+ QueryId queryId,
+ FederationQueryMode queryMode,
+ boolean planCacheHit,
+ long planningNanos,
+ long executionNanos,
+ long firstRowNanos,
+ long returnedRows,
+ long returnedBytes,
+ long intermediateRows,
+ long intermediateBytes,
+ boolean complete,
+ boolean cancelled,
+ List fragments
+ ) {
+ this(
+ queryId,
+ queryMode,
+ planCacheHit,
+ planningNanos,
+ 0,
+ 0,
+ 0,
+ 0,
+ executionNanos,
+ firstRowNanos,
+ returnedRows,
+ returnedBytes,
+ intermediateRows,
+ intermediateBytes,
+ complete,
+ cancelled,
+ false,
+ false,
+ "",
+ fragments,
+ List.of()
+ );
+ }
+
+ /**
+ * 返回第三方 Adapter 尚未接入指标时的空快照。
+ *
+ * @return 空快照
+ */
+ public static FederationQueryMetricsSnapshot unavailable() {
+ return new FederationQueryMetricsSnapshot(
+ null,
+ FederationQueryMode.SINGLE_SOURCE,
+ false,
+ 0,
+ 0,
+ 0,
+ 0,
+ 0,
+ 0,
+ -1,
+ 0,
+ -1,
+ 0,
+ -1,
+ false,
+ false,
+ false,
+ false,
+ "",
+ List.of(),
+ List.of()
+ );
+ }
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryPermit.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryPermit.java
new file mode 100644
index 0000000..9010a62
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationQueryPermit.java
@@ -0,0 +1,14 @@
+package com.easyagents.federation.sql.execute;
+
+/**
+ * 查询准入许可,关闭时释放并发配额。
+ */
+@FunctionalInterface
+public interface FederationQueryPermit extends AutoCloseable {
+
+ /**
+ * 释放准入许可。
+ */
+ @Override
+ void close();
+}
diff --git a/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationResultCursor.java b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationResultCursor.java
new file mode 100644
index 0000000..bffa59e
--- /dev/null
+++ b/easy-agents-federation-sql/easy-agents-federation-sql-core/src/main/java/com/easyagents/federation/sql/execute/FederationResultCursor.java
@@ -0,0 +1,84 @@
+package com.easyagents.federation.sql.execute;
+
+import java.io.InputStream;
+import java.io.Reader;
+import java.util.List;
+
+/**
+ * 按行消费且必须关闭的流式结果游标。
+ */
+public interface FederationResultCursor extends AutoCloseable {
+
+ /**
+ * 返回查询标识。
+ *
+ * @return 查询标识
+ */
+ QueryId queryId();
+
+ /**
+ * 返回结果列元数据。
+ *
+ * @return 不可变列列表
+ */
+ List columns();
+
+ /**
+ * 返回查询当前或关闭后的消耗指标快照。
+ *
+ * @return 不可变指标快照
+ */
+ default FederationQueryMetricsSnapshot metrics() {
+ return FederationQueryMetricsSnapshot.unavailable();
+ }
+
+ /**
+ * 移动至下一行。
+ *
+ * @return 是否存在下一行
+ */
+ boolean next();
+
+ /**
+ * 按 JDBC 列序号读取当前行。
+ *
+ * @param columnIndex 从 1 开始的列序号
+ * @return 列值
+ */
+ Object getObject(int columnIndex);
+
+ /**
+ * 以流方式读取二进制列,避免调用方为大字段一次性分配完整字节数组。
+ *
+ * @param columnIndex 从 1 开始的列序号
+ * @return 二进制流;SQL NULL 返回 null
+ * @throws UnsupportedOperationException Adapter 不支持流式列读取
+ */
+ default InputStream getBinaryStream(int columnIndex) {
+ throw new UnsupportedOperationException("binary stream access is not supported by this adapter");
+ }
+
+ /**
+ * 以流方式读取字符列,避免调用方为大字段一次性分配完整字符串。
+ *
+ * @param columnIndex 从 1 开始的列序号
+ * @return 字符流;SQL NULL 返回 null
+ * @throws UnsupportedOperationException Adapter 不支持流式列读取
+ */
+ default Reader getCharacterStream(int columnIndex) {
+ throw new UnsupportedOperationException("character stream access is not supported by this adapter");
+ }
+
+ /**
+ * 将当前行复制为不可变列表。
+ *
+ * @return 当前行列值
+ */
+ List