diff --git a/documents/ADRs/adr-010-hybrid-jmx-otel-driver-metrics.md b/documents/ADRs/adr-010-hybrid-jmx-otel-driver-metrics.md new file mode 100644 index 000000000..b4b85820e --- /dev/null +++ b/documents/ADRs/adr-010-hybrid-jmx-otel-driver-metrics.md @@ -0,0 +1,60 @@ +# ADR 010: Hybrid JMX core + opt-in OpenTelemetry adapter for driver-side metrics + +In the context of the OJP JDBC driver (`ojp-jdbc-driver`), +facing the need to expose client-side throttling metrics (in-flight, proactive/reactive/effective limits, acquired/rejected counters, server-overload events, AIMD limit changes) without forcing observability dependencies onto every host application that uses the driver, + +we decided for a **hybrid approach**: ship **JMX MBeans in the driver core** (zero new runtime dependencies) and publish a **separate optional adapter module** (`ojp-jdbc-driver-otel-metrics`) that maps the same metrics onto OpenTelemetry for users who want native dimensional metrics. +A small internal abstraction (`ClientThrottleMetrics` interface, with `NoOp`, `Jmx`, and — later — `OpenTelemetry` implementations selected by `ClientThrottleMetricsFactory` via the `ojp.jdbc.metrics` system property) lets the binding be chosen at runtime, mirroring the server-side `SqlStatementMetrics` / `NoOpSqlStatementMetrics` / `OpenTelemetrySqlStatementMetrics` pattern already in `ojp-server`. + +We neglected: +- **JMX-only.** Rejected because it forces operators on an OTel pipeline to bolt on the `jmx_exporter` agent plus YAML rules and gives up native dimensional labels, histograms, and trace exemplars. +- **OpenTelemetry-API-only in the driver core.** Rejected because adding `io.opentelemetry:opentelemetry-api` as a runtime dependency of a JDBC driver creates real classpath / version-skew risk (estimated ~70% chance of biting at least one user during the 0.x lifecycle), introduces the silent no-op trap when a host has the API but no SDK, and is essentially irreversible once shipped under our metric names. Shading the API would defeat the integration value that motivates choosing OTel in the first place. +- **Micrometer in the driver core.** Rejected for the same dependency-cost reason; Micrometer is a strictly larger commitment than the OTel API. + +to achieve: +- **Zero new mandatory runtime dependencies** for the driver — it stays a polite guest in any host JVM. +- **Universal out-of-the-box consumption** of metrics via JConsole / VisualVM / `jmx_exporter` / any APM agent that already scrapes JMX. +- **Native OpenTelemetry integration as an opt-in** for users who want it, without paying any cost when they don't. +- **A reversible, low-risk Phase 1** whose architectural artefact (the `ClientThrottleMetrics` interface) is exactly what an eventual Phase 2 OTel adapter needs. +- **Cultural alignment** with how HikariCP (the de facto reference for JDBC-driver-adjacent metrics) evolved: JMX in core, Micrometer / Dropwizard as separately-published optional adapters. + +accepting: +- **Two implementations to maintain** (JMX and, in Phase 2, OTel). This cost is bounded by the small surface of the `ClientThrottleMetrics` interface. +- **JMX exposes scalars only.** The proposed driver metrics in `documents/analysis/THROTTLING_METRICS_ANALYSIS.md` §5 are all gauges and counters — no histograms — so this is not currently a functional limitation. +- **Operators on Prometheus via JMX need a `jmx_exporter` configuration**. We will ship a sample rules YAML alongside the documentation. +- **MBean lifecycle is real work.** Registration must be idempotent per `connHash`, unregister-on-close must be best-effort, and duplicate-name registration must be tolerated. Encapsulated in `JmxClientThrottleMetrics`. +- **Phase 2 timing is not committed by this ADR.** The OTel adapter module will be authored under a follow-up issue once Phase 1 adoption produces a real signal that it is wanted. + +because: +- The driver lives in the host application's classpath, where the dominant operational risk is dependency conflict, not metric expressiveness; JMX has zero such risk. +- The hybrid keeps the door open for OpenTelemetry-native consumption without forcing it on anyone. +- The internal interface makes the decision reversible: if the maintainers later prefer to consolidate on OTel, the JMX implementation can be deprecated without touching the throttle-management code paths. + +## Scope + +- **In scope.** Driver-side client throttling metrics defined in `THROTTLING_METRICS_ANALYSIS.md` §5 (`ojp.client.throttle.*`). +- **Out of scope.** Server-side `ojp.throttle.*` metrics (gRPC interceptor and `SlotManager`); the server already standardises on OpenTelemetry per ADR-005 and that does not change here. + +## Configuration + +A single system property controls the binding: + +``` +-Dojp.jdbc.metrics=jmx # default — register JMX MBeans +-Dojp.jdbc.metrics=none # disable; use NoOpClientThrottleMetrics +-Dojp.jdbc.metrics=otel # Phase 2 — only available with ojp-jdbc-driver-otel-metrics on the classpath +``` + +If `otel` is requested but the OTel adapter is not on the classpath, the factory falls back to `jmx` and logs a single warning. + +## Rollout + +1. **Phase 1 (this ADR).** Land `ClientThrottleMetrics` interface, `NoOp` and `JMX` implementations, factory, and wiring in `ClientThrottleManager` / `Connection`. Default: `jmx`. No new runtime dependencies in `ojp-jdbc-driver`. **Status: implemented.** +2. **Phase 2 (this ADR).** `ojp-jdbc-driver-otel-metrics` Maven module: depends on `opentelemetry-api` (scope `provided`), implements `ClientThrottleMetrics` against `GlobalOpenTelemetry.get()`, and registers an `OpenTelemetryClientThrottleMetricsProvider` via `META-INF/services` so the driver-core factory discovers it through `ServiceLoader` when `-Dojp.jdbc.metrics=otel` is set. Users who want OTel add the adapter jar; users who don't, pay nothing. **Status: implemented.** See `ojp-jdbc-driver-otel-metrics/`. + +| Status | PROPOSED | +|---------------|---------------------| +| Proposer(s) | GitHub Copilot, on direction from Rogerio Robetti | +| Proposal date | 23/05/2026 | +| Approver(s) | | +| Approval date | | diff --git a/documents/analysis/README.md b/documents/analysis/README.md index ccc08d952..b9b89425a 100644 --- a/documents/analysis/README.md +++ b/documents/analysis/README.md @@ -4,6 +4,17 @@ This directory contains technical analysis documents for various OJP features an ## Latest Analysis (May 2026) +### 🆕 Throttling Metrics Exposure + +**Question:** OJP throttles requests in three independent layers but exposes no metrics. What should be exposed? + +**Quick Answer:** Introduce an `ojp.throttle.*` OpenTelemetry meter scope covering the global gRPC gate, per-datasource `SlotManager` (fast/slow lanes, queue depth, observed peak, AIMD), and the driver-side `ClientThrottleManager` (in-flight, proactive/reactive limits, server-overload events). Driver-side transport (JMX vs OTel API) is left as an open decision. + +**Document:** [THROTTLING_METRICS_ANALYSIS.md](./THROTTLING_METRICS_ANALYSIS.md) + +--- + + ### 🆕 Prepared Statement Cache Server Settings Design **Question:** How should OJP expose default-enabled prepared statement caching using standard OJP server properties and runtime datasource translation? diff --git a/documents/analysis/THROTTLING_METRICS_ANALYSIS.md b/documents/analysis/THROTTLING_METRICS_ANALYSIS.md new file mode 100644 index 000000000..7624f6dd5 --- /dev/null +++ b/documents/analysis/THROTTLING_METRICS_ANALYSIS.md @@ -0,0 +1,394 @@ +# Throttling Metrics Analysis + +**Status:** Analysis only — no implementation yet, awaiting approval. +**Author:** Investigation requested 2026-05-23. +**Scope:** Identify which metrics OJP should expose about its existing throttling +behaviour so operators can observe, alert on, and tune it. + +--- + +## 1. Why this analysis exists + +OJP enforces back-pressure in **three independent layers**, but none of them +currently emit metrics. The only observable signals today are `log.warn` / +`log.debug` lines and the human-readable `SlotManager.getStatus()` string. + +Consequence: when OJP starts rejecting work, an operator cannot tell: + +- whether requests are being shed at all, +- which layer is doing the shedding, +- whether the configured limits are too tight or too loose, +- whether the AIMD feedback loops are converging or oscillating. + +This document proposes a concrete metrics surface that closes that gap, and +flags the design decisions that need a human call before implementation begins. + +**Confidence in the gap analysis: High** — verified by reading each class +directly; no metric registration exists in any of them. + +--- + +## 2. The three throttling layers today + +| Layer | Class | Scope | Rejection signal | +|---|---|---|---| +| **A. Global gRPC concurrency gate** | `ConcurrencyThrottleInterceptor` (server) | Server-wide, all in-flight gRPC calls | Closes call with `RESOURCE_EXHAUSTED` "too many concurrent requests" | +| **B. Admission control / SQS slot manager** | `SlotManager` (server) | Per-datasource, fast/slow lanes + borrowing + AIMD `observedPeak` | `acquireFastSlot` / `acquireSlowSlot` return `false` after timeout or when wait-queue cap is reached | +| **C. Client-side reactive throttle** | `ClientThrottleManager` (driver) | Per JVM, per `connHash` | Local fail-fast via CAS on an `AtomicInteger`; no network round-trip | + +The layers compose: a request can be rejected at C (cheapest, no network), at A +(once it reaches the server), or at B (once a session is established and a slot +is needed for a query). Healthy steady state should have C dominating +rejections, with A and B contributing only during true overload. + +--- + +## 3. Operator-facing questions the metrics must answer + +Each proposed metric is justified by which of these questions it answers: + +1. *Is OJP rejecting requests right now? Which layer?* +2. *How close to capacity am I — scale OJP horizontally, or raise pool size?* +3. *Is SQS doing its job (fast lane stays fast even when slow queries pile up)?* +4. *Is the client-side throttle protecting the server, or over-tightening?* +5. *Is AIMD converging on a stable `observedPeak` / `reactiveLimit`, or oscillating?* +6. *Per datasource / per database — which is the hot spot?* + +--- + +## 4. Proposed metrics — Server side + +All metrics live under a new OpenTelemetry meter scope **`ojp.throttle`**, +following the existing `ojp.sql`, `ojp.hikari.pool`, `ojp.xa.pool` convention. +Implementation should mirror the existing pattern: interface + +OpenTelemetry impl + NoOp impl + factory (see +`OpenTelemetrySqlStatementMetrics`, `OpenTelemetryPoolMetrics`). + +### 4.1 Global gRPC throttle (`ConcurrencyThrottleInterceptor`) + +| Metric | Type | Attributes | Question answered | +|---|---|---|---| +| `ojp.throttle.grpc.inflight` | Gauge | – | Live concurrency level | +| `ojp.throttle.grpc.limit` | Gauge | – | Configured limit visibility | +| `ojp.throttle.grpc.utilization` | Gauge (0–100) | – | One-glance "how close to the wall" | +| `ojp.throttle.grpc.rejected.total` | Counter | `grpc.method` | Are we shedding load? Which RPCs? | +| `ojp.throttle.grpc.accepted.total` | Counter | `grpc.method` | Denominator for rejection ratio | +| `ojp.throttle.grpc.hold.time` | Histogram (ms) | – | Detect "blocked on backend" vs "burst arrivals" | + +Rejection ratio in PromQL: `rate(rejected) / (rate(rejected) + rate(accepted))`. + +### 4.2 Admission control / SlotManager (per datasource) + +Attributes carry `datasource` and, where meaningful, `lane = fast|slow`. + +| Metric | Type | Attributes | Question answered | +|---|---|---|---| +| `ojp.throttle.slots.total` | Gauge | `datasource`, `lane` | Configured lane size | +| `ojp.throttle.slots.active` | Gauge | `datasource`, `lane` | Live in-use count | +| `ojp.throttle.slots.available` | Gauge | `datasource`, `lane` | `semaphore.availablePermits()` | +| `ojp.throttle.slots.utilization` | Gauge 0–100 | `datasource`, `lane` | Single-pane lane saturation | +| `ojp.throttle.slots.queue.depth` | Gauge | `datasource`, `lane` | `semaphore.getQueueLength()` — earliest warning, rises *before* timeouts | +| `ojp.throttle.slots.queue.max` | Gauge | `datasource` | Configured `maxWaitQueueDepth` | +| `ojp.throttle.slots.acquired.total` | Counter | `datasource`, `lane`, `path=immediate\|wait\|borrowed` | Fast-path hit rate; lane borrowing frequency | +| `ojp.throttle.slots.rejected.total` | Counter | `datasource`, `lane`, `reason=timeout\|queue_full` | Why are we shedding? Tune queue vs slots | +| `ojp.throttle.slots.wait.time` | Histogram (ms) | `datasource`, `lane` | p95/p99 admission latency | +| `ojp.throttle.slots.borrowed` | Gauge | `datasource`, `direction=slow_to_fast\|fast_to_slow` | Are lanes well-sized? | +| `ojp.throttle.slots.observedPeak` | Gauge | `datasource` | AIMD-tracked peak sent to clients | +| `ojp.throttle.slots.enabled` | Gauge (1/0) | `datasource` | Avoid "all clear but feature was off" trap | + +> **Note:** `observedPeak` is the value the server ships to clients to size +> their throttles. Exposing it makes the otherwise-invisible AIMD loop +> observable — essential for confidence in the auto-tuning. + +--- + +## 5. Proposed metrics — Driver / client side + +The driver has no metrics infrastructure today. Two pure options and one +hybrid; see **Appendix §9** for the full deep dive. + +- **Option A: JMX MBeans.** Zero new driver dependency, universal consumer + support (JConsole, `jmx_exporter`, any APM agent). Scalars only; needs + exporter config for Prometheus. +- **Option B: OpenTelemetry API in the driver.** Native dimensional model + and histograms, but introduces a real classpath dependency for every host + application — including those that don't use OTel — and creates a non-zero + version-skew risk. +- **Option C (recommended): Hybrid.** JMX in the driver core, plus an + optional `ojp-jdbc-driver-otel-metrics` adapter jar published separately. + Mirrors HikariCP's history (JMX core + opt-in Micrometer/Dropwizard + adapters). Decision is reversible and Phase 2 is gated on real adoption + signal. + +Metrics themselves (regardless of transport), attributes `connHash`, `mode`: + +| Metric | Type | Question answered | +|---|---|---| +| `ojp.client.throttle.inflight` | Gauge | Live in-flight at the driver | +| `ojp.client.throttle.proactiveLimit` | Gauge | Current proactive limit (from `SessionInfo`) | +| `ojp.client.throttle.reactiveLimit` | Gauge | Current reactive limit (AIMD) | +| `ojp.client.throttle.effectiveLimit` | Gauge | `min(proactive, reactive)` actually enforced | +| `ojp.client.throttle.rejected.total` | Counter | Driver-local fail-fast count (no server round-trip) | +| `ojp.client.throttle.acquired.total` | Counter | Denominator for ratio | +| `ojp.client.throttle.serverOverload.events.total` | Counter | Every `notifyServerOverload()` (RESOURCE_EXHAUSTED) | +| `ojp.client.throttle.limitChanges.total` | Counter, `direction=increase\|decrease` | AIMD stability — frequent flips → unstable cluster sizing | + +Healthy steady state: `ojp.client.throttle.rejected.total` dominates +`ojp.throttle.grpc.rejected.total` and `ojp.throttle.slots.rejected.total`, +because client-side rejection avoids a network round-trip. + +--- + +## 6. Cross-cutting recommendations + +1. **Reuse existing `OpenTelemetryHolder`** that already feeds `ojp.sql` and + `ojp.*.pool` scopes — do not introduce a second OTel bootstrap path. +2. **NoOp impl by default** in unit tests, matching `NoOpSqlStatementMetrics` / + `NoOpPoolMetrics`. Throttle logic must not be coupled to OTel directly. +3. **Cardinality discipline.** Allowed labels: `datasource`, `lane`, + `grpc.method`, `mode`, `path`, `reason`, `direction`. Disallowed: + `client.uuid`, `session.id`, `sql.statement` — these would explode + cardinality. The proposed labels are all bounded by configuration or by a + small fixed enum. +4. **Document a recommended alerting set** in + `documents/telemetry/README.md`. At minimum: + - sustained `utilization > 80%` → scale out; + - any non-zero `rejected` rate → page; + - `observedPeak < totalSlots * 0.5` for >10 minutes → DB is the + bottleneck, not OJP. +5. **Reduce `SlotManager.getStatus()` log spam to DEBUG** once equivalent + metrics exist — the periodic INFO line becomes redundant. + +--- + +## 7. Suggested implementation sequence + +1. **Layer B — SlotManager.** Biggest operator value, no new transport, fits + the existing OTel scope pattern. Ship first. +2. **Layer A — gRPC interceptor.** Small surface; counters and gauges around + the existing `AtomicInteger`. +3. **Layer C — driver.** Last, with its own ADR for the JMX-vs-OTel-API + transport decision. + +Each layer can ship independently; each is observable from PromQL the same day. + +--- + +## 8. Open decisions (need approval before implementing) + +- **Driver-side transport.** Decided in **ADR-010**: hybrid (JMX in driver + core + opt-in `ojp-jdbc-driver-otel-metrics` adapter module). Phase 1 (JMX + + `ClientThrottleMetrics` interface + `ClientThrottleMetricsFactory`) is + implemented in this PR. Phase 2 (OTel adapter module) deferred to a + follow-up issue. +- **`datasource` attribute identity.** The hash (`connHash`) is already + available in code; a human-friendly datasource name would require config + plumbing. Recommend starting with the hash and revisiting once needs are + clearer. +- **`grpc.method` cardinality.** ~10 RPC methods today; confirm this label is + acceptable before adopting it on `ojp.throttle.grpc.*`. +- **Overlap with existing `ojp.sql.slow.executions.total`.** Intentionally + kept `lane` attribution only on *slot* metrics (admission concern); SQL + statement metrics keep their existing shape. Confirm this split is right. + +--- + +## 9. Appendix — Deep dive: JMX vs OpenTelemetry for **driver-side** metrics + +> Scope: this appendix is **driver-only**. The server already uses +> OpenTelemetry and that is not in question. The asymmetry matters: a JDBC +> driver is loaded into someone else's JVM, so it must be a polite guest in a +> way the server never has to be. + +### 9.1 Baseline facts that frame the decision + +- **Driver has zero OTel dependency today** (`grep opentelemetry + ojp-jdbc-driver/pom.xml` → 0; server → 21). Adding OTel is therefore an + *introduction*, not an extension. +- **Driver minimum runtime is Java 11**, server is Java 21. Anything pulled + into the driver must compile and run on 11. +- **The driver ships into the host application's classpath.** It does not + control which JVM, which OTel version (if any), or which exporter the host + application uses. +- **Driver metric volume is tiny** (~8 series per `connHash`, no per-query + cardinality). The technical demand is well below what JMX can handle. + +### 9.2 Option A — JMX MBeans + +**How it would work.** Register one `ThrottleManagerMXBean` per `connHash` +under e.g. `org.openjproxy:type=ClientThrottle,connHash=` with attributes +mirroring §5 (inflight, limits, counters). No new runtime dependency. Host +apps consume via: +- JConsole / VisualVM out of the box, +- `jmx_exporter` Java agent → Prometheus, +- any APM agent (New Relic, AppDynamics, Datadog) that already scrapes JMX + (they all do). + +**Benefits** +- **Zero new dependencies.** Stays inside `java.management`. No classloader, + no shading, no version-skew risk for host applications. +- **Universally consumable.** Every Java monitoring tool understands JMX. + Operators don't need to install or configure anything OJP-specific. +- **Java 11 friendly with no caveats.** No multi-release jar concerns, no + conditional bytecode. +- **Cheap to remove or evolve.** MBeans are a thin façade over the existing + `AtomicInteger`s; if we change our mind in 6 months, deletion is trivial and + affects no public Java API. +- **Matches the cultural norm for JDBC drivers.** HikariCP exposes pool + metrics via JMX by default (it added Micrometer/Dropwizard as *optional* + adapters). Operators expect this. +- **No risk of "double OTel".** If the host app already uses OTel with a + different version, OJP's metrics still work — there is no classpath + collision to resolve. + +**Implications and downsides** +- **Pull-only and rate-unaware by itself.** JMX exports raw values; rates + (`rate(rejected[5m])`) happen on the scraper side. This is exactly what + `jmx_exporter` is built for, but it is one more moving part operators must + configure. +- **Naming / hierarchy maps awkwardly to dimensional metrics.** Prometheus + labels via `jmx_exporter` require a regex-based YAML config; operators must + copy or be given a sample rules file. We should ship one. +- **No histograms natively.** JMX attributes are scalars. We would expose + `count` + `sum` + a few percentiles computed locally (e.g. via HdrHistogram) + if needed. For the proposed driver metric set this is acceptable because + there are no histograms — wait time lives only on the server side. +- **No trace correlation.** JMX cannot attach exemplars or link a metric + point to a trace span. For a driver that does no tracing today, this is + zero practical loss. +- **MBean lifecycle = a small footgun.** We must unregister on + `Connection`/throttle-manager teardown, and re-register safely if the same + `connHash` reappears. A static set of registered ObjectNames plus a JVM + shutdown hook handles this; it is a known pattern but real work. +- **Security manager / module-system friction is essentially nil today** but + worth noting: `java.management` is a base module; no `--add-opens` + gymnastics. + +### 9.3 Option B — OpenTelemetry **API** in the driver + +> Important distinction: **API**, not **SDK**. The driver would call +> `GlobalOpenTelemetry.get().meterBuilder(...)`, and the host application +> supplies the SDK and exporter. If the host app has no SDK installed, OTel's +> API returns a no-op meter and nothing breaks. + +**Benefits** +- **Native dimensional model.** Attributes (`connHash`, `mode`, …) map 1:1 to + the proposed metric design. No YAML translation rules to maintain. +- **Histograms first-class.** If we later want client-side wait-time + histograms (we currently don't), they are free. +- **Trace correlation / exemplars.** A future "trace each throttled + rejection" feature lands without extra plumbing. +- **One observability stack across server and driver.** A single dashboard + template, one mental model for operators who are already on OTel. +- **Modern, where the industry is heading.** Most new-issue observability work + in Java assumes OTel. + +**Implications and downsides** +- **Real dependency on the host application's classpath.** Even just the API + jar (`io.opentelemetry:opentelemetry-api`) means: + - One more jar in apps that don't use OTel. + - **Version-skew risk** is the dominant concern. If the host already pulls + a different OTel API version, the usual "newest wins" classpath rule + applies. OTel's API has been stable since 1.0 and tries hard to avoid + breaks, but transitive dependency conflicts are the #1 reason JDBC + drivers get blamed for "weird" startup failures. *Confidence this will + bite us at least once: ~70%.* + - We may eventually need to **shade and relocate** the OTel API + (`org.openjproxy.shaded.io.opentelemetry`) to be safe — but that breaks + the "metrics flow into the host's OTel pipeline" property, which is the + main reason to choose OTel in the first place. Shading defeats the + purpose; not shading risks conflicts. There is no clean middle. +- **Silent no-op trap.** If the host application has the API but not an SDK + installed (a very common state — many apps include OTel transitively + without configuring it), `GlobalOpenTelemetry.get()` returns a no-op meter + and OJP metrics silently disappear. Operators will file bugs. +- **API stability is *good*, not perfect.** The metrics API was incubating + for longer than tracing. We must pin a minimum version (≥ 1.31 is safe) + and document it. +- **Larger driver jar / longer cold start.** Marginal (~200 KB, a few classes + loaded), but JDBC drivers compete on smallness. +- **Java 11 baseline is fine** for current OTel API versions, but we lose + some freedom — if upstream raises its baseline to 17, we either pin an + older OTel or raise the driver's baseline. Either is awkward. +- **Harder to remove.** Once shipped, removing the OTel dependency is a + breaking change for any user who wired exporters to our metric names. + +### 9.4 Hybrid — JMX now, optional OTel adapter later + +A pragmatic third path: + +1. **Phase 1:** Ship JMX. Counters/gauges live behind a tiny + `ClientThrottleMetrics` interface with a `JmxClientThrottleMetrics` + implementation, plus a `NoOp` (same pattern as + `OpenTelemetrySqlStatementMetrics` / `NoOpSqlStatementMetrics` on the + server side). +2. **Phase 2 (optional, separate release):** Add a *separate Maven + coordinate* — e.g. `ojp-jdbc-driver-otel-metrics` — that depends on the + OTel API and implements the same interface. Users who want OTel add the + adapter jar; users who don't, don't pay any cost. No shading needed + because users opt in. No version conflicts for non-OTel users. + +**Why this is attractive** +- The hard architectural decision (`ClientThrottleMetrics` interface) is + needed under *any* option. Doing it once unlocks both. +- Mirrors HikariCP's history: JMX in core; Micrometer/Dropwizard as + separately-published adapters that nobody is forced to depend on. +- Lets the maintainers see real adoption signal before committing the driver + to an OTel dependency. +- Reversible. Phase 2 can be cancelled with zero impact on Phase 1 users. + +### 9.5 Comparison summary + +| Dimension | JMX | OTel API in driver | Hybrid (JMX core + adapter) | +|---|---|---|---| +| New driver dependency | None | `opentelemetry-api` | None (adapter is optional jar) | +| Risk of host-classpath conflicts | None | Medium-High | None for core; opt-in for adapter | +| Out-of-the-box for OTel users | Needs `jmx_exporter` or APM agent | Native | Native (with adapter) | +| Out-of-the-box for non-OTel users | Native (JConsole, APMs) | They pay the dep cost anyway | Native | +| Dimensional labels | Via `jmx_exporter` rules | Native | Both, depending on jar | +| Histograms | Scalars only (fine for §5) | Native | Both, depending on jar | +| Trace exemplars | No | Yes | Adapter-only | +| Driver jar size impact | ~0 | ~200 KB | ~0 for core | +| Reversibility | High | Low | High | +| Cultural fit for a JDBC driver | Strong (Hikari precedent) | Mixed | Strong + future-proof | +| Implementation effort | Low | Low | Low + Low (separate module) | + +### 9.6 Recommendation + +**Adopt the hybrid.** Concretely: + +1. Introduce `org.openjproxy.jdbc.metrics.ClientThrottleMetrics` in + `ojp-jdbc-driver` with `NoOp` and `Jmx` implementations. +2. Wire `ClientThrottleManager` and `Connection` to the interface; default + binding selected by a system property (`ojp.jdbc.metrics=jmx|none`, + default `jmx`). +3. Open a follow-up issue to evaluate an `ojp-jdbc-driver-otel-metrics` + adapter module after one minor release of real-world usage. + +**Confidence: ~80% (High).** The main residual uncertainty is whether enough +users will ask for native OTel in the driver to justify Phase 2; that +question answers itself with telemetry-of-the-telemetry over the next couple +of releases. + +### 9.7 Things that would *change* the recommendation + +- If maintainers already plan to add tracing to the driver, jumping straight + to OTel becomes more defensible — counters and spans want to share the + same context. +- If we want exemplar-linked debugging (`rejected.total{trace_id=...}`) as a + first-class operator feature, OTel becomes strongly preferred. +- If the project decides to raise the driver's minimum Java to 17+, the + classpath-conflict risk shrinks (more apps will have aligned to modern OTel + baselines), nudging toward straight OTel. + +--- + +## 10. Code references + +- `ojp-server/src/main/java/org/openjproxy/grpc/server/ConcurrencyThrottleInterceptor.java` +- `ojp-server/src/main/java/org/openjproxy/grpc/server/SlotManager.java` +- `ojp-server/src/main/java/org/openjproxy/grpc/server/GrpcServer.java` (interceptor wiring) +- `ojp-server/src/main/java/org/openjproxy/grpc/server/metrics/` (existing OTel metrics pattern to follow) +- `ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/ClientThrottleManager.java` +- `documents/telemetry/README.md` (where the new metrics should be documented) +- `documents/analysis/CLIENT_REACTIVE_THROTTLING_ANALYSIS.md` (related background) diff --git a/ojp-jdbc-driver-otel-metrics/pom.xml b/ojp-jdbc-driver-otel-metrics/pom.xml new file mode 100644 index 000000000..7c2756431 --- /dev/null +++ b/ojp-jdbc-driver-otel-metrics/pom.xml @@ -0,0 +1,148 @@ + + + 4.0.0 + + OJP JDBC Driver OpenTelemetry Metrics Adapter + ojp-jdbc-driver-otel-metrics + 0.4.18-SNAPSHOT + + Optional OpenTelemetry adapter for the OJP JDBC driver's client-side + throttling metrics. Drop this jar on the classpath and set + -Dojp.jdbc.metrics=otel to publish ojp.client.throttle.* metrics through + the host application's OpenTelemetry pipeline. See ADR-010. + + + + org.openjproxy + ojp-parent + 0.4.18-SNAPSHOT + ../pom.xml + + + + 1.62.0 + + + + + + org.openjproxy + ojp-jdbc-driver + 0.4.18-SNAPSHOT + + + + + io.opentelemetry + opentelemetry-api + ${opentelemetry.version} + provided + + + + + org.slf4j + slf4j-api + ${slf4j.version} + provided + + + + + org.junit.jupiter + junit-jupiter + 5.14.3 + test + + + io.opentelemetry + opentelemetry-sdk + ${opentelemetry.version} + test + + + io.opentelemetry + opentelemetry-sdk-testing + ${opentelemetry.version} + test + + + org.slf4j + slf4j-simple + ${slf4j.version} + test + + + + + + + org.apache.maven.plugins + maven-surefire-plugin + 3.2.5 + + + **/*Test.java + + + + + org.apache.maven.plugins + maven-source-plugin + 3.2.1 + + + attach-sources + + jar + + + + + + org.apache.maven.plugins + maven-javadoc-plugin + 3.6.3 + + false + none + + + + attach-javadocs + + jar + + + + + + + + + https://github.com/Open-J-Proxy/ojp + scm:git:https://github.com/Open-J-Proxy/ojp.git + scm:git:ssh://git@github.com:Open-J-Proxy/ojp.git + HEAD + + + + + The Apache License, Version 2.0 + http://www.apache.org/licenses/LICENSE-2.0.txt + repo + + + + + + rrobetti + Rogerio Robetti + + + + https://github.com/Open-J-Proxy/ojp + + diff --git a/ojp-jdbc-driver-otel-metrics/src/main/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetrics.java b/ojp-jdbc-driver-otel-metrics/src/main/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetrics.java new file mode 100644 index 000000000..5a519f8e8 --- /dev/null +++ b/ojp-jdbc-driver-otel-metrics/src/main/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetrics.java @@ -0,0 +1,198 @@ +package org.openjproxy.jdbc.metrics.otel; + +import io.opentelemetry.api.GlobalOpenTelemetry; +import io.opentelemetry.api.OpenTelemetry; +import io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.metrics.LongCounter; +import io.opentelemetry.api.metrics.Meter; +import io.opentelemetry.api.metrics.ObservableLongGauge; +import lombok.extern.slf4j.Slf4j; +import org.openjproxy.jdbc.metrics.ClientThrottleMetrics; +import org.openjproxy.jdbc.metrics.ClientThrottleStateProvider; + +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * OpenTelemetry-backed {@link ClientThrottleMetrics}. + * + *

Publishes the following instruments under meter scope + * {@code ojp.client.throttle}, all carrying the attribute + * {@code ojp.connection.hash}:

+ * + * + * + *

Selected by {@code -Dojp.jdbc.metrics=otel} when this adapter is on the + * classpath. See ADR-010.

+ */ +@Slf4j +public final class OpenTelemetryClientThrottleMetrics implements ClientThrottleMetrics { + + static final String METER_NAME = "ojp.client.throttle"; + static final AttributeKey CONN_HASH_KEY = AttributeKey.stringKey("ojp.connection.hash"); + static final AttributeKey DIRECTION_KEY = AttributeKey.stringKey("direction"); + + /** One-shot guard for the "GlobalOpenTelemetry is noop" warning (process-wide). */ + private static final AtomicBoolean NOOP_WARNED = new AtomicBoolean(false); + + private final Attributes baseAttrs; + private final Attributes increaseAttrs; + private final Attributes decreaseAttrs; + + private final LongCounter acquiredCounter; + private final LongCounter rejectedCounter; + private final LongCounter serverOverloadCounter; + private final LongCounter limitChangeCounter; + + private final ObservableLongGauge inFlightGauge; + private final ObservableLongGauge proactiveLimitGauge; + private final ObservableLongGauge reactiveLimitGauge; + private final ObservableLongGauge effectiveLimitGauge; + + private volatile boolean closed; + + public OpenTelemetryClientThrottleMetrics(String connHash, ClientThrottleStateProvider stateProvider) { + this(resolveGlobalOpenTelemetryWithWarning(), connHash, stateProvider); + } + + /** + * Resolve {@link GlobalOpenTelemetry#get()} and log a single WARN the first time we observe + * that the global is the no-op instance. Without an SDK registered (e.g. via the OTel + * Java agent, {@code opentelemetry-sdk-extension-autoconfigure}, or a manual + * {@code GlobalOpenTelemetry.set(...)}), every counter and gauge below silently drops on the + * floor. This warning surfaces the misconfiguration so the operator knows why no metrics appear. + */ + private static OpenTelemetry resolveGlobalOpenTelemetryWithWarning() { + OpenTelemetry otel = GlobalOpenTelemetry.get(); + if (otel == OpenTelemetry.noop() && NOOP_WARNED.compareAndSet(false, true)) { + log.warn("ojp.jdbc.metrics=otel selected, but GlobalOpenTelemetry is the no-op instance — " + + "no client throttle metrics will be exported. Register an OpenTelemetry SDK before " + + "the first JDBC connection (e.g. via the OpenTelemetry Java agent, " + + "opentelemetry-sdk-extension-autoconfigure, or GlobalOpenTelemetry.set(sdk))."); + } + return otel; + } + + public OpenTelemetryClientThrottleMetrics(OpenTelemetry openTelemetry, String connHash, + ClientThrottleStateProvider stateProvider) { + if (openTelemetry == null) { + throw new IllegalArgumentException("openTelemetry must not be null"); + } + if (connHash == null || connHash.isEmpty()) { + throw new IllegalArgumentException("connHash must not be null or empty"); + } + if (stateProvider == null) { + throw new IllegalArgumentException("stateProvider must not be null"); + } + + this.baseAttrs = Attributes.of(CONN_HASH_KEY, connHash); + this.increaseAttrs = Attributes.of(CONN_HASH_KEY, connHash, DIRECTION_KEY, "increase"); + this.decreaseAttrs = Attributes.of(CONN_HASH_KEY, connHash, DIRECTION_KEY, "decrease"); + + Meter meter = openTelemetry.getMeter(METER_NAME); + + this.acquiredCounter = meter.counterBuilder("ojp.client.throttle.acquired.total") + .setDescription("Successful client-side throttle acquisitions.") + .setUnit("requests") + .build(); + this.rejectedCounter = meter.counterBuilder("ojp.client.throttle.rejected.total") + .setDescription("Client-side fail-fast rejections.") + .setUnit("requests") + .build(); + this.serverOverloadCounter = meter.counterBuilder("ojp.client.throttle.server.overload.total") + .setDescription("Server-overload notifications (RESOURCE_EXHAUSTED) received from OJP server.") + .setUnit("events") + .build(); + this.limitChangeCounter = meter.counterBuilder("ojp.client.throttle.limit.changes.total") + .setDescription("AIMD limit changes, tagged by direction (increase|decrease).") + .setUnit("events") + .build(); + + // Capture state provider in async callbacks; the manager (and so the provider) lives + // for the JVM lifetime, mirroring ClientThrottleManager in driver core. + this.inFlightGauge = meter.gaugeBuilder("ojp.client.throttle.inflight") + .ofLongs() + .setDescription("Current in-flight client-side throttled requests.") + .setUnit("requests") + .buildWithCallback(obs -> obs.record(stateProvider.getInFlight(), baseAttrs)); + this.proactiveLimitGauge = meter.gaugeBuilder("ojp.client.throttle.limit.proactive") + .ofLongs() + .setDescription("Current proactive limit (derived from server SessionInfo).") + .setUnit("requests") + .buildWithCallback(obs -> obs.record(saturate(stateProvider.getProactiveLimit()), baseAttrs)); + this.reactiveLimitGauge = meter.gaugeBuilder("ojp.client.throttle.limit.reactive") + .ofLongs() + .setDescription("Current reactive limit (AIMD on server overload).") + .setUnit("requests") + .buildWithCallback(obs -> obs.record(saturate(stateProvider.getReactiveLimit()), baseAttrs)); + this.effectiveLimitGauge = meter.gaugeBuilder("ojp.client.throttle.limit.effective") + .ofLongs() + .setDescription("Current effective limit (min of proactive and reactive in COMBINED mode).") + .setUnit("requests") + .buildWithCallback(obs -> obs.record(saturate(stateProvider.getEffectiveLimit()), baseAttrs)); + + log.debug("OpenTelemetry client throttle metrics initialised for connHash={}", connHash); + } + + /** + * Driver represents "no limit" as {@link Integer#MAX_VALUE}; emit 0 for the gauge so dashboards + * do not show a 2-billion bar before SessionInfo arrives. + */ + private static long saturate(int value) { + return value == Integer.MAX_VALUE ? 0L : value; + } + + @Override + public void recordAcquired() { + acquiredCounter.add(1, baseAttrs); + } + + @Override + public void recordRejected() { + rejectedCounter.add(1, baseAttrs); + } + + @Override + public void recordServerOverload() { + serverOverloadCounter.add(1, baseAttrs); + } + + @Override + public void recordLimitChange(LimitChangeDirection direction) { + limitChangeCounter.add(1, direction == LimitChangeDirection.INCREASE ? increaseAttrs : decreaseAttrs); + } + + @Override + public synchronized void close() { + if (closed) { + return; + } + closed = true; + // Best-effort: unregister callbacks so a closed manager stops contributing gauge values. + closeQuietly(inFlightGauge); + closeQuietly(proactiveLimitGauge); + closeQuietly(reactiveLimitGauge); + closeQuietly(effectiveLimitGauge); + } + + private static void closeQuietly(ObservableLongGauge gauge) { + if (gauge == null) { + return; + } + try { + gauge.close(); + } catch (Exception e) { + log.debug("Closing observable gauge failed: {}", e.getMessage()); + } + } +} diff --git a/ojp-jdbc-driver-otel-metrics/src/main/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsProvider.java b/ojp-jdbc-driver-otel-metrics/src/main/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsProvider.java new file mode 100644 index 000000000..ed25ab9df --- /dev/null +++ b/ojp-jdbc-driver-otel-metrics/src/main/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsProvider.java @@ -0,0 +1,42 @@ +package org.openjproxy.jdbc.metrics.otel; + +import lombok.extern.slf4j.Slf4j; +import org.openjproxy.jdbc.metrics.ClientThrottleMetrics; +import org.openjproxy.jdbc.metrics.ClientThrottleMetricsProvider; +import org.openjproxy.jdbc.metrics.ClientThrottleStateProvider; +import org.openjproxy.jdbc.metrics.NoOpClientThrottleMetrics; + +/** + * {@link ClientThrottleMetricsProvider} that binds to {@code ojp.jdbc.metrics=otel}. + * + *

Registered via {@code META-INF/services/org.openjproxy.jdbc.metrics.ClientThrottleMetricsProvider}. + * Constructs an {@link OpenTelemetryClientThrottleMetrics} backed by + * {@link io.opentelemetry.api.GlobalOpenTelemetry#get()}.

+ * + *

If the OpenTelemetry API is missing at runtime (NoClassDefFoundError) or any unexpected + * failure occurs, returns {@link NoOpClientThrottleMetrics#INSTANCE} rather than letting the + * driver propagate a failure — metrics must never break JDBC functionality.

+ */ +@Slf4j +public final class OpenTelemetryClientThrottleMetricsProvider implements ClientThrottleMetricsProvider { + + @Override + public String name() { + return "otel"; + } + + @Override + public ClientThrottleMetrics create(String connHash, ClientThrottleStateProvider stateProvider) { + try { + return new OpenTelemetryClientThrottleMetrics(connHash, stateProvider); + } catch (NoClassDefFoundError e) { + log.warn("OpenTelemetry API not available at runtime; falling back to no-op metrics: {}", + e.getMessage()); + return NoOpClientThrottleMetrics.INSTANCE; + } catch (RuntimeException e) { + log.warn("Failed to initialise OpenTelemetry client throttle metrics for connHash={}: {}", + connHash, e.getMessage()); + return NoOpClientThrottleMetrics.INSTANCE; + } + } +} diff --git a/ojp-jdbc-driver-otel-metrics/src/main/resources/META-INF/services/org.openjproxy.jdbc.metrics.ClientThrottleMetricsProvider b/ojp-jdbc-driver-otel-metrics/src/main/resources/META-INF/services/org.openjproxy.jdbc.metrics.ClientThrottleMetricsProvider new file mode 100644 index 000000000..3ea9f5aeb --- /dev/null +++ b/ojp-jdbc-driver-otel-metrics/src/main/resources/META-INF/services/org.openjproxy.jdbc.metrics.ClientThrottleMetricsProvider @@ -0,0 +1 @@ +org.openjproxy.jdbc.metrics.otel.OpenTelemetryClientThrottleMetricsProvider diff --git a/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/FactorySelectsOtelWhenAdapterPresentTest.java b/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/FactorySelectsOtelWhenAdapterPresentTest.java new file mode 100644 index 000000000..fdb7e4a09 --- /dev/null +++ b/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/FactorySelectsOtelWhenAdapterPresentTest.java @@ -0,0 +1,44 @@ +package org.openjproxy.jdbc.metrics.otel; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.openjproxy.jdbc.metrics.ClientThrottleMetrics; +import org.openjproxy.jdbc.metrics.ClientThrottleMetricsFactory; +import org.openjproxy.jdbc.metrics.ClientThrottleStateProvider; + +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Verifies that with this adapter on the classpath, the driver-core factory + * picks the {@link OpenTelemetryClientThrottleMetrics} when + * {@code -Dojp.jdbc.metrics=otel} is set. + */ +class FactorySelectsOtelWhenAdapterPresentTest { + + private static final class FixedState implements ClientThrottleStateProvider { + @Override public String getMode() { return "OFF"; } + @Override public int getInFlight() { return 0; } + @Override public int getProactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getReactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getEffectiveLimit() { return Integer.MAX_VALUE; } + } + + @AfterEach + void tearDown() { + System.clearProperty(ClientThrottleMetricsFactory.PROPERTY); + } + + @Test + void shouldReturnOtelMetricsWhenPropertyIsOtelAndAdapterOnClasspath() { + System.setProperty(ClientThrottleMetricsFactory.PROPERTY, "otel"); + ClientThrottleMetrics m = ClientThrottleMetricsFactory.create("h-factory", new FixedState()); + try { + assertNotNull(m); + assertTrue(m instanceof OpenTelemetryClientThrottleMetrics, + "Expected OpenTelemetryClientThrottleMetrics, got " + m.getClass().getName()); + } finally { + m.close(); + } + } +} diff --git a/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsProviderTest.java b/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsProviderTest.java new file mode 100644 index 000000000..ef2a4d186 --- /dev/null +++ b/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsProviderTest.java @@ -0,0 +1,48 @@ +package org.openjproxy.jdbc.metrics.otel; + +import org.junit.jupiter.api.Test; +import org.openjproxy.jdbc.metrics.ClientThrottleMetrics; +import org.openjproxy.jdbc.metrics.ClientThrottleMetricsProvider; +import org.openjproxy.jdbc.metrics.ClientThrottleStateProvider; + +import java.util.ServiceLoader; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class OpenTelemetryClientThrottleMetricsProviderTest { + + private static final class FixedState implements ClientThrottleStateProvider { + @Override public String getMode() { return "OFF"; } + @Override public int getInFlight() { return 0; } + @Override public int getProactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getReactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getEffectiveLimit() { return Integer.MAX_VALUE; } + } + + @Test + void shouldExposeOtelNameAndBeDiscoverableViaServiceLoader() { + boolean found = false; + for (ClientThrottleMetricsProvider p : ServiceLoader.load(ClientThrottleMetricsProvider.class)) { + if (p instanceof OpenTelemetryClientThrottleMetricsProvider) { + assertEquals("otel", p.name()); + found = true; + } + } + assertTrue(found, "OpenTelemetryClientThrottleMetricsProvider must be registered via ServiceLoader"); + } + + @Test + void shouldReturnMetricsInstanceWhenCreated() { + OpenTelemetryClientThrottleMetricsProvider provider = new OpenTelemetryClientThrottleMetricsProvider(); + ClientThrottleMetrics m = provider.create("h-x", new FixedState()); + assertNotNull(m); + // Recording must not throw regardless of which concrete type was returned (real OTel or no-op fallback). + m.recordAcquired(); + m.recordRejected(); + m.recordServerOverload(); + m.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.INCREASE); + m.close(); + } +} diff --git a/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsTest.java b/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsTest.java new file mode 100644 index 000000000..9e9367e1f --- /dev/null +++ b/ojp-jdbc-driver-otel-metrics/src/test/java/org/openjproxy/jdbc/metrics/otel/OpenTelemetryClientThrottleMetricsTest.java @@ -0,0 +1,155 @@ +package org.openjproxy.jdbc.metrics.otel; + +import io.opentelemetry.api.OpenTelemetry; +import io.opentelemetry.sdk.OpenTelemetrySdk; +import io.opentelemetry.sdk.metrics.SdkMeterProvider; +import io.opentelemetry.sdk.metrics.data.MetricData; +import io.opentelemetry.sdk.testing.exporter.InMemoryMetricReader; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.openjproxy.jdbc.metrics.ClientThrottleMetrics; +import org.openjproxy.jdbc.metrics.ClientThrottleStateProvider; + +import java.util.Collection; +import java.util.Optional; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class OpenTelemetryClientThrottleMetricsTest { + + private InMemoryMetricReader reader; + private OpenTelemetry openTelemetry; + private OpenTelemetryClientThrottleMetrics metrics; + + private static final class FixedState implements ClientThrottleStateProvider { + @Override public String getMode() { return "COMBINED"; } + @Override public int getInFlight() { return 4; } + @Override public int getProactiveLimit() { return 10; } + @Override public int getReactiveLimit() { return 7; } + @Override public int getEffectiveLimit() { return 7; } + } + + @BeforeEach + void setUp() { + reader = InMemoryMetricReader.create(); + SdkMeterProvider meterProvider = SdkMeterProvider.builder().registerMetricReader(reader).build(); + openTelemetry = OpenTelemetrySdk.builder().setMeterProvider(meterProvider).build(); + } + + @AfterEach + void tearDown() { + if (metrics != null) { + metrics.close(); + } + } + + private Optional metric(String name) { + Collection all = reader.collectAllMetrics(); + return all.stream().filter(m -> m.getName().equals(name)).findFirst(); + } + + private long sumCounter(String name) { + return metric(name) + .map(m -> m.getLongSumData().getPoints().stream().mapToLong(p -> p.getValue()).sum()) + .orElse(0L); + } + + private long firstGauge(String name) { + return metric(name) + .map(m -> m.getLongGaugeData().getPoints().iterator().next().getValue()) + .orElseThrow(() -> new AssertionError("missing gauge: " + name)); + } + + @Test + void shouldRejectNullArguments() { + assertThrows(IllegalArgumentException.class, + () -> new OpenTelemetryClientThrottleMetrics(null, "h", new FixedState())); + assertThrows(IllegalArgumentException.class, + () -> new OpenTelemetryClientThrottleMetrics(openTelemetry, null, new FixedState())); + assertThrows(IllegalArgumentException.class, + () -> new OpenTelemetryClientThrottleMetrics(openTelemetry, "", new FixedState())); + assertThrows(IllegalArgumentException.class, + () -> new OpenTelemetryClientThrottleMetrics(openTelemetry, "h", null)); + } + + @Test + void shouldIncrementCountersWhenEventsRecorded() { + metrics = new OpenTelemetryClientThrottleMetrics(openTelemetry, "h-A", new FixedState()); + + metrics.recordAcquired(); + metrics.recordAcquired(); + metrics.recordAcquired(); + metrics.recordRejected(); + metrics.recordServerOverload(); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.INCREASE); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); + + assertEquals(3L, sumCounter("ojp.client.throttle.acquired.total")); + assertEquals(1L, sumCounter("ojp.client.throttle.rejected.total")); + assertEquals(1L, sumCounter("ojp.client.throttle.server.overload.total")); + // limit.changes.total is split by direction; total across points should be 3. + assertEquals(3L, sumCounter("ojp.client.throttle.limit.changes.total")); + } + + @Test + void shouldEmitGaugesFromStateProvider() { + metrics = new OpenTelemetryClientThrottleMetrics(openTelemetry, "h-B", new FixedState()); + + assertEquals(4L, firstGauge("ojp.client.throttle.inflight")); + assertEquals(10L, firstGauge("ojp.client.throttle.limit.proactive")); + assertEquals(7L, firstGauge("ojp.client.throttle.limit.reactive")); + assertEquals(7L, firstGauge("ojp.client.throttle.limit.effective")); + } + + @Test + void shouldSaturateIntMaxValueToZeroForGauges() { + ClientThrottleStateProvider unlimited = new ClientThrottleStateProvider() { + @Override public String getMode() { return "OFF"; } + @Override public int getInFlight() { return 0; } + @Override public int getProactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getReactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getEffectiveLimit() { return Integer.MAX_VALUE; } + }; + metrics = new OpenTelemetryClientThrottleMetrics(openTelemetry, "h-C", unlimited); + + assertEquals(0L, firstGauge("ojp.client.throttle.limit.proactive")); + assertEquals(0L, firstGauge("ojp.client.throttle.limit.reactive")); + assertEquals(0L, firstGauge("ojp.client.throttle.limit.effective")); + } + + @Test + void shouldNotThrowOnDoubleClose() { + metrics = new OpenTelemetryClientThrottleMetrics(openTelemetry, "h-D", new FixedState()); + metrics.close(); + metrics.close(); + // After close, callbacks unregistered — gauges may or may not be present, but recording must not throw. + metrics.recordAcquired(); + assertNotNull(metrics); + metrics = null; // prevent tearDown re-close + } + + @Test + void shouldTagLimitChangesByDirection() { + metrics = new OpenTelemetryClientThrottleMetrics(openTelemetry, "h-E", new FixedState()); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.INCREASE); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); + + MetricData m = metric("ojp.client.throttle.limit.changes.total") + .orElseThrow(() -> new AssertionError("missing")); + long increase = m.getLongSumData().getPoints().stream() + .filter(p -> "increase".equals(p.getAttributes().get(OpenTelemetryClientThrottleMetrics.DIRECTION_KEY))) + .mapToLong(p -> p.getValue()).sum(); + long decrease = m.getLongSumData().getPoints().stream() + .filter(p -> "decrease".equals(p.getAttributes().get(OpenTelemetryClientThrottleMetrics.DIRECTION_KEY))) + .mapToLong(p -> p.getValue()).sum(); + assertEquals(1L, increase); + assertEquals(2L, decrease); + assertTrue(m.getLongSumData().isMonotonic()); + } +} diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/ClientThrottleManager.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/ClientThrottleManager.java index 23907db4a..95f23955f 100644 --- a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/ClientThrottleManager.java +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/ClientThrottleManager.java @@ -2,6 +2,8 @@ import com.openjproxy.grpc.SessionInfo; import lombok.extern.slf4j.Slf4j; +import org.openjproxy.jdbc.metrics.ClientThrottleMetrics; +import org.openjproxy.jdbc.metrics.NoOpClientThrottleMetrics; import java.util.concurrent.atomic.AtomicInteger; @@ -27,6 +29,21 @@ public class ClientThrottleManager { private volatile int reactiveLimit = Integer.MAX_VALUE; private volatile int lastProactiveLimit = Integer.MAX_VALUE; private volatile int lastReactiveLimit = Integer.MAX_VALUE; + private volatile ClientThrottleMetrics metrics = NoOpClientThrottleMetrics.INSTANCE; + + /** + * Attach a metrics sink. Idempotent for the No-op default; for a real + * sink, the caller (typically {@link Connection}) is responsible for + * lifecycle (close on connection-pool eviction is not required because + * the manager lives for the JVM lifetime in {@code THROTTLE_MANAGERS}). + */ + public void setMetrics(ClientThrottleMetrics metrics) { + this.metrics = metrics == null ? NoOpClientThrottleMetrics.INSTANCE : metrics; + } + + public ClientThrottleMetrics getMetrics() { + return metrics; + } /** * Update limits from a fresh SessionInfo. @@ -59,8 +76,10 @@ public void updateFromSessionInfo(SessionInfo sessionInfo) { } if (newProactive < lastProactiveLimit) { proactiveLimit = newProactive; + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); } else if (newProactive > lastProactiveLimit) { proactiveLimit = lastProactiveLimit + 1; + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.INCREASE); } lastProactiveLimit = proactiveLimit; } @@ -74,8 +93,10 @@ public void updateFromSessionInfo(SessionInfo sessionInfo) { } if (newReactive < lastReactiveLimit) { reactiveLimit = newReactive; + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); } else if (newReactive > lastReactiveLimit) { reactiveLimit = lastReactiveLimit + 1; + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.INCREASE); } lastReactiveLimit = reactiveLimit; } else { @@ -121,12 +142,19 @@ public boolean tryAcquire(ClientThrottleMode mode, boolean inTransaction) { int effectiveLimit = getEffectiveLimit(mode); if (effectiveLimit == Integer.MAX_VALUE) { + // No limit configured yet (e.g. before first SessionInfo arrives, or REACTIVE mode + // pre-overload). Still increment inFlight so that the observability gauge reflects + // actual concurrent driver load, and so acquire/release remain symmetric (release() + // always decrements). The increment is cheap (one CAS) and never blocks the request. + inFlight.incrementAndGet(); + metrics.recordAcquired(); return true; } int current = inFlight.get(); if (current >= effectiveLimit) { log.debug("Client throttle rejected: inFlight={}, effectiveLimit={}, mode={}", current, effectiveLimit, mode); + metrics.recordRejected(); return false; } // CAS loop: atomically check-and-increment to avoid exceeding the limit due to races @@ -134,9 +162,11 @@ public boolean tryAcquire(ClientThrottleMode mode, boolean inTransaction) { int cur = inFlight.get(); if (cur >= effectiveLimit) { log.debug("Client throttle rejected (CAS): inFlight={}, effectiveLimit={}, mode={}", cur, effectiveLimit, mode); + metrics.recordRejected(); return false; } if (inFlight.compareAndSet(cur, cur + 1)) { + metrics.recordAcquired(); return true; } } @@ -153,7 +183,7 @@ public void release(ClientThrottleMode mode, boolean inTransaction) { inFlight.updateAndGet(v -> Math.max(0, v - 1)); } - private int getEffectiveLimit(ClientThrottleMode mode) { + int getEffectiveLimit(ClientThrottleMode mode) { switch (mode) { case PROACTIVE: return proactiveLimit; case REACTIVE: return reactiveLimit; @@ -196,6 +226,10 @@ public void notifyServerOverload() { } reactiveLimit = newLimit; lastReactiveLimit = newLimit; + metrics.recordServerOverload(); + if (current != newLimit) { + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); + } log.debug("ClientThrottleManager notifyServerOverload: reactiveLimit {} -> {}", current, newLimit); } diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/Connection.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/Connection.java index 406b6e47d..1f17c0b8d 100644 --- a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/Connection.java +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/Connection.java @@ -15,6 +15,9 @@ import org.openjproxy.constants.CommonConstants; import org.openjproxy.grpc.ProtoConverter; import org.openjproxy.grpc.client.StatementService; +import org.openjproxy.jdbc.metrics.ClientThrottleMetrics; +import org.openjproxy.jdbc.metrics.ClientThrottleMetricsFactory; +import org.openjproxy.jdbc.metrics.ClientThrottleStateProvider; import java.sql.SQLClientInfoException; import java.sql.SQLException; @@ -92,14 +95,51 @@ public void setSession(SessionInfo session) { private void updateThrottleLimits(SessionInfo info) { if (info != null && !info.getConnHash().isEmpty()) { - THROTTLE_MANAGERS.computeIfAbsent(info.getConnHash(), k -> new ClientThrottleManager()) - .updateFromSessionInfo(info); + ClientThrottleManager manager = THROTTLE_MANAGERS.computeIfAbsent(info.getConnHash(), + k -> createThrottleManager(k)); + manager.updateFromSessionInfo(info); } } + private ClientThrottleManager createThrottleManager(String connHash) { + ClientThrottleManager manager = new ClientThrottleManager(); + // Snapshot the mode to avoid the state provider holding a reference back to this Connection + // (this throttle manager lives in a static map for the lifetime of the JVM). + final ClientThrottleMode capturedMode = throttleMode == null ? ClientThrottleMode.OFF : throttleMode; + ClientThrottleStateProvider stateProvider = new ClientThrottleStateProvider() { + @Override + public String getMode() { + return capturedMode.name(); + } + + @Override + public int getInFlight() { + return manager.getInFlight(); + } + + @Override + public int getProactiveLimit() { + return manager.getProactiveLimit(); + } + + @Override + public int getReactiveLimit() { + return manager.getReactiveLimit(); + } + + @Override + public int getEffectiveLimit() { + return manager.getEffectiveLimit(capturedMode); + } + }; + ClientThrottleMetrics metrics = ClientThrottleMetricsFactory.create(connHash, stateProvider); + manager.setMetrics(metrics); + return manager; + } + ClientThrottleManager getThrottleManager() { if (session != null && !session.getConnHash().isEmpty()) { - return THROTTLE_MANAGERS.computeIfAbsent(session.getConnHash(), k -> new ClientThrottleManager()); + return THROTTLE_MANAGERS.computeIfAbsent(session.getConnHash(), k -> createThrottleManager(k)); } return null; } diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetrics.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetrics.java new file mode 100644 index 000000000..cce4c4ee3 --- /dev/null +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetrics.java @@ -0,0 +1,37 @@ +package org.openjproxy.jdbc.metrics; + +/** + * Driver-side throttling metrics sink. + * + *

Implementations may publish to JMX, OpenTelemetry, or do nothing. + * Selected at runtime by {@link ClientThrottleMetricsFactory} based on the + * {@code ojp.jdbc.metrics} system property.

+ * + *

All methods must be safe to call from many JDBC client threads + * concurrently and must not throw.

+ * + *

See ADR-010 and {@code documents/analysis/THROTTLING_METRICS_ANALYSIS.md} §5.

+ */ +public interface ClientThrottleMetrics { + + /** Record one successful client-side throttle acquisition. */ + void recordAcquired(); + + /** Record one client-side fail-fast rejection (no server round-trip). */ + void recordRejected(); + + /** Record a server-overload notification (RESOURCE_EXHAUSTED from the server). */ + void recordServerOverload(); + + /** Record an AIMD limit change. */ + void recordLimitChange(LimitChangeDirection direction); + + /** Release any registered resources (e.g. unregister MBeans). Must be idempotent. */ + void close(); + + /** Direction tag for {@link #recordLimitChange(LimitChangeDirection)}. */ + enum LimitChangeDirection { + INCREASE, + DECREASE + } +} diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsFactory.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsFactory.java new file mode 100644 index 000000000..027fa25f5 --- /dev/null +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsFactory.java @@ -0,0 +1,91 @@ +package org.openjproxy.jdbc.metrics; + +import lombok.extern.slf4j.Slf4j; + +import java.util.Locale; +import java.util.ServiceLoader; + +/** + * Selects the {@link ClientThrottleMetrics} implementation based on the + * {@code ojp.jdbc.metrics} system property. + * + *
    + *
  • {@code jmx} — default; registers a {@link JmxClientThrottleMetrics} MBean.
  • + *
  • {@code none} — returns {@link NoOpClientThrottleMetrics#INSTANCE}.
  • + *
  • Any other value — looked up via {@link ServiceLoader} for a + * {@link ClientThrottleMetricsProvider} whose {@link ClientThrottleMetricsProvider#name() name} + * matches (case-insensitive). For example, the {@code ojp-jdbc-driver-otel-metrics} + * adapter registers a provider named {@code "otel"}. If no matching provider is + * found, falls back to {@code jmx} with a single WARN log.
  • + *
+ * + *

See ADR-010.

+ */ +@Slf4j +public final class ClientThrottleMetricsFactory { + + /** System property name. */ + public static final String PROPERTY = "ojp.jdbc.metrics"; + + /** Default binding when the property is unset or unrecognised. */ + public static final String DEFAULT = "jmx"; + + private ClientThrottleMetricsFactory() { + // utility + } + + /** + * Create a metrics sink for the given {@code connHash}. + * + * @param connHash connection-hash identifying the throttle scope (must not be null) + * @param stateProvider live snapshot provider for gauge attributes (must not be null) + */ + public static ClientThrottleMetrics create(String connHash, ClientThrottleStateProvider stateProvider) { + if (connHash == null || connHash.isEmpty() || stateProvider == null) { + return NoOpClientThrottleMetrics.INSTANCE; + } + String binding = resolveBinding(); + switch (binding) { + case "none": + return NoOpClientThrottleMetrics.INSTANCE; + case "jmx": + return new JmxClientThrottleMetrics(connHash, stateProvider); + default: + return loadProvider(binding, connHash, stateProvider); + } + } + + private static ClientThrottleMetrics loadProvider(String binding, String connHash, + ClientThrottleStateProvider stateProvider) { + try { + for (ClientThrottleMetricsProvider provider : ServiceLoader.load(ClientThrottleMetricsProvider.class)) { + if (provider != null && binding.equalsIgnoreCase(provider.name())) { + ClientThrottleMetrics metrics = provider.create(connHash, stateProvider); + if (metrics != null) { + return metrics; + } + } + } + } catch (Throwable t) { + // ServiceLoader.iterator() can throw ServiceConfigurationError; never let metrics break JDBC. + log.warn("ServiceLoader lookup for ojp.jdbc.metrics={} failed: {}; falling back to JMX.", + binding, t.getMessage()); + return new JmxClientThrottleMetrics(connHash, stateProvider); + } + log.warn("ojp.jdbc.metrics={} requires a matching ClientThrottleMetricsProvider on the classpath " + + "(e.g. ojp-jdbc-driver-otel-metrics for 'otel'); falling back to JMX. See ADR-010.", binding); + return new JmxClientThrottleMetrics(connHash, stateProvider); + } + + private static String resolveBinding() { + String raw = System.getProperty(PROPERTY, DEFAULT); + if (raw == null) { + return DEFAULT; + } + String value = raw.trim().toLowerCase(Locale.ROOT); + if (value.isEmpty()) { + return DEFAULT; + } + return value; + } +} diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsMXBean.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsMXBean.java new file mode 100644 index 000000000..92eb21be6 --- /dev/null +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsMXBean.java @@ -0,0 +1,45 @@ +package org.openjproxy.jdbc.metrics; + +/** + * JMX-visible view of one connection-hash's client throttling state. + * + *

Bound under {@code org.openjproxy:type=ClientThrottle,connHash=}.

+ * + *

Attributes are live readings (counters are monotonic; gauges reflect + * current state at read time).

+ */ +public interface ClientThrottleMetricsMXBean { + + /** Connection-hash identifying this throttle scope. */ + String getConnHash(); + + /** Configured throttle mode (e.g. {@code OFF}, {@code PROACTIVE}, {@code REACTIVE}, {@code COMBINED}). */ + String getMode(); + + /** Live in-flight count at this driver instance. */ + int getInFlight(); + + /** Current proactive limit (derived from {@code SessionInfo}). */ + int getProactiveLimit(); + + /** Current reactive limit (AIMD). */ + int getReactiveLimit(); + + /** Effective limit actually enforced — min of proactive and reactive under {@code COMBINED}. */ + int getEffectiveLimit(); + + /** Total client-side fail-fast rejections. */ + long getRejectedTotal(); + + /** Total successful client-side acquisitions. */ + long getAcquiredTotal(); + + /** Total server-overload notifications received (RESOURCE_EXHAUSTED). */ + long getServerOverloadEventsTotal(); + + /** Total AIMD limit-increase events. */ + long getLimitIncreaseTotal(); + + /** Total AIMD limit-decrease events. */ + long getLimitDecreaseTotal(); +} diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsProvider.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsProvider.java new file mode 100644 index 000000000..3e6a3f579 --- /dev/null +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsProvider.java @@ -0,0 +1,29 @@ +package org.openjproxy.jdbc.metrics; + +/** + * SPI for pluggable {@link ClientThrottleMetrics} implementations. + * + *

Implementations are discovered by {@link ClientThrottleMetricsFactory} + * via {@link java.util.ServiceLoader} when the driver is configured with + * {@code -Dojp.jdbc.metrics=otel} (or any future non-built-in binding name).

+ * + *

An adapter module — for example {@code ojp-jdbc-driver-otel-metrics} — + * registers an implementation under + * {@code META-INF/services/org.openjproxy.jdbc.metrics.ClientThrottleMetricsProvider}.

+ * + *

See ADR-010.

+ */ +public interface ClientThrottleMetricsProvider { + + /** + * The binding name this provider serves (matches the {@code ojp.jdbc.metrics} + * system-property value, e.g. {@code "otel"}). Must not be {@code null} or empty. + */ + String name(); + + /** + * Create a metrics sink for the given {@code connHash}. + * Must not throw; implementations should log and return a no-op when in doubt. + */ + ClientThrottleMetrics create(String connHash, ClientThrottleStateProvider stateProvider); +} diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleStateProvider.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleStateProvider.java new file mode 100644 index 000000000..72bddf41f --- /dev/null +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/ClientThrottleStateProvider.java @@ -0,0 +1,24 @@ +package org.openjproxy.jdbc.metrics; + +/** + * Read-only snapshot provider for the gauges that back + * {@link ClientThrottleMetricsMXBean}. Decouples the metrics implementation + * from {@code ClientThrottleManager} so the dependency arrow points one way. + */ +public interface ClientThrottleStateProvider { + + /** Configured mode label (e.g. {@code COMBINED}). */ + String getMode(); + + /** Current in-flight request count. */ + int getInFlight(); + + /** Current proactive limit. */ + int getProactiveLimit(); + + /** Current reactive limit. */ + int getReactiveLimit(); + + /** Current effective limit ({@code min(proactive, reactive)} under {@code COMBINED}). */ + int getEffectiveLimit(); +} diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/JmxClientThrottleMetrics.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/JmxClientThrottleMetrics.java new file mode 100644 index 000000000..7b04135de --- /dev/null +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/JmxClientThrottleMetrics.java @@ -0,0 +1,169 @@ +package org.openjproxy.jdbc.metrics; + +import lombok.extern.slf4j.Slf4j; + +import javax.management.InstanceAlreadyExistsException; +import javax.management.InstanceNotFoundException; +import javax.management.MBeanServer; +import javax.management.MalformedObjectNameException; +import javax.management.ObjectName; +import java.lang.management.ManagementFactory; +import java.util.concurrent.atomic.AtomicLong; + +/** + * JMX-backed {@link ClientThrottleMetrics}. Registers one MBean per + * {@code connHash} under {@code org.openjproxy:type=ClientThrottle,connHash=}. + * + *

Counters are {@link AtomicLong}s; gauges are read live through the + * supplied {@link ClientThrottleStateProvider}. MBean registration is + * best-effort: failure to register is logged at WARN and does not throw, so a + * locked-down host SecurityManager cannot break JDBC functionality.

+ * + *

See ADR-010.

+ */ +@Slf4j +public final class JmxClientThrottleMetrics implements ClientThrottleMetrics, ClientThrottleMetricsMXBean { + + static final String DOMAIN = "org.openjproxy"; + + private final String connHash; + private final ClientThrottleStateProvider stateProvider; + private final ObjectName objectName; + private final boolean registered; + + private final AtomicLong acquired = new AtomicLong(); + private final AtomicLong rejected = new AtomicLong(); + private final AtomicLong serverOverload = new AtomicLong(); + private final AtomicLong limitIncrease = new AtomicLong(); + private final AtomicLong limitDecrease = new AtomicLong(); + + public JmxClientThrottleMetrics(String connHash, ClientThrottleStateProvider stateProvider) { + this.connHash = connHash; + this.stateProvider = stateProvider; + this.objectName = buildObjectName(connHash); + this.registered = register(); + } + + private static ObjectName buildObjectName(String connHash) { + try { + // ObjectName values are quoted to tolerate any characters the connHash may legally contain. + return new ObjectName(DOMAIN + ":type=ClientThrottle,connHash=" + ObjectName.quote(connHash)); + } catch (MalformedObjectNameException e) { + log.warn("Could not build JMX ObjectName for connHash={}: {}", connHash, e.getMessage()); + return null; + } + } + + private boolean register() { + if (objectName == null) { + return false; + } + try { + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + server.registerMBean(this, objectName); + return true; + } catch (InstanceAlreadyExistsException e) { + // Another driver instance in the same JVM already registered this hash — that's fine. + log.debug("JMX MBean already registered for {}", objectName); + return false; + } catch (Exception e) { + log.warn("Could not register JMX MBean for connHash={}: {}", connHash, e.getMessage()); + return false; + } + } + + @Override + public void recordAcquired() { + acquired.incrementAndGet(); + } + + @Override + public void recordRejected() { + rejected.incrementAndGet(); + } + + @Override + public void recordServerOverload() { + serverOverload.incrementAndGet(); + } + + @Override + public void recordLimitChange(LimitChangeDirection direction) { + if (direction == LimitChangeDirection.INCREASE) { + limitIncrease.incrementAndGet(); + } else { + limitDecrease.incrementAndGet(); + } + } + + @Override + public void close() { + if (!registered || objectName == null) { + return; + } + try { + ManagementFactory.getPlatformMBeanServer().unregisterMBean(objectName); + } catch (InstanceNotFoundException e) { + // Already gone — fine. + } catch (Exception e) { + log.warn("Could not unregister JMX MBean {}: {}", objectName, e.getMessage()); + } + } + + // ---- ClientThrottleMetricsMXBean ---- + + @Override + public String getConnHash() { + return connHash; + } + + @Override + public String getMode() { + return stateProvider.getMode(); + } + + @Override + public int getInFlight() { + return stateProvider.getInFlight(); + } + + @Override + public int getProactiveLimit() { + return stateProvider.getProactiveLimit(); + } + + @Override + public int getReactiveLimit() { + return stateProvider.getReactiveLimit(); + } + + @Override + public int getEffectiveLimit() { + return stateProvider.getEffectiveLimit(); + } + + @Override + public long getRejectedTotal() { + return rejected.get(); + } + + @Override + public long getAcquiredTotal() { + return acquired.get(); + } + + @Override + public long getServerOverloadEventsTotal() { + return serverOverload.get(); + } + + @Override + public long getLimitIncreaseTotal() { + return limitIncrease.get(); + } + + @Override + public long getLimitDecreaseTotal() { + return limitDecrease.get(); + } +} diff --git a/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/NoOpClientThrottleMetrics.java b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/NoOpClientThrottleMetrics.java new file mode 100644 index 000000000..131bcdc76 --- /dev/null +++ b/ojp-jdbc-driver/src/main/java/org/openjproxy/jdbc/metrics/NoOpClientThrottleMetrics.java @@ -0,0 +1,39 @@ +package org.openjproxy.jdbc.metrics; + +/** + * No-op {@link ClientThrottleMetrics}. Used when {@code ojp.jdbc.metrics=none} + * or when no metrics backend is available. + */ +public final class NoOpClientThrottleMetrics implements ClientThrottleMetrics { + + public static final NoOpClientThrottleMetrics INSTANCE = new NoOpClientThrottleMetrics(); + + private NoOpClientThrottleMetrics() { + // singleton + } + + @Override + public void recordAcquired() { + // no-op + } + + @Override + public void recordRejected() { + // no-op + } + + @Override + public void recordServerOverload() { + // no-op + } + + @Override + public void recordLimitChange(LimitChangeDirection direction) { + // no-op + } + + @Override + public void close() { + // no-op + } +} diff --git a/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/ClientThrottleManagerMetricsIntegrationTest.java b/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/ClientThrottleManagerMetricsIntegrationTest.java new file mode 100644 index 000000000..83aeb2f40 --- /dev/null +++ b/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/ClientThrottleManagerMetricsIntegrationTest.java @@ -0,0 +1,96 @@ +package org.openjproxy.jdbc.metrics; + +import com.openjproxy.grpc.SessionInfo; +import org.junit.jupiter.api.Test; +import org.openjproxy.jdbc.ClientThrottleManager; +import org.openjproxy.jdbc.ClientThrottleMode; + +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Verifies that {@link ClientThrottleManager} invokes the attached + * {@link ClientThrottleMetrics} on each relevant event. + */ +class ClientThrottleManagerMetricsIntegrationTest { + + private static final class CountingMetrics implements ClientThrottleMetrics { + final AtomicInteger acquired = new AtomicInteger(); + final AtomicInteger rejected = new AtomicInteger(); + final AtomicInteger overload = new AtomicInteger(); + final AtomicInteger increase = new AtomicInteger(); + final AtomicInteger decrease = new AtomicInteger(); + + @Override public void recordAcquired() { acquired.incrementAndGet(); } + @Override public void recordRejected() { rejected.incrementAndGet(); } + @Override public void recordServerOverload() { overload.incrementAndGet(); } + @Override public void recordLimitChange(LimitChangeDirection direction) { + if (direction == LimitChangeDirection.INCREASE) { + increase.incrementAndGet(); + } else { + decrease.incrementAndGet(); + } + } + @Override public void close() { /* no-op */ } + } + + @Test + void shouldRecordAcquiredAndRejectedAndOverload() { + ClientThrottleManager manager = new ClientThrottleManager(); + CountingMetrics metrics = new CountingMetrics(); + manager.setMetrics(metrics); + + // Seed limit to 2 via SessionInfo. ClientThrottleManager applies a 0.9 safety margin + // after ceiling division: ceil(30/10) = 3; (int)(3 * 0.9) = 2. + manager.updateFromSessionInfo(SessionInfo.newBuilder() + .setConnHash("h") + .setMaxAdmission(30) + .setClientCount(10) + .build()); + + // 2 acquires succeed, 3rd is rejected. + assertEquals(true, manager.tryAcquire(ClientThrottleMode.PROACTIVE, false)); + assertEquals(true, manager.tryAcquire(ClientThrottleMode.PROACTIVE, false)); + assertEquals(false, manager.tryAcquire(ClientThrottleMode.PROACTIVE, false)); + + assertEquals(2, metrics.acquired.get()); + assertEquals(1, metrics.rejected.get()); + + // Server overload halves the reactive limit and records the event + a DECREASE. + manager.notifyServerOverload(); + assertEquals(1, metrics.overload.get()); + // Total decreases: 1 from initial proactive seed (newProactive < MAX_VALUE) + 1 from overload. + assertEquals(2, metrics.decrease.get()); + } + + @Test + void shouldNotRecordWhenModeOff() { + ClientThrottleManager manager = new ClientThrottleManager(); + CountingMetrics metrics = new CountingMetrics(); + manager.setMetrics(metrics); + + manager.tryAcquire(ClientThrottleMode.OFF, false); + manager.tryAcquire(ClientThrottleMode.OFF, false); + assertEquals(0, metrics.acquired.get()); + assertEquals(0, metrics.rejected.get()); + } + + @Test + void shouldTrackInFlightEvenWhenLimitIsUnlimited() { + // Before any SessionInfo arrives the effective limit is Integer.MAX_VALUE. The inflight + // gauge must still reflect concurrent driver load and stay symmetric with release(). + ClientThrottleManager manager = new ClientThrottleManager(); + manager.setMetrics(new CountingMetrics()); + + assertEquals(0, manager.getInFlight()); + assertTrue(manager.tryAcquire(ClientThrottleMode.COMBINED, false)); + assertTrue(manager.tryAcquire(ClientThrottleMode.COMBINED, false)); + assertEquals(2, manager.getInFlight()); + + manager.release(ClientThrottleMode.COMBINED, false); + manager.release(ClientThrottleMode.COMBINED, false); + assertEquals(0, manager.getInFlight()); + } +} diff --git a/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsFactoryTest.java b/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsFactoryTest.java new file mode 100644 index 000000000..07b6dd1d9 --- /dev/null +++ b/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/ClientThrottleMetricsFactoryTest.java @@ -0,0 +1,65 @@ +package org.openjproxy.jdbc.metrics; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class ClientThrottleMetricsFactoryTest { + + private static final class FixedState implements ClientThrottleStateProvider { + @Override public String getMode() { return "COMBINED"; } + @Override public int getInFlight() { return 0; } + @Override public int getProactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getReactiveLimit() { return Integer.MAX_VALUE; } + @Override public int getEffectiveLimit() { return Integer.MAX_VALUE; } + } + + @AfterEach + void clearProperty() { + System.clearProperty(ClientThrottleMetricsFactory.PROPERTY); + } + + @Test + void shouldReturnNoOpWhenConnHashIsNull() { + ClientThrottleMetrics m = ClientThrottleMetricsFactory.create(null, new FixedState()); + assertSame(NoOpClientThrottleMetrics.INSTANCE, m); + } + + @Test + void shouldReturnNoOpWhenStateProviderIsNull() { + ClientThrottleMetrics m = ClientThrottleMetricsFactory.create("h", null); + assertSame(NoOpClientThrottleMetrics.INSTANCE, m); + } + + @Test + void shouldReturnNoOpWhenPropertySetToNone() { + System.setProperty(ClientThrottleMetricsFactory.PROPERTY, "none"); + ClientThrottleMetrics m = ClientThrottleMetricsFactory.create("h-none", new FixedState()); + assertSame(NoOpClientThrottleMetrics.INSTANCE, m); + } + + @Test + void shouldReturnJmxByDefault() { + ClientThrottleMetrics m = ClientThrottleMetricsFactory.create("h-default", new FixedState()); + try { + assertNotNull(m); + assertTrue(m instanceof JmxClientThrottleMetrics); + } finally { + m.close(); + } + } + + @Test + void shouldFallBackToJmxWhenOtelRequestedAndAdapterMissing() { + System.setProperty(ClientThrottleMetricsFactory.PROPERTY, "otel"); + ClientThrottleMetrics m = ClientThrottleMetricsFactory.create("h-otel-fallback", new FixedState()); + try { + assertTrue(m instanceof JmxClientThrottleMetrics); + } finally { + m.close(); + } + } +} diff --git a/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/JmxClientThrottleMetricsTest.java b/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/JmxClientThrottleMetricsTest.java new file mode 100644 index 000000000..52e4e111e --- /dev/null +++ b/ojp-jdbc-driver/src/test/java/org/openjproxy/jdbc/metrics/JmxClientThrottleMetricsTest.java @@ -0,0 +1,100 @@ +package org.openjproxy.jdbc.metrics; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import javax.management.MBeanServer; +import javax.management.ObjectName; +import java.lang.management.ManagementFactory; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class JmxClientThrottleMetricsTest { + + private JmxClientThrottleMetrics metrics; + + private static final class FixedState implements ClientThrottleStateProvider { + @Override public String getMode() { return "COMBINED"; } + @Override public int getInFlight() { return 3; } + @Override public int getProactiveLimit() { return 10; } + @Override public int getReactiveLimit() { return 8; } + @Override public int getEffectiveLimit() { return 8; } + } + + @AfterEach + void tearDown() { + if (metrics != null) { + metrics.close(); + } + } + + @Test + void shouldRegisterMBeanAndExposeAttributesWhenCreated() throws Exception { + metrics = new JmxClientThrottleMetrics("hash-A", new FixedState()); + + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + ObjectName name = new ObjectName(JmxClientThrottleMetrics.DOMAIN + + ":type=ClientThrottle,connHash=" + ObjectName.quote("hash-A")); + assertTrue(server.isRegistered(name)); + + assertEquals("hash-A", server.getAttribute(name, "ConnHash")); + assertEquals("COMBINED", server.getAttribute(name, "Mode")); + assertEquals(3, server.getAttribute(name, "InFlight")); + assertEquals(10, server.getAttribute(name, "ProactiveLimit")); + assertEquals(8, server.getAttribute(name, "ReactiveLimit")); + assertEquals(8, server.getAttribute(name, "EffectiveLimit")); + } + + @Test + void shouldIncrementCountersWhenEventsRecorded() throws Exception { + metrics = new JmxClientThrottleMetrics("hash-B", new FixedState()); + + metrics.recordAcquired(); + metrics.recordAcquired(); + metrics.recordRejected(); + metrics.recordServerOverload(); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.INCREASE); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); + metrics.recordLimitChange(ClientThrottleMetrics.LimitChangeDirection.DECREASE); + + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + ObjectName name = new ObjectName(JmxClientThrottleMetrics.DOMAIN + + ":type=ClientThrottle,connHash=" + ObjectName.quote("hash-B")); + + assertEquals(2L, server.getAttribute(name, "AcquiredTotal")); + assertEquals(1L, server.getAttribute(name, "RejectedTotal")); + assertEquals(1L, server.getAttribute(name, "ServerOverloadEventsTotal")); + assertEquals(1L, server.getAttribute(name, "LimitIncreaseTotal")); + assertEquals(2L, server.getAttribute(name, "LimitDecreaseTotal")); + } + + @Test + void shouldUnregisterMBeanWhenClosed() throws Exception { + metrics = new JmxClientThrottleMetrics("hash-C", new FixedState()); + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + ObjectName name = new ObjectName(JmxClientThrottleMetrics.DOMAIN + + ":type=ClientThrottle,connHash=" + ObjectName.quote("hash-C")); + assertTrue(server.isRegistered(name)); + + metrics.close(); + assertFalse(server.isRegistered(name)); + + // Idempotent: a second close must not throw. + metrics.close(); + metrics = null; // prevent tearDown double-close + } + + @Test + void shouldNotThrowWhenSameConnHashRegisteredTwice() { + metrics = new JmxClientThrottleMetrics("hash-D", new FixedState()); + // Second registration of the same connHash is tolerated. + JmxClientThrottleMetrics dup = new JmxClientThrottleMetrics("hash-D", new FixedState()); + assertNotNull(dup); + // Recording on the duplicate must still be safe. + dup.recordAcquired(); + dup.close(); + } +} diff --git a/pom.xml b/pom.xml index 055e71aa4..db34746db 100644 --- a/pom.xml +++ b/pom.xml @@ -15,6 +15,7 @@ ojp-grpc-commons ojp-jdbc-driver + ojp-jdbc-driver-otel-metrics ojp-server ojp-datasource-api ojp-datasource-hikari