diff --git a/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottleAutoConfiguration.java b/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottleAutoConfiguration.java index b84d6df1be..107d742ff8 100644 --- a/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottleAutoConfiguration.java +++ b/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottleAutoConfiguration.java @@ -9,25 +9,23 @@ */ package org.eclipse.hawkbit.autoconfigure.throttle; -import java.util.Optional; - import javax.sql.DataSource; +import io.micrometer.core.instrument.MeterRegistry; import lombok.extern.slf4j.Slf4j; +import org.eclipse.hawkbit.throttle.Config; import org.eclipse.hawkbit.throttle.Throttle; -import org.eclipse.hawkbit.throttle.ThrottleProperties; -import org.eclipse.hawkbit.throttle.ThrottleProperties.ThrottleConfig; +import org.jspecify.annotations.NonNull; +import org.springframework.beans.BeansException; import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; /** - * Wires the per-tenant DB connection throttle only if explicitly enabled. - * - *

Default configuration for every data source is read from {@code hawkbit.throttle.db.*}. Per-data-source overrides are read from - * {@code hawkbit.throttle.db-.*}. Overrides are self-contained override - i.e. doesn't inherit anything from default. + * Wires the DB connection throttle only if explicitly enabled. * *

A {@link BeanPostProcessor} wraps the application {@link DataSource} in a {@link ThrottlingDataSourceDecorator}, mirroring the proven * {@code QueryCountConfiguration} pattern. Pool capacity is auto-detected from Hikari's {@code maximumPoolSize} (so granted permits never @@ -35,54 +33,61 @@ */ @Slf4j @AutoConfiguration +@ConditionalOnProperty(prefix = "hawkbit.throttle", name = "enabled", havingValue = "true") @EnableConfigurationProperties(ThrottleProperties.class) @SuppressWarnings("java:S1118") // false positive - auto config instantiated by Spring public class ThrottleAutoConfiguration { - private static final String DB_DOMAIN = "db"; private static final int DEFAULT_CAPACITY = 10; // Hikari default; used only if auto-detect fails @Bean - static BeanPostProcessor throttlingDataSourcePostProcessor(final ObjectProvider provider) { + static BeanPostProcessor throttlingDataSourcePostProcessor( + final ObjectProvider provider, final ObjectProvider meterRegistryProvider) { return new BeanPostProcessor() { @Override - public Object postProcessAfterInitialization(final Object bean, final String beanName) { + public Object postProcessAfterInitialization(@NonNull final Object bean, @NonNull final String beanName) { if (!(bean instanceof DataSource dataSource) || bean instanceof ThrottlingDataSourceDecorator) { return bean; } - final ThrottleConfig props = Optional.ofNullable(provider.getIfAvailable()) - .map(properties -> Optional.ofNullable(properties.getThrottle().get(DB_DOMAIN + "-" + beanName)) - .orElseGet(() -> properties.getThrottle().get(DB_DOMAIN))) - .orElse(null); - if (props != null && props.isEnabled()) { - final Throttle.Policy policy = props.toPolicy(resolveCapacity(props, dataSource)); - log.info( - "hawkBit connection throttle enabled (bean '{}'): capacity={}, limit={}, timeout={}, threshold={}, systemFloor={} slots", - beanName, policy.capacity(), props.getLimit(), props.getTimeout(), - props.getThreshold() == -1 ? " -1 (always enforce limit)" : " " + props.getThreshold() + " slots", - policy.priorityFloor(null)); - return new ThrottlingDataSourceDecorator(dataSource, new Throttle(policy), props.getTimeout()); - } else { - return bean; // the throttle is not enabled + final ThrottleProperties props = provider.getIfAvailable(); + if (props == null) { // should never happen, but just in case + log.warn("hawkBit throttle enabled (for bean '{}'): but no ThrottleProperties available", beanName); + return bean; } + + final Config config = props.toConfig(resolveCapacity(props.getCapacity(), dataSource)); + log.info("hawkBit connection throttle enabled (bean '{}'): props: {}, config: {}", beanName, props, config); + return new ThrottlingDataSourceDecorator( + dataSource, new Throttle(config), props.getTimeout(), + () -> meterRegistry(meterRegistryProvider)); } }; } - private static int resolveCapacity(final ThrottleConfig props, final DataSource dataSource) { - final int configured = props.getCapacity(); - if (configured > 0) { - return configured; + // the data source is decorated long before the metrics infrastructure is up, so the decorator asks for the + // registry lazily. A failure here means the context is not far enough along to have one - report no metrics + // rather than let a monitoring concern break a connection acquisition + private static MeterRegistry meterRegistry(final ObjectProvider meterRegistryProvider) { + try { + return meterRegistryProvider.getIfAvailable(); + } catch (final BeansException e) { + log.debug("MeterRegistry not (yet) resolvable, throttle metrics are skipped", e); + return null; + } + } + + private static int resolveCapacity(final int configuredCapacity, final DataSource dataSource) { + if (configuredCapacity >= 0) { + return configuredCapacity; } - // -1 or 0 → auto-detect + // < 0 → auto-detect final Integer poolSize = hikariMaximumPoolSize(dataSource); - if (poolSize != null && poolSize > 0) { + if (poolSize != null && poolSize >= 0) { return poolSize; } - log.warn("Could not auto-detect DB pool size for throttling; set hawkbit.throttle.{}.capacity. Falling back to {}", - DB_DOMAIN, DEFAULT_CAPACITY); + log.warn("Could not auto-detect DB pool size for throttling; set hawkbit.throttle.capacity. Falling back to {}", DEFAULT_CAPACITY); return DEFAULT_CAPACITY; } diff --git a/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottleProperties.java b/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottleProperties.java new file mode 100644 index 0000000000..b74b5aa078 --- /dev/null +++ b/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottleProperties.java @@ -0,0 +1,171 @@ +/** + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.hawkbit.autoconfigure.throttle; + +import static org.eclipse.hawkbit.throttle.Config.PendingRemoveStrategy.FIFO; +import static org.eclipse.hawkbit.throttle.Config.Priority.HIGH; +import static org.eclipse.hawkbit.throttle.Config.Priority.NORMAL; + +import java.time.Duration; +import java.util.EnumMap; +import java.util.Locale; +import java.util.Map; +import java.util.stream.Collectors; + +import lombok.Data; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.hawkbit.throttle.Config; +import org.eclipse.hawkbit.throttle.Config.Policy; +import org.eclipse.hawkbit.throttle.ThrottledException; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.util.LinkedCaseInsensitiveMap; + +/** + * Operator-owned per-tenant throttling configuration. Caps how many of the shared db connection a single tenant may hold at once. + * Per-tenant overrides live in the {@link #tenants} map. + * + *

Configuration Examples: + *

{@code
+ * # Mode 1: Disabled (default)
+ * hawkbit.throttle.enabled=false
+ *
+ * # Mode 2: Hard cap, fast-reject
+ * hawkbit.throttle.enabled=true
+ * hawkbit.throttle.limit=5           # each tenant max 5 connections
+ * hawkbit.throttle.timeout=0         # reject immediately when over limit
+ * hawkbit.throttle.threshold=-1      # no threshold (always enforce limit)
+ *
+ * # Mode 3: Contention-aware burst allowance
+ * hawkbit.throttle.enabled=true
+ * hawkbit.throttle.limit=10          # ceiling per tenant
+ * hawkbit.throttle.timeout=20s       # wait up to 20s before HTTP 429
+ * hawkbit.throttle.max-pending=190   # since we block thread (20 sec) we reject if threads become too many.
+ *                                    # over all tenants, so capacity + max-pending stays below the container's limit
+ * hawkbit.throttle.threshold=7       # fairness kicks in at 7/10 pool utilization
+ * }
+ */ +@Slf4j +@Data +@ConfigurationProperties("hawkbit.throttle") // real prefix is "hawkbit.throttle" but the rest comes from the single value - throttle map +public class ThrottleProperties { + + /** + * Master switch. Default OFF ⇒ no DataSource wrapping, no filter, byte-for-byte current behavior. + * + *

Default: {@code false}, disabled + */ + private boolean enabled = false; + + /** + * Total permits. {@code -1} = auto-detect the Hikari {@code maximumPoolSize}; set explicitly only when auto-detect is not possible. + * + *

Default: {@code -1}, auto-detect + */ + private int capacity = -1; + + /** + * Global per-tenant ceiling on concurrent connection borrows. + * + *

Default: {@code -1}, no limit + */ + private int limit = -1; + + /** + * Max wait for a permit before {@link ThrottledException} (→ HTTP 429). + * {@code 0} fast-rejects instead of blocking; a positive value is wait-then-served backpressure. + * + *

Default: 0, fast-reject + */ + private Duration timeout = Duration.ZERO; + + /** + * Max total pending requests, {@code -1} = no limit, limit it when no virtual threads. {@code 0} = no pending. + * + *

Default: {@code -1}, no limit + */ + private int maxPending = -1; + + /** + * Contention threshold (absolute slot count). When global in-use reaches this threshold, fairness (per key/tenant) enforcement activates: + * each tenant is served in round-robin, capped by its limit. Below the threshold, any tenant may burst freely up to pool capacity. + * {@code <= 0} = no threshold (always enforce hard cap on tenant limit). + * + *

Default: {@code -1}, hard cap mode + */ + private int threshold = -1; + + /** + * Capacity that internal, non-tenant work is served ahead of the queue for ({@link Config.Priority#HIGH}. + * Internal work is charged to the {@code null} key — schedulers, migrations and health checks run without a tenant context. + * + *

A floor it cannot be starved below, not a cap: beyond it system work competes as an ordinary key, so it + * can neither be crowded out by tenant load nor crowd tenants out itself. Costs no reserved capacity — tenants + * use those slots freely whenever system work is not. {@code 0} disables the priority. Must resolve to fewer + * slots than {@link #capacity}, or tenant work may never be reached; the policy warns if it does not. + * + *

Default: {@code 4}, system work is served ahead of the queue for 4 slots + */ + private int systemGranted = 4; + + /** Per-tenant limit overrides */ + private final LinkedCaseInsensitiveMap tenants = new LinkedCaseInsensitiveMap<>(Locale.ROOT); + + @Data + public static class TenantConfig { + + private int limit = -1; + private Duration timeout; + private int maxPending = -1; + } + + public Config toConfig(final int capacity) { + if (systemGranted >= capacity) { + log.warn(""" + Throttle system granted ({}) is not below capacity ({}): priority admissions can consume the whole \ + resource, so tenant work may never be served. Lower `hawkbit.throttle.system-granted`.""", + systemGranted, capacity); + } + // pre init system and per tenant policy overrides in order to do not build for avery key (tenant) + final Policy systemPolicy = Policy.builder() + .permits(capacity) + .timeout(timeout) + .pendingRemoveStrategy(FIFO) + .priority(HIGH) // set high priority to use high priority granted + .build(); + final Policy defaultTenantPolicy = Policy.builder() // policy if not explicitly overridden + .permits(limit < 0 ? capacity : limit) + .timeout(timeout) + .pendingRemoveStrategy(FIFO) + .priority(NORMAL) + .build(); + final Map tenantToPolicy = tenants.entrySet().stream() // overrides + .collect(Collectors.toMap(Map.Entry::getKey, e -> { + final TenantConfig tenantConfig = e.getValue(); + return Policy.builder() + .permits(tenantConfig.limit == -1 ? capacity : tenantConfig.limit) + .timeout(tenantConfig.timeout == null ? timeout : tenantConfig.timeout) + .maxPending(tenantConfig.maxPending) + .pendingRemoveStrategy(FIFO) + .priority(NORMAL) + .build(); + }, (policy, duplicate) -> policy, () -> new LinkedCaseInsensitiveMap<>(Locale.ROOT))); // case-insensitive map + return Config.builder() + .priorityToGranted(new EnumMap<>(Map.of(HIGH, systemGranted))) // reserve system priority granted + .permits(capacity) + .timeout(timeout) + .maxPending(maxPending) + // we want to remove the first pending - best chance to have been already expired on caller + .pendingRemoveStrategy(FIFO) + .keyPolicyThreshold(Math.min(threshold, capacity)) + // tenant null <=> system. It still could be a tenant work, but we still don't know it + .keyPolicyFn(tenant -> tenant == null ? systemPolicy : tenantToPolicy.getOrDefault(tenant, defaultTenantPolicy)) + .build(); + } +} diff --git a/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecorator.java b/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecorator.java index d06598fe5f..acc6a847b9 100644 --- a/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecorator.java +++ b/hawkbit-autoconfigure/src/main/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecorator.java @@ -9,15 +9,22 @@ */ package org.eclipse.hawkbit.autoconfigure.throttle; +import static org.eclipse.hawkbit.tenancy.DefaultTenantConfiguration.TENANT_TAG; +import static org.eclipse.hawkbit.tenancy.DefaultTenantConfiguration.TENANT_TAG_VALUE_PROVIDER; + import java.sql.Connection; import java.sql.SQLException; import java.time.Duration; +import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; import javax.sql.DataSource; +import io.micrometer.core.instrument.MeterRegistry; import lombok.experimental.Delegate; import org.eclipse.hawkbit.context.AccessContext; import org.eclipse.hawkbit.throttle.Throttle; +import org.eclipse.hawkbit.throttle.ThrottledException; import org.eclipse.hawkbit.throttle.Permit; import org.springframework.jdbc.datasource.DelegatingDataSource; @@ -30,16 +37,35 @@ */ class ThrottlingDataSourceDecorator extends DelegatingDataSource { + // how long a caller waited to actually get a connection - the throttle queue plus the pool itself, which is the + // latency the request pays. Only successful acquisitions are recorded; refusals are counted by REJECTED_METER + private static final String CONNECTION_METER = "hawkbit.throttle.connection"; + // refused acquisitions, split by why - see ThrottledException.Reason + private static final String REJECTED_METER = "hawkbit.throttle.rejected"; + private static final String REASON_TAG = "reason"; + // per instance - for every data source private final ThreadLocal rootPermit = new ThreadLocal<>(); private final Throttle throttle; private final Duration timeout; + private final Supplier meterRegistrySupplier; + + // the registry is resolved on first use, not injected: this data source is built by a BeanPostProcessor, so at + // construction time the metrics infrastructure does not exist yet. null until then, and metrics are simply skipped + private volatile MeterRegistry meterRegistry; ThrottlingDataSourceDecorator(final DataSource targetDataSource, final Throttle throttle, final Duration timeout) { + this(targetDataSource, throttle, timeout, () -> null); + } + + ThrottlingDataSourceDecorator( + final DataSource targetDataSource, final Throttle throttle, final Duration timeout, + final Supplier meterRegistrySupplier) { super(targetDataSource); this.throttle = throttle; this.timeout = timeout; + this.meterRegistrySupplier = meterRegistrySupplier; } @Override @@ -53,14 +79,20 @@ public Connection getConnection(final String username, final String password) th } private Connection acquire(final Throttle throttle, final Duration timeout, final ConnectionSupplier raw) throws SQLException { + final long startNano = System.nanoTime(); final Permit root = rootPermit.get(); final Permit permit; - if (root == null) { - // AccessContext.tenant() is null for internal work, which is exactly the throttle's key for it - permit = throttle.acquire(AccessContext.tenant(), timeout); - rootPermit.set(permit); - } else { - permit = root.acquire(timeout); + try { + if (root == null) { + // AccessContext.tenant() is null for internal work, which is exactly the throttle's key for it + permit = throttle.acquire(AccessContext.tenant(), timeout); + rootPermit.set(permit); + } else { + permit = root.acquire(timeout); + } + } catch (final ThrottledException e) { + countRejected(e); + throw e; } final Connection connection; @@ -70,9 +102,35 @@ private Connection acquire(final Throttle throttle, final Duration timeout, fina release(permit); throw e; } + recordConnectionTime(System.nanoTime() - startNano); return new ThrottledConnection(connection, permit); // wrapper that will release permit when closed } + private void countRejected(final ThrottledException e) { + final MeterRegistry registry = meterRegistry(); + if (registry != null) { + registry.counter(REJECTED_METER, TENANT_TAG, TENANT_TAG_VALUE_PROVIDER.get(), REASON_TAG, e.getReason().name()).increment(); + } + } + + private void recordConnectionTime(final long nanos) { + final MeterRegistry registry = meterRegistry(); + if (registry != null) { + registry.timer(CONNECTION_METER, TENANT_TAG, TENANT_TAG_VALUE_PROVIDER.get()).record(nanos, TimeUnit.NANOSECONDS); + } + } + + // resolved once, on the first acquisition after the registry bean exists; before that every call re-asks, which + // is a map lookup in the bean factory and only happens while the context is still coming up + private MeterRegistry meterRegistry() { + MeterRegistry registry = meterRegistry; + if (registry == null) { + registry = meterRegistrySupplier.get(); + meterRegistry = registry; + } + return registry; + } + private void release(final Permit permit) { try { permit.close(); diff --git a/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecoratorTest.java b/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecoratorTest.java index 5e69a89e40..0ae380c134 100644 --- a/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecoratorTest.java +++ b/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingDataSourceDecoratorTest.java @@ -24,7 +24,6 @@ import com.zaxxer.hikari.HikariDataSource; import org.eclipse.hawkbit.throttle.Permit; import org.eclipse.hawkbit.throttle.Throttle; -import org.eclipse.hawkbit.throttle.ThrottleProperties.ThrottleConfig; import org.eclipse.hawkbit.throttle.ThrottledException; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -182,12 +181,12 @@ private interface ThrowingRunnable { void run() throws Exception; } - private static Throttle throttle(final int capacity, final Consumer customizer) { - final ThrottleConfig config = new ThrottleConfig(); - config.setThreshold(0); // always contended: fair share enforced from the first permit - config.setSystemFloor(0); // priority off unless a test opts in - customizer.accept(config); - return new Throttle(config.toPolicy(capacity)); + private static Throttle throttle(final int capacity, final Consumer customizer) { + final ThrottleProperties props = new ThrottleProperties(); + props.setThreshold(0); // always contended: fair share enforced from the first permit + props.setSystemGranted(0); // priority off unless a test opts in + customizer.accept(props); + return new Throttle(props.toConfig(capacity)); } private ThrottlingDataSourceDecorator dataSource(final Throttle throttle) { diff --git a/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingRequiresNewTransactionTest.java b/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingRequiresNewTransactionTest.java index 0570533f72..dbf91860a8 100644 --- a/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingRequiresNewTransactionTest.java +++ b/hawkbit-autoconfigure/src/test/java/org/eclipse/hawkbit/autoconfigure/throttle/ThrottlingRequiresNewTransactionTest.java @@ -19,7 +19,6 @@ import com.zaxxer.hikari.HikariConfig; import com.zaxxer.hikari.HikariDataSource; import org.eclipse.hawkbit.throttle.Throttle; -import org.eclipse.hawkbit.throttle.ThrottleProperties.ThrottleConfig; import org.eclipse.hawkbit.throttle.ThrottledException; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -42,6 +41,8 @@ @EnableTransactionManagement class ThrottlingRequiresNewTransactionTest { + private static final Duration TIMEOUT = Duration.ofSeconds(1); + private HikariDataSource rawDataSource; @BeforeEach @@ -132,15 +133,17 @@ private static TransactionTemplate requiresNew(final PlatformTransactionManager return template; } - private static Throttle throttle(final int capacity, final Consumer customizer) { - final ThrottleConfig config = new ThrottleConfig(); - config.setThreshold(0); - config.setSystemFloor(0); - customizer.accept(config); - return new Throttle(config.toPolicy(capacity)); + private static Throttle throttle(final int capacity, final Consumer customizer) { + final ThrottleProperties props = new ThrottleProperties(); + props.setThreshold(0); + props.setSystemGranted(0); + // the engine gates waiting on the configured timeout too, so it must allow at least what txManager() passes + props.setTimeout(TIMEOUT); + customizer.accept(props); + return new Throttle(props.toConfig(capacity)); } private PlatformTransactionManager txManager(final Throttle throttle) { - return new DataSourceTransactionManager(new ThrottlingDataSourceDecorator(rawDataSource, throttle, Duration.ofSeconds(1))); + return new DataSourceTransactionManager(new ThrottlingDataSourceDecorator(rawDataSource, throttle, TIMEOUT)); } } diff --git a/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/Config.java b/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/Config.java new file mode 100644 index 0000000000..e2a3fbe3ab --- /dev/null +++ b/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/Config.java @@ -0,0 +1,106 @@ +/** + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.hawkbit.throttle; + +import java.time.Duration; +import java.util.EnumMap; +import java.util.Optional; +import java.util.function.Function; + +import lombok.AccessLevel; +import lombok.Builder; +import lombok.Getter; +import lombok.ToString; +import lombok.Value; +import lombok.experimental.Accessors; + +@Builder +@Value +@Accessors(fluent = true) +public class Config { + + // free permits are granted by priority, up to the granted. after granted the priority is not considered + @Builder.Default + @Getter(AccessLevel.PRIVATE) + EnumMap priorityToGranted = new EnumMap<>(Priority.class); + + int permits; // capacity, max permits to be granted + @Builder.Default + @Getter(AccessLevel.PRIVATE) + Duration timeout = Duration.ZERO; // non-positive timeout means direct reject. In this case maxPending is ignored + @Builder.Default + int maxPending = -1; // max pending requests over all keys together, -1 means no limit + @Builder.Default + PendingRemoveStrategy pendingRemoveStrategy = PendingRemoveStrategy.FIFO; // strategy to remove pending requests when maxPending is reached + int keyPolicyThreshold; // shall be less than permits, when lest that number permits are acquired - doesn't consider key specific permits limit + @Builder.Default + @ToString.Exclude + @Getter(AccessLevel.PRIVATE) + Function keyPolicyFn = key -> null; // function to get the policy for a given key + + public int priorityGranted(final String key) { + final Policy keyPolicy = keyPolicyFn.apply(key); + return Optional.ofNullable(priorityToGranted.get(keyPolicy == null ? Priority.NORMAL : keyPolicy.priority())).orElse(0); + } + + //--- per key --- + + public int permits(final String key) { + final Policy keyPolicy = keyPolicyFn.apply(key); + if (keyPolicy == null) { + return permits; + } else { + final int keyPermits = keyPolicy.permits(); + return keyPermits >= 0 ? keyPermits : permits; // if key policy has permits >= 0, use it, otherwise use default permits + } + } + + public Duration timeout(final String key) { + return Optional.ofNullable(keyPolicyFn.apply(key)).map(Policy::timeout).orElse(timeout); + } + + @SuppressWarnings("java:S3358") // easy readable that way + public int maxPending(final String key) { + final Policy keyPolicy = keyPolicyFn.apply(key); + return keyPolicy == null || keyPolicy.maxPending() < 0 ? + maxPending : + maxPending < 0 ? keyPolicy.maxPending() : Math.min(maxPending, keyPolicy.maxPending()); + } + + public PendingRemoveStrategy pendingRemoveStrategy(final String key) { + return Optional.ofNullable(keyPolicyFn.apply(key)).map(Policy::pendingRemoveStrategy).orElse(pendingRemoveStrategy); + } + + public enum Priority { + HIGH, + NORMAL, + LOW + } + + public enum PendingRemoveStrategy { + LIFO, // last in first out + FIFO // fist in first out + } + + @Builder + @Value + public static class Policy { + + @Builder.Default + int permits = -1; // capacity, max permits to be granted, <= 0 means no permits (blocked) + Duration timeout; // non-positive timeout means direct reject. In this case maxPending is ignored + @Builder.Default + int maxPending = -1; // max pending requests for this key, -1 means no key bound (only Config#maxPending), 0 - no pending + @Builder.Default + PendingRemoveStrategy pendingRemoveStrategy = PendingRemoveStrategy.FIFO; // strategy to remove pending requests when maxPending is reached + @Builder.Default + Priority priority = Priority.NORMAL; // priority of the policy, default is MEDIUM + } +} diff --git a/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/Throttle.java b/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/Throttle.java index cfd72cb087..c65c27ac78 100644 --- a/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/Throttle.java +++ b/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/Throttle.java @@ -9,39 +9,43 @@ */ package org.eclipse.hawkbit.throttle; +import static org.eclipse.hawkbit.throttle.Config.PendingRemoveStrategy.LIFO; + import java.time.Duration; import java.util.ArrayDeque; import java.util.Collections; import java.util.Deque; import java.util.HashMap; +import java.util.Iterator; import java.util.LinkedList; import java.util.Map; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.ReentrantLock; import lombok.extern.slf4j.Slf4j; +import org.eclipse.hawkbit.throttle.ThrottledException.Reason; /** * Contention-aware, priority- and fair-share admission engine over a single finite resource (e.g. the DB connection pool, REST APIs). * *

Every borrow of the guarded resource takes a permit, with no exempt path. That makes granted permits an exact mirror - * of resource usage, which in turn lets this engine be the only queue in front of the resource: if {@link Policy#capacity()} equals - * the real resource size, the resource itself never has to block. + * of resource usage, which in turn lets this engine be the only queue in front of the resource: + * if {@link org.eclipse.hawkbit.throttle.Config.Policy#permits()} equals the real resource size, the resource itself never has to block. * *

A permit is charged to a key — e.g. a tenant id, or {@code null} for internal, non-tenant work (schedulers, health checks). * - *

Permits form a tree per unit of work: {@link #acquire(String, Duration)} starts one, {@link Permit#acquire(Duration)} - * joins it. A child request comes from a caller that already holds part of the resource, which is the only way this - * engine can deadlock — see {@link Permit} for the full parent/child contract. + *

Permits form a tree per unit of work: {@link #acquire(String, Duration)} starts one, {@link Permit#acquire(Duration)} joins it. + * A child request comes from a caller that already holds part of the resource, which is the only way this engine can deadlock — see + * {@link Permit} for the full parent/child contract. * *

Admission order (first match wins): *

    *
  1. Children — bypass fair share entirely. A caller blocked while holding part of the resource is what deadlocks it, so unblocking * work already in flight always outranks admitting new work.
  2. - *
  3. Keys below their {@link Policy#priorityFloor(String)} — a floor internal work cannot be starved below, costing no reserved capacity, - * since tenants use those slots freely whenever internal work is not.
  4. - *
  5. Below {@link Policy#burstThreshold()} — anyone may burst.
  6. - *
  7. At or above the threshold — round-robin + ceiling bound share.
  8. + *
  9. Keys below their ({@link Config#priorityGranted(String)} granted — a floor internal work cannot be starved below, costing no + * reserved capacity, since tenants use those slots freely whenever internal work is not.
  10. + *
  11. Below {@link Config#keyPolicyThreshold()} — anyone may burst.
  12. + *
  13. At or above the threshold — round-robin + per key max permit bound share.
  14. *
* *

Fairness across keys is round-robin, not first-come-first-served by key. Waiters queue per key, and the keys @@ -50,6 +54,13 @@ * Only service rotates a key: a waiter that times out, and a key skipped for being at its fair share, both * leave the order untouched, so neither is penalised. * + *

Waiting is bounded too, not just admission. A parked waiter holds its caller's thread, so an unbounded + * queue simply moves exhaustion from the guarded resource to the thread pool in front of it. Two bounds apply, both + * only to requests that would otherwise park: {@link Config#maxPending()} over all keys together, and + * {@link Config#maxPending(String)} per key. On overflow one waiter is dropped rather than the arrival being + * refused outright — for the global bound it is taken from the longest queue, so a key flooding the queue + * pays for it instead of pushing a quiet key's waiters out. + * *

Deadlock is avoided rather than prevented, so no capacity is reserved. A child request may block only while * at least one other unit of work is still running and will therefore release. The arrival that would leave every unit * waiting on another is refused instead of parked. @@ -60,31 +71,31 @@ @Slf4j public class Throttle { - private final Policy policy; + private final Config config; private final ReentrantLock lock = new ReentrantLock(); - private int globalInUse; - private final Map perKeyInUse = new HashMap<>(); + private int globalGranted; + private final Map perKeyGranted = new HashMap<>(); // permits held per unit of work, keyed on the tree's root permit; entry removed at zero. Self-bounding: every unit // owns at least one permit, so this can never exceed capacity entries — it cannot leak the way a caller-side // thread registry can, because an entry only exists while a permit does. private final Map permitsByRoot = new HashMap<>(); - // wait state: per-key FIFO of waiters + the keys themselves in service order (rotated on every admission). - private final Map> waitersByKey = new HashMap<>(); + private int globalPending; + // wait state: per-key FIFO of pending + the keys themselves in service order (rotated on every admission). + private final Map> perKeyPending = new HashMap<>(); // LinkedList, not ArrayDeque: the system key is null and ArrayDeque rejects null elements. - private final Deque waitingKeys = new LinkedList<>(); - - // child waiters across all keys, in arrival order; drives both admission tier 1 and the cycle check - private final Deque childWaiters = new ArrayDeque<>(); + private final Deque pendingKeys = new LinkedList<>(); + // pending children across all keys, in arrival order; drives both admission tier 1 and the cycle check + private final Deque pendingChildren = new ArrayDeque<>(); // contention tracking for WARN-on-transition logging private boolean wasAboveThreshold; private boolean wasAtCapacity; - public Throttle(final Policy policy) { - this.policy = policy; + public Throttle(final Config config) { + this.config = config; } /** @@ -109,7 +120,7 @@ public Stats stats() { lock.lock(); try { // HashMap copy, not Map.copyOf: the system key is null, which Map.copyOf rejects - return new Stats(globalInUse, Collections.unmodifiableMap(new HashMap<>(perKeyInUse)), permitsByRoot.size()); + return new Stats(globalGranted, Collections.unmodifiableMap(new HashMap<>(perKeyGranted)), permitsByRoot.size(), globalPending); } finally { lock.unlock(); } @@ -118,15 +129,17 @@ public Stats stats() { // single entry point; parent == null starts a new unit of work, otherwise the permit joins the parent's // java:S2093 - permit must be closed by caller site // java:S899 - awaitNanos result intentionally ignored; the deadline is re-checked by the loop - @SuppressWarnings({ "java:S2093", "java:S899", "java:S3776" }) + // java:S3776 - better readable in one place + // java:S2274 - fine, false positive of sonar requiring a `while` + @SuppressWarnings({ "java:S2093", "java:S899", "java:S3776", "java:S2274" }) private Permit acquire(final String key, final PermitImpl parent, final Duration timeout) { lock.lock(); try { - if (parent == null && ceiling(key) == 0) { - throw new ThrottledException("'" + key + "' is throttled (blocked)"); + if (parent == null && config.permits(key) == 0) { // handle blocked + throw new ThrottledException(Reason.BLOCKED, "'%s' is throttled (blocked)".formatted(key)); } - if (parent != null && !permitsByRoot.containsKey(parent.root)) { - throw new IllegalStateException("cannot acquire from a stale permit (key '" + key + "')"); + if (parent != null && !permitsByRoot.containsKey(parent.root)) { // handle stale parent + throw new IllegalStateException("cannot acquire from a stale permit (key '%s')".formatted(key)); } final PermitImpl permit = new PermitImpl(key, parent, lock.newCondition()); @@ -134,31 +147,59 @@ private Permit acquire(final String key, final PermitImpl parent, final Duration try { grantNextEligible(); // may grant this permit immediately if a slot is free - if (!permit.admitted) { - // Refuse rather than park when parking would leave every unit of work waiting on another. - if (permit.isChild() && childWaiters.size() >= permitsByRoot.size()) { // cycle break + if (permit.state != PermitImpl.State.GRANTED) { + if (!timeout.isPositive()) { // no wait requestedd, fast-reject + throw new ThrottledException(Reason.LIMIT, "'%s' is throttled (limit)".formatted(key)); + } + + final Duration keyTimeout = config.timeout(key); + if (!keyTimeout.isPositive()) { // no wait is allowed, fast-reject + throw new ThrottledException(Reason.FAST_REJECT, "'%s' is throttled (limit, fast-reject)".formatted(key)); + } // else eligible to wait + + // handle children / deadlocks, refuse rather than park when parking would leave every unit of work waiting on another. + if (permit.isChild() && pendingChildren.size() >= permitsByRoot.size()) { // cycle break log.debug("{} is throttled: nested request refused, waiting would deadlock all {} units of work holding {}/{} slots", - key, permitsByRoot.size(), globalInUse, policy.capacity()); - throw new ThrottledException("'" + key + "' is throttled (nested / deadlock)"); + key, permitsByRoot.size(), globalGranted, config.permits(key)); + throw new ThrottledException(Reason.DEADLOCK, "'%s' is throttled (nested / deadlock)".formatted(key)); } - if (timeout.isPositive()) { - final long deadline = System.nanoTime() + timeout.toNanos(); - long remaining; - while (!permit.admitted && (remaining = deadline - System.nanoTime()) > 0) { - try { - permit.condition.awaitNanos(remaining); - } catch (final InterruptedException e) { - Thread.currentThread().interrupt(); - break; // granted just before the interrupt is honored by the caller's admitted check - } + // handle max pending, per key first - an eviction there also relieves the global bound, + // since the evicted permit is counted by both + final Deque keyPending = perKeyPending.get(key); + final int keyMaxPending = config.maxPending(key); + if (keyMaxPending >= 0 && keyPending.size() > keyMaxPending + && evictOrReject(keyPending, permit, config.pendingRemoveStrategy(key))) { + log.trace("{} is throttled: key pending queue full, this", key); + // throw, finally deque will remove it + throw new ThrottledException(Reason.KEY_QUEUE_FULL, "'%s' is throttled (key pending queue full, this)".formatted(key)); + } + + final int maxPending = config.maxPending(); + if (maxPending >= 0 && globalPending > maxPending + && evictOrReject(longestPending(keyPending), permit, config.pendingRemoveStrategy())) { + log.trace("{} is throttled: pending queue full, this", key); + // throw, finally deque will remove it + throw new ThrottledException(Reason.QUEUE_FULL, "'%s' is throttled (pending queue full, this)".formatted(key)); + } + + final long deadline = System.nanoTime() + Math.min(timeout.toNanos(), keyTimeout.toNanos()); + for (long remaining; permit.state == PermitImpl.State.WAITING && (remaining = deadline - System.nanoTime()) > 0; ) { + try { + permit.condition.awaitNanos(remaining); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + break; // granted just before the interrupt is honored by the caller's admitted check } - } // else → fast reject; caller throws - } + } - if (!permit.admitted) { - throw new ThrottledException("'" + key + "' is throttled (" + (timeout.isPositive() ? "timeout)" : "limit, fast-reject)")); + if (permit.state != PermitImpl.State.GRANTED) { + // woken before the deadline means someone took this slot in the queue, not that time ran out + final Reason reason = deadline - System.nanoTime() > 0 ? Reason.EVICTED : Reason.TIMEOUT; + throw new ThrottledException(reason, "'%s' is throttled (%s)".formatted(key, reason)); + } } + return permit; } finally { permit.dequeue(); // idempotent: no-op if already removed by grantNextEligible @@ -169,115 +210,123 @@ private Permit acquire(final String key, final PermitImpl parent, final Duration } } + /** + * Frees one slot in {@code queue} following {@code strategy}, skipping children - a child is a caller that already + * holds part of the resource, so dropping it throws away work that is about to release rather than relieving + * pressure. When every waiter is a child the first one scanned is taken anyway, since the bound must hold. + * + * @return {@code true} when the arriving permit is itself the victim and the caller must reject it + */ + // must be called under lock + private boolean evictOrReject(final Deque queue, final PermitImpl permit, final Config.PendingRemoveStrategy strategy) { + final Iterator it = strategy == LIFO ? queue.descendingIterator() : queue.iterator(); + PermitImpl victim = it.next(); // may be child, must have at least one - this has been just enqueued + while (victim.isChild() && it.hasNext()) { // try to find first non-child victim + final PermitImpl next = it.next(); + if (!next.isChild()) { + victim = next; + break; + } + } + if (victim == permit) { + return true; + } + log.trace("{} evicted from the pending queue to make room", victim.key); + victim.interrupt(); // dequeues (so both counters drop) and wakes it to be rejected on its own thread + return false; + } + + /** + * The longest pending queue, i.e. the key responsible for the global overflow, so a key flooding the queue cannot + * push other keys' waiters out. + * + *

Must be called under lock. + */ + private Deque longestPending(final Deque preferred) { + Deque longest = preferred; + for (final Deque queue : perKeyPending.values()) { + if (queue.size() > longest.size()) { + longest = queue; + } + } + return longest; + } + // grants at most one waiter; must be called under lock private void grantNextEligible() { - final int global = globalInUse; - final int capacity = policy.capacity(); + final int global = globalGranted; + final int capacity = config.permits(); if (global >= capacity) { log.debug("grantNextEligible() → no admit (resource full: {}/{})", global, capacity); return; } - if (waitingKeys.isEmpty()) { + if (pendingKeys.isEmpty()) { return; } // tier 1: children bypass fair share — unblocking a unit of work releases more than it takes - final PermitImpl child = childWaiters.peekFirst(); + final PermitImpl child = pendingChildren.peekFirst(); if (child != null) { log.debug("grantNextEligible[{}] → admit (child, {}/{})", child.key, global, capacity); - child.admit(); + child.grant(); return; } // Pick first, admit after: admit() rotates the served key, which mutates waitingKeys — doing that while an // iterator over it is still live would be a ConcurrentModificationException waiting to happen. - final PermitImpl selected = selectWaiter(global, capacity); + final PermitImpl selected = selectPending(global, capacity); if (selected != null) { - selected.admit(); + selected.grant(); } } /** * The waiter to serve, or {@code null} if nobody is admissible. Returns the waiter rather than its key because - * {@code null} is itself a valid key. Scans {@link #waitingKeys} in service order, so a key rotated to the back by + * {@code null} is itself a valid key. Scans {@link #pendingKeys} in service order, so a key rotated to the back by * an earlier grant is considered last. Must be called under lock. */ - private PermitImpl selectWaiter(final int global, final int capacity) { - // tier 2: keys still below their priority floor - for (final String key : waitingKeys) { - final int floor = policy.priorityFloor(key); - if (floor > 0 && perKeyInUse.getOrDefault(key, 0) < floor) { - log.debug("grantNextEligible[{}] → admit (priority floor {}, {}/{})", key, floor, global, capacity); - return waitersByKey.get(key).peekFirst(); + private PermitImpl selectPending(final int global, final int capacity) { + // tier 2: keys still below their priority granted + for (final String key : pendingKeys) { + final int granted = config.priorityGranted(key); + if (granted > 0 && perKeyGranted.getOrDefault(key, 0) < granted) { + log.debug("grantNextEligible[{}] → admit (priority granted {}, {}/{})", key, granted, global, capacity); + return perKeyPending.get(key).peekFirst(); } } // tier 3: below the contention threshold everyone may burst. - // threshold == -1 (or 0) means always enforce the fair share (no burst allowance) - final int threshold = policy.burstThreshold(); + // threshold <= 0 means always enforce the fair share (no burst allowance) + final int threshold = config.keyPolicyThreshold(); if (threshold > 0 && global < threshold) { - final String key = waitingKeys.peekFirst(); + final String key = pendingKeys.peekFirst(); log.debug("grantNextEligible[{}] → admit (burst: {}/{}, threshold={})", key, global, capacity, threshold); - return waitersByKey.get(key).peekFirst(); + return perKeyPending.get(key).peekFirst(); } // tier 4: up to ceiling share - for (final String key : waitingKeys) { - final int keyCeiling = ceiling(key); - final int keyCurrent = perKeyInUse.getOrDefault(key, 0); - final boolean admit = keyCurrent < keyCeiling; - log.debug("grantNextEligible[{}] → {} (current={}, ceiling={}, global={}/{}, threshold={})", - key, admit, keyCurrent, keyCeiling, global, capacity, threshold); - if (admit) { - return waitersByKey.get(key).peekFirst(); + for (final String key : pendingKeys) { + final int keyPermits = config.permits(key); + final int keyCurrent = perKeyGranted.getOrDefault(key, 0); + final boolean grant = keyCurrent < keyPermits; + log.debug("grantNextEligible[{}] → {} (current={}, permits={}, global={}/{}, threshold={})", + key, grant, keyCurrent, keyPermits, global, capacity, threshold); + if (grant) { + return perKeyPending.get(key).peekFirst(); } } return null; } - private int ceiling(final String key) { - final int capacity = policy.capacity(); - final int limit = policy.ceiling(key); - return limit < 0 ? capacity : Math.min(limit, capacity); - } - - /** Admission policy. Everything operator-tunable lives here so the engine itself stays pure. */ - public interface Policy { - - /** Total permits, i.e. the size of the guarded resource. */ - int capacity(); - - /** - * Global in-use count at which fair-share enforcement activates. Below it any key may burst freely. - * {@code -1} disables bursting, i.e. always enforce the fair share. - */ - int burstThreshold(); - - /** - * Per-key ceiling on concurrently held permits; {@code -1} = unlimited (bounded only by capacity). - * - * @param key the key, {@code null} for internal non-tenant work - */ - int ceiling(String key); - - /** - * Permits this key is served ahead of the queue for, while holding fewer than that many. {@code 0} = no priority permits. - * Acts as a floor the key cannot be starved below, not as a cap — beyond it the key competes normally, which is what keeps - * a privileged key from crowding everyone else out. - * - * @param key the key, {@code null} for internal non-tenant work - */ - int priorityFloor(String key); - } - /** * A consistent snapshot of resource usage, taken under the engine's lock. * * @param inUse total permits held * @param perKey permits held per key; keys with none are absent. The {@code null} key is internal, non-tenant work * @param units distinct units of work holding at least one permit — fewer than {@code inUse} when nesting is in play + * @param pending waiters parked over all keys together; the value bounded by the global max pending */ - public record Stats(int inUse, Map perKey, int units) { + public record Stats(int inUse, Map perKey, int units, int pending) { /** Permits held by {@code key}; {@code null} for internal, non-tenant work. */ public int inUse(final String key) { @@ -287,11 +336,14 @@ public int inUse(final String key) { private final class PermitImpl implements Permit { + private enum State { + WAITING, GRANTED, CLOSED /* rejected or granted and then closed*/ + } + private final String key; private final PermitImpl root; // this, for a root permit private final Condition condition; - private boolean admitted; - private boolean closed; + private State state = State.WAITING; private PermitImpl(final String key, final PermitImpl parent, final Condition condition) { this.key = key; @@ -323,12 +375,13 @@ private boolean isChild() { public void close() { lock.lock(); try { - if (!admitted || closed) { - return; // idempotent, and a permit that was never granted owns nothing to release + if (state != State.GRANTED) { + return; // idempotent, only GRANTED -> CLOSED transition } - closed = true; - globalInUse--; - perKeyInUse.compute(key, (k, count) -> (count == null || count <= 1) ? null : count - 1); + + state = State.CLOSED; + globalGranted--; + perKeyGranted.compute(key, (k, count) -> (count == null || count <= 1) ? null : count - 1); permitsByRoot.compute(root, (r, count) -> (count == null || count <= 1) ? null : count - 1); logContentionStateIfChanged(); grantNextEligible(); @@ -337,70 +390,89 @@ public void close() { } } + // must be called under lock + private void interrupt() { + if (state != State.WAITING) { + throw new IllegalStateException("Should be called only in WAITING state"); + } + + state = State.CLOSED; + // remove now (earlier, in interrupter thread) to do not wait for this permit thread to get the lock to dequeue (and count it by then) + dequeue(); + condition.signal(); + } + // must be called just after construction, under lock private void enqueue() { - final Deque queue = waitersByKey.computeIfAbsent(key, k -> new ArrayDeque<>()); + final Deque queue = perKeyPending.computeIfAbsent(key, k -> new ArrayDeque<>()); if (queue.isEmpty()) { - waitingKeys.addLast(key); + pendingKeys.addLast(key); } queue.addLast(this); + globalPending++; if (isChild()) { - childWaiters.addLast(this); + pendingChildren.addLast(this); } } - // must be called under lock + // must be called under lock; idempotent - the removal from the key queue is what decides whether this + // permit was still pending, so every other bit of wait state is unwound exactly once with it private void dequeue() { + final Deque queue = perKeyPending.get(key); + if (queue == null || !queue.remove(this)) { + return; // already dequeued by grant() or interrupt() + } + + globalPending--; if (isChild()) { - childWaiters.remove(this); + pendingChildren.remove(this); } - final Deque queue = waitersByKey.get(key); - if (queue != null && queue.remove(this) && queue.isEmpty()) { - waitersByKey.remove(key); - waitingKeys.remove(key); + if (queue.isEmpty()) { + perKeyPending.remove(key); + pendingKeys.remove(key); } } // must be called under lock - private void admit() { + private void grant() { dequeue(); // remove now so it cannot be granted twice // served -> behind every other waiting key - if (waitersByKey.containsKey(key) && waitingKeys.remove(key)) { - waitingKeys.addLast(key); + if (perKeyPending.containsKey(key) && pendingKeys.remove(key)) { + pendingKeys.addLast(key); } - globalInUse++; - perKeyInUse.merge(key, 1, Integer::sum); + globalGranted++; + perKeyGranted.merge(key, 1, Integer::sum); permitsByRoot.merge(root, 1, Integer::sum); logContentionStateIfChanged(); - admitted = true; + state = State.GRANTED; condition.signal(); } // must be called under lock private void logContentionStateIfChanged() { - final int global = globalInUse; - final int capacity = policy.capacity(); - final int threshold = policy.burstThreshold(); - final boolean nowAboveThreshold = threshold > 0 && global >= threshold; - final boolean nowAtCapacity = global >= capacity; + final int global = globalGranted; + final int permits = config.permits(); + final int keyPolicyThreshold = config.keyPolicyThreshold(); + final boolean nowAboveThreshold = keyPolicyThreshold > 0 && global >= keyPolicyThreshold; + final boolean nowAtCapacity = global >= permits; // log threshold crossing (fairness activation) if (nowAboveThreshold && !wasAboveThreshold) { log.warn( - "Throttle contention threshold reached: {}/{} slots in use (threshold={} slots), fairness enforcement active. Active keys: {}", - global, capacity, threshold, perKeyInUse.keySet()); + "Throttle contention threshold reached: {}/{} slots in use (key policy threshold={} slots), fairness enforcement active. Active keys: {}", + global, permits, keyPolicyThreshold, perKeyGranted.keySet()); wasAboveThreshold = true; } else if (!nowAboveThreshold && wasAboveThreshold) { - log.info("Throttle contention below threshold: {}/{} slots in use, fairness relaxed", global, capacity); + log.info("Throttle contention below threshold: {}/{} slots in use, fairness relaxed", global, permits); wasAboveThreshold = false; } // log capacity exhaustion if (nowAtCapacity && !wasAtCapacity) { - log.warn("Throttle capacity exhausted: {}/{} slots in use. Per-key breakdown: {}", global, capacity, perKeyInUse); + log.warn("Throttle capacity exhausted: {}/{} slots in use. Per-key breakdown: {}", global, permits, perKeyGranted); wasAtCapacity = true; } else if (!nowAtCapacity && wasAtCapacity) { - log.info("Throttle capacity freed: {}/{} slots in use", global, capacity); + log.info("Throttle capacity freed: {}/{} slots in use", global, permits); wasAtCapacity = false; } } diff --git a/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/ThrottleProperties.java b/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/ThrottleProperties.java deleted file mode 100644 index 1c97b3793f..0000000000 --- a/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/ThrottleProperties.java +++ /dev/null @@ -1,158 +0,0 @@ -/** - * Copyright (c) 2026 Contributors to the Eclipse Foundation - * - * This program and the accompanying materials are made - * available under the terms of the Eclipse Public License 2.0 - * which is available at https://www.eclipse.org/legal/epl-2.0/ - * - * SPDX-License-Identifier: EPL-2.0 - */ -package org.eclipse.hawkbit.throttle; - -import java.time.Duration; -import java.util.HashMap; -import java.util.Locale; -import java.util.Map; -import java.util.function.ToIntFunction; - -import lombok.Data; -import lombok.extern.slf4j.Slf4j; -import org.eclipse.hawkbit.throttle.Throttle.Policy; -import org.springframework.boot.context.properties.ConfigurationProperties; -import org.springframework.util.LinkedCaseInsensitiveMap; - -/** - * Operator-owned per-tenant throttling configuration. Caps how many of the shared object per domain (e.g. DB connection pool's slots) - * a single tenant may hold at once. Per-tenant overrides live in the {@link ThrottleConfig#tenants} map. - * - *

Configuration Examples: - *

{@code
- * # Mode 1: Disabled (default)
- * hawkbit.throttle.db.enabled=false
- *
- * # Mode 2: Hard cap, fast-reject
- * hawkbit.throttle.db.enabled=true
- * hawkbit.throttle.db.limit=5           # each tenant max 5 connections
- * hawkbit.throttle.db.timeout=0         # reject immediately when over limit
- * hawkbit.throttle.db.threshold=-1      # no threshold (always enforce limit)
- *
- * # Mode 3: Contention-aware burst allowance
- * hawkbit.throttle.db.enabled=true
- * hawkbit.throttle.db.limit=10          # ceiling per tenant
- * hawkbit.throttle.db.timeout=20s       # wait up to 20s before HTTP 429
- * hawkbit.throttle.db.threshold=8       # fairness kicks in at 8/10 pool utilization
- * spring.threads.virtual.enabled=true   # consider it for timeout>0
- * }
- */ -@Data -@ConfigurationProperties("hawkbit") // real prefix is "hawkbit.throttle" but the rest comes from the single value - throttle map -public class ThrottleProperties { - - // throttle config domain to throttle config mapping - private Map throttle = new HashMap<>(); - - @Data - public static class ThrottleConfig { - - /** - * Master switch. Default OFF ⇒ no DataSource wrapping, no filter, byte-for-byte current behavior. - */ - private boolean enabled = false; - - /** - * Total permits. {@code -1} or {@code 0} = auto-detect the Hikari {@code maximumPoolSize}; set explicitly only - * when auto-detect is not possible. - * - *

Because every borrow takes a permit — tenant, system and nested alike — permits map one-to-one onto - * pooled connections, so this should equal the pool size. No headroom is needed: reentrant requests that would - * deadlock are refused rather than parked. - */ - private int capacity = -1; - - /** - * Contention threshold (absolute slot count). When global in-use reaches this threshold, fairness - * enforcement activates: each tenant is served in round-robin, capped by its limit. - * Below the threshold, any tenant may burst freely up to pool capacity. - * {@code -1} = no threshold (always enforce hard cap on tenant limit, 0 or any < -1 is essentially the same). - * Default: -1 (hard cap mode). - */ - private int threshold = -1; - - /** - * Global per-tenant ceiling on concurrent connection borrows. {@code -1} = no ceiling. - */ - private int limit = -1; - - /** - * Capacity that internal, non-tenant work is served ahead of the queue for. Internal work - * is charged to the {@code null} key — schedulers, migrations and health checks run without a tenant context. - * - *

A floor it cannot be starved below, not a cap: beyond it system work competes as an ordinary key, so it - * can neither be crowded out by tenant load nor crowd tenants out itself. Costs no reserved capacity — tenants - * use those slots freely whenever system work is not. {@code 0} disables the priority. Must resolve to fewer - * slots than {@link #capacity}, or tenant work may never be reached; the policy warns if it does not. - */ - private int systemFloor = 1; // default - at least 1 slot with priority for internal work - - /** - * Max wait for a permit before {@link ThrottledException} (→ HTTP 429). {@code 0} fast-rejects - * instead of blocking (safe on platform threads); a positive value is wait-then-serve - * backpressure (requires virtual threads enabled). Default: 0 (safe-first; switch to 20s after - * virtual thread validation). - */ - private Duration timeout = Duration.ZERO; - - /** Operator-only per-tenant limit overrides (tenant id → limit; {@code -1} = opt out). */ - private final LinkedCaseInsensitiveMap tenants = new LinkedCaseInsensitiveMap<>(Locale.ROOT); - - /** - * Resolve the effective limit for a specific tenant. - * - * @param tenant the tenant identifier - * @return the per-tenant ceiling: the tenant override if set in {@link #tenants}, else the global - * {@link #limit}. {@code -1} = unlimited (bounded only by capacity). - */ - public int limit(final String tenant) { - return tenants.getOrDefault(tenant, limit); - } - - public Policy toPolicy(final int capacity) { - return new PolicyImpl(capacity, getThreshold(), this::limit, Math.clamp(systemFloor, 0, capacity)); - } - } - - /** - * Default {@link Policy}: per-key ceilings from a lookup function, and a priority floor granted to internal, non-tenant work — - * the {@code null} key. - * - * @param capacity total permits (should equal the real size of the guarded resource) - * @param burstThreshold in-use count at which fair-share enforcement activates; {@code -1} to always enforce - * @param ceilings key → ceiling on concurrently held permits; {@code -1} = unlimited - * @param systemFloor permits the {@code null} key is served ahead of the queue for; {@code 0} = no priority - */ - @Slf4j - private record PolicyImpl(int capacity, int burstThreshold, ToIntFunction ceilings, int systemFloor) implements Policy { - - public PolicyImpl { - // Priority-floor admissions are served before the fair-share pass, so if the granted floor can cover the whole - // resource there is no arrival order in which a non-priority key is reached. Warn rather than reject: the - // operator may be running a deployment that is genuinely internal-work-only. - if (systemFloor >= capacity) { - log.warn(""" - Throttle system floor ({}) is not below capacity ({}): priority admissions can consume the whole \ - resource, so tenant work may never be served. Lower hawkbit.throttle..system-floor.""", - systemFloor, capacity); - } - } - - @Override - public int ceiling(final String key) { - return ceilings.applyAsInt(key); - } - - @Override - public int priorityFloor(final String key) { - return key == null ? systemFloor : 0; - } - } -} diff --git a/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/ThrottledException.java b/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/ThrottledException.java index 253ba8bc74..d85a748113 100644 --- a/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/ThrottledException.java +++ b/hawkbit-core/src/main/java/org/eclipse/hawkbit/throttle/ThrottledException.java @@ -27,7 +27,38 @@ public class ThrottledException extends AbstractServerRtException { private static final SpServerError THIS_ERROR = SpServerError.SP_THROTTLED; - public ThrottledException(final String message) { + /** + * Why admission was refused. Low cardinality on purpose - it is meant to be a metric tag, so that the rejections + * that mean "the resource is busy" can be told apart from the ones that mean "this deployment is misconfigured". + */ + public enum Reason { + + /** The key is configured with no permits at all, so nothing it asks for can ever be admitted. */ + BLOCKED, + /** No slot free and the caller asked not to wait. */ + LIMIT, + /** No slot free and waiting is disabled for this key by configuration. */ + FAST_REJECT, + /** A nested request that could not be admitted and could not be parked either - parking it would deadlock. */ + DEADLOCK, + /** The key's own pending queue is full and this request is the one to drop. */ + KEY_QUEUE_FULL, + /** The pending queue over all keys is full and this request is the one to drop. */ + QUEUE_FULL, + /** Parked, then dropped to make room for another waiter. */ + EVICTED, + /** Parked, but no slot became available within the timeout. */ + TIMEOUT + } + + private final Reason reason; + + public ThrottledException(final Reason reason, final String message) { super(THIS_ERROR, message); + this.reason = reason; + } + + public Reason getReason() { + return reason; } } diff --git a/hawkbit-core/src/test/java/org/eclipse/hawkbit/throttle/ThrottleTest.java b/hawkbit-core/src/test/java/org/eclipse/hawkbit/throttle/ThrottleTest.java index cd0ff0cccc..733f9b26e8 100644 --- a/hawkbit-core/src/test/java/org/eclipse/hawkbit/throttle/ThrottleTest.java +++ b/hawkbit-core/src/test/java/org/eclipse/hawkbit/throttle/ThrottleTest.java @@ -15,13 +15,17 @@ import java.time.Duration; import java.util.ArrayList; import java.util.Collections; +import java.util.EnumMap; import java.util.List; +import java.util.Map; +import java.util.TreeMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; -import org.eclipse.hawkbit.throttle.ThrottleProperties.ThrottleConfig; +import org.eclipse.hawkbit.throttle.Config.Policy; +import org.eclipse.hawkbit.throttle.Config.Priority; import org.junit.jupiter.api.Test; class ThrottleTest { @@ -396,9 +400,9 @@ void childStillParksWhileAnotherUnitCanRelease() throws Exception { } /** - * Builds the engine through the production {@link ThrottleConfig#toPolicy} path rather than a hand-rolled policy, - * so the per-tenant limit lookup is covered too. Threshold 0 means always contended, i.e. ceilings are enforced - * from the first permit — the interesting mode for every test here. + * Builds the engine through the same key → {@link Policy} wiring the operator-facing configuration produces, + * so the per-tenant limit lookup is covered too rather than a hand-rolled policy. Threshold 0 means always + * contended, i.e. ceilings are enforced from the first permit — the interesting mode for every test here. */ private static Throttle throttle(final int capacity) { return throttle(capacity, config -> { }); @@ -409,10 +413,62 @@ private static Throttle throttle(final int capacity, final Consumer ceiling on concurrently held permits; -1 (and any tenant not listed) falls back to `limit` + private final Map tenants = new TreeMap<>(String.CASE_INSENSITIVE_ORDER); + + private int threshold = -1; + private int limit = -1; // global per-tenant ceiling, -1 = unlimited (bounded only by capacity) + private int systemFloor; + // the engine gates waiting on the configured timeout as well as the per-call one; keep it out of the way + // so that, as in these tests, the value passed to acquire() is what decides + private Duration timeout = Duration.ofMinutes(1); + + private Map getTenants() { + return tenants; + } + + private void setThreshold(final int threshold) { + this.threshold = threshold; + } + + private void setLimit(final int limit) { + this.limit = limit; + } + + private void setSystemFloor(final int systemFloor) { + this.systemFloor = systemFloor; + } + + private void setTimeout(final Duration timeout) { + this.timeout = timeout; + } + + private Config toConfig(final int capacity) { + // null key <=> internal, non-tenant work - the only key granted the priority floor + final Policy systemPolicy = Policy.builder().permits(capacity).timeout(timeout).priority(Priority.HIGH).build(); + return Config.builder() + .permits(capacity) + .timeout(timeout) + .keyPolicyThreshold(Math.min(threshold, capacity)) + .priorityToGranted(new EnumMap<>(Map.of(Priority.HIGH, Math.clamp(systemFloor, 0, capacity)))) + .keyPolicyFn(key -> key == null + ? systemPolicy + : Policy.builder().permits(tenants.getOrDefault(key, limit)).timeout(timeout).build()) + .build(); + } + } } diff --git a/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/management/JpaDeploymentManagement.java b/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/management/JpaDeploymentManagement.java index 4c69013204..bc8a5e4ad9 100644 --- a/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/management/JpaDeploymentManagement.java +++ b/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/management/JpaDeploymentManagement.java @@ -481,7 +481,7 @@ public void startScheduledActionsByRolloutGroupParent(final long rolloutId, fina startScheduledActions0(groupScheduledActions.getContent()); return groupScheduledActions.getTotalElements(); } - }, txManager) > 0) ; + }, txManager) > 0); } @Override @@ -838,8 +838,7 @@ private DistributionSetAssignmentResult assignDistributionSetToTargets( } else { return distributionSetManagement.lock(entityManager.merge(dsValidAndComplete)); } - }, txManager - ); + }, txManager); } else { distributionSet = dsValidAndComplete; } diff --git a/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaAutoAssignHandler.java b/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaAutoAssignHandler.java index a47a366a52..f7a7a7809d 100644 --- a/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaAutoAssignHandler.java +++ b/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaAutoAssignHandler.java @@ -74,19 +74,19 @@ public class JpaAutoAssignHandler implements AutoAssignHandler { private final AutoAssignmentManagement autoAssignmentManagement; private final TargetManagement targetManagement; private final DeploymentManagement deploymentManagement; - private final PlatformTransactionManager transactionManager; + private final PlatformTransactionManager txManager; private final LockRegistry lockRegistry; private final Optional meterRegistry; public JpaAutoAssignHandler( final AutoAssignmentManagement autoAssignmentManagement, final TargetManagement targetManagement, final DeploymentManagement deploymentManagement, - final PlatformTransactionManager transactionManager, final LockRegistry lockRegistry, + final PlatformTransactionManager txManager, final LockRegistry lockRegistry, final Optional meterRegistry) { this.autoAssignmentManagement = autoAssignmentManagement; this.targetManagement = targetManagement; this.deploymentManagement = deploymentManagement; - this.transactionManager = transactionManager; + this.txManager = txManager; this.lockRegistry = lockRegistry; this.meterRegistry = meterRegistry; } @@ -269,7 +269,7 @@ private int runTransactionalAssignment(final AutoAssignment autoAssignment, fina () -> deploymentManagement.assignDistributionSets(deploymentRequests, actionMessage)); } return count; - }, transactionManager); + }, txManager); } /** diff --git a/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaRolloutExecutor.java b/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaRolloutExecutor.java index c42115b534..fd1560e4d6 100644 --- a/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaRolloutExecutor.java +++ b/hawkbit-repository/hawkbit-repository-jpa/src/main/java/org/eclipse/hawkbit/repository/jpa/scheduler/JpaRolloutExecutor.java @@ -558,14 +558,14 @@ private JpaRolloutGroup fillRolloutGroupWithTargets( final long targetsInGroupFilter; if (!RolloutHelper.isRolloutRetried(rollout.getTargetFilterQuery())) { // default case targetsInGroupFilter = DeploymentHelper.runInNewTransaction( - "countByRsqlAndNotInRolloutGroupsAndCompatibleAndUpdatable", count -> countByRsqlAndNotInRolloutGroupsAndCompatibleAndUpdatable( - groupTargetFilter, readyGroups, rollout.getDistributionSet().getTypeId()), txManager - ); + "countByRsqlAndNotInRolloutGroupsAndCompatibleAndUpdatable", + count -> countByRsqlAndNotInRolloutGroupsAndCompatibleAndUpdatable( + groupTargetFilter, readyGroups, rollout.getDistributionSet().getTypeId()), txManager); } else { // if it is a rollout retry targetsInGroupFilter = DeploymentHelper.runInNewTransaction( - "countByFailedRolloutAndNotInRolloutGroupsAndCompatible", count -> countByFailedRolloutAndNotInRolloutGroups( - RolloutHelper.getIdFromRetriedTargetFilter(rollout.getTargetFilterQuery()), readyGroups), txManager - ); + "countByFailedRolloutAndNotInRolloutGroupsAndCompatible", + count -> countByFailedRolloutAndNotInRolloutGroups( + RolloutHelper.getIdFromRetriedTargetFilter(rollout.getTargetFilterQuery()), readyGroups), txManager); } final double percentFromTheRest; @@ -577,8 +577,7 @@ private JpaRolloutGroup fillRolloutGroupWithTargets( final long expectedInGroup = Math.round(percentFromTheRest * targetsInGroupFilter / 100); long targetsLeftToAdd = expectedInGroup - DeploymentHelper.runInNewTransaction( - "countRolloutTargetGroupByRolloutGroup", count -> rolloutTargetGroupRepository.countByRolloutGroup(group), txManager - ); + "countRolloutTargetGroupByRolloutGroup", count -> rolloutTargetGroupRepository.countByRolloutGroup(group), txManager); try { while (targetsLeftToAdd > 0) { // Add up to TRANSACTION_TARGETS of the left targets. In case a TransactionException is thrown this loop aborts diff --git a/hawkbit-repository/hawkbit-repository-jpa/src/test/java/org/eclipse/hawkbit/repository/jpa/scheduler/AutoAssignHandlerTest.java b/hawkbit-repository/hawkbit-repository-jpa/src/test/java/org/eclipse/hawkbit/repository/jpa/scheduler/AutoAssignHandlerTest.java index 1eea955674..4bcc97f22a 100644 --- a/hawkbit-repository/hawkbit-repository-jpa/src/test/java/org/eclipse/hawkbit/repository/jpa/scheduler/AutoAssignHandlerTest.java +++ b/hawkbit-repository/hawkbit-repository-jpa/src/test/java/org/eclipse/hawkbit/repository/jpa/scheduler/AutoAssignHandlerTest.java @@ -55,7 +55,7 @@ class AutoAssignHandlerTest { @Mock private DeploymentManagement deploymentManagement; @Mock - private PlatformTransactionManager transactionManager; + private PlatformTransactionManager txManager; @Mock LockRegistry lockRegistry; @@ -65,7 +65,7 @@ class AutoAssignHandlerTest { @BeforeEach void before() { autoAssignHandler = new JpaAutoAssignHandler( - autoAssignmentManagement, targetManagement, deploymentManagement, transactionManager, lockRegistry, Optional.empty()); + autoAssignmentManagement, targetManagement, deploymentManagement, txManager, lockRegistry, Optional.empty()); } /**