Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 12 additions & 11 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,24 +4,25 @@ This file provides guidance to Claude Code (claude.ai/code) when working with co

## Development Commands

Install the development dependency group with `uv`:
Use both dependency groups: `develop` supplies tools and `test` includes all
adapter and OTLP regression dependencies. Pass the groups to `uv run` too.

```bash
uv sync --group develop
uv sync --group develop --group test

# Run all tests, or a selected file/pattern
uv run pytest tests
uv run pytest tests/test_utils.py -k "connection_string"
uv run --group develop --group test pytest tests
uv run --group develop --group test pytest tests/test_utils.py -k "connection_string"

# Run static checks
uv run ruff check hasql tests
uv run mypy --install-types --non-interactive hasql tests
uv run --group develop --group test ruff check hasql tests
uv run --group develop --group test mypy --install-types --non-interactive hasql tests

# Run tox environments
uv run tox -e lint
uv run tox -e mypy
uv run tox -e py310 # Also available: py311, py312, py313, py314
uv run tox -e py310-uvloop # Also available for py311–py314
uv run --group develop --group test tox -e lint
uv run --group develop --group test tox -e mypy
uv run --group develop --group test tox -e py310 # Also py311–py314
uv run --group develop --group test tox -e py310-uvloop # Also py311–py314
```

The `justfile` contains only `chaos-*` recipes for the optional test cluster;
Expand Down Expand Up @@ -142,7 +143,7 @@ All balancer policies live in `hasql/balancer_policy/`. The module structure:
- `_get_candidates(read_only, fallback_master, choose_master_as_replica)` — builds candidates in fallback order: fresh replicas, optional master, known stale replicas, then wait
- `_get_pool(...)` — **abstract method**, the only thing subclasses must implement

**Import architecture:** The circular import between `pool_manager` and `balancer_policy` is broken by the `PoolStateProvider` protocol in `pool_state.py`. Balancer policies depend on `PoolStateProvider` (direct import from `pool_state.py`, no `TYPE_CHECKING` needed). `BasePoolManager` passes its `_pool_state` attribute (which implements `PoolStateProvider`) to the balancer constructor. Import graph (no cycles): `pool_state.py` → `utils.py`, `abc.py`, `acquire.py`, `metrics.py`, `staleness.py`; `balancer_policy/base.py` → `pool_state.py`; `pool_manager.py` → `pool_state.py`, `balancer_policy/`, `health.py`, `staleness.py`; `health.py` → `pool_manager.py` (TYPE_CHECKING only).
**Import architecture:** The circular import between `pool_manager` and `balancer_policy` is broken by the `PoolStateProvider` protocol in `pool_state.py`. Balancer policies depend on `PoolStateProvider` (direct import from `pool_state.py`, no `TYPE_CHECKING` needed). `BasePoolManager` passes its `_pool_state` attribute (which implements `PoolStateProvider`) to the balancer constructor. Import graph (no cycles): `pool_state.py` → `utils.py`, `abc.py`, `acquire.py`, `metrics.py`, `staleness.py`; `balancer_policy/base.py` → `pool_state.py`; `pool_manager.py` → `pool_state.py`, `balancer_policy/`, `health.py`, `staleness.py`; `health.py` → `pool_state.py` (runtime import).

**Policies:**

Expand Down
69 changes: 52 additions & 17 deletions README.rst
Original file line number Diff line number Diff line change
Expand Up @@ -491,20 +491,53 @@ Use the helper from ``example/otlp/common.py``:

.. code-block:: python

import asyncio

from hasql.driver.asyncpg import PoolManager
from example.otlp.common import (
register_hasql_metrics,
observe_hasql_metrics,
setup_meter_provider,
)

provider = setup_meter_provider(export_interval_ms=10_000)
async def main(dsn):
provider = setup_meter_provider(export_interval_ms=10_000)
try:
pool = PoolManager(dsn, fallback_master=True)
try:
await pool.ready()
async with observe_hasql_metrics(pool, sample_interval=1.0):
while True:
async with pool.acquire_master() as conn:
await conn.fetchval("SELECT 1")
await asyncio.sleep(1)
finally:
await pool.close()
finally:
await asyncio.to_thread(provider.shutdown)

The helper samples ``metrics()`` once on the owning event loop before
registration and then every ``sample_interval`` seconds (default: 1).
OTel callbacks read only a detached immutable snapshot, never the live manager.
Export runs independently (default: 10 seconds); both intervals must be positive
and finite. Manual collection also reads the latest sample, not live state.
Each callback retains one snapshot; a collection across instruments is not an
atomic batch. A blocked event loop delays sampling. Sampling errors are logged
and clear observations until the next successful sample. Acquire counters stay
cumulative; missing lag and selected extra keys produce no observation.

Pass ``extra_keys=("overflow",)`` for SQLAlchemy extras, or selected psycopg3
keys, to the same helper. It stops and awaits its sampler before pool cleanup;
the provider is owned by the caller and shut down off-loop even if cleanup
fails. Final SDK collection may use the last detached sample after sampling
stops. Export calls have a 10-second timeout; SDK shutdown defaults to 30 seconds.

Run the scripts as modules from the repository checkout (direct script paths
can shadow driver packages):

.. code-block:: bash

pool = PoolManager(dsn, fallback_master=True)
await pool.ready()
python -m example.otlp.asyncpg --dsn postgresql://u:p@db1,db2/mydb

# Registers observable gauges — OTel calls metrics()
# automatically at each export interval
register_hasql_metrics(pool)

Exported OTel instruments
*************************
Expand Down Expand Up @@ -572,19 +605,18 @@ Driver-specific extras
**********************

Some drivers expose additional pool internals via ``PoolMetrics.extra``.
Use ``register_extra_gauges()`` to export them as OTel gauges:
Pass ``extra_keys`` to the observation context in the quick start above:

.. code-block:: python

from example.otlp.common import register_extra_gauges

# psycopg3: queue depth, error counters, etc.
register_extra_gauges(pool, [
"pool_size", "requests_waiting", "connections_errors",
])
extra_keys = ("pool_size", "requests_waiting", "connections_errors")

# Or, for SQLAlchemy: overflow connections
extra_keys = ("overflow",)

# SQLAlchemy: overflow connections
register_extra_gauges(pool, ["overflow"])
async with observe_hasql_metrics(pool, extra_keys=extra_keys):
... # application workload

Per-driver examples live in ``example/otlp/``.

Expand Down Expand Up @@ -642,8 +674,11 @@ The exported metrics map well to Grafana / Datadog dashboard panels:
**Alerting rules**

* ``db.pool.masters == 0`` — **critical**: no master available
* ``db.pool.replicas == 0`` — **warning**: all reads will fall back
to master (if ``fallback_master=True``) or fail
* ``db.pool.replicas == 0`` — **warning**: no fresh replicas; reads use
an available master if fallback is enabled, otherwise an available known
stale replica. With no candidates, acquisition waits up to its timeout.
Separately, ``master_as_replica_weight`` can include a master alongside
fresh replicas
* ``db.pool.connections.used / db.pool.connections.max > 0.9`` —
**warning**: pool near exhaustion
* ``db.pool.health_check.duration > threshold`` — **warning**: host
Expand Down
2 changes: 1 addition & 1 deletion docs/class-scheme.md
Original file line number Diff line number Diff line change
Expand Up @@ -225,7 +225,7 @@ No additional methods or overrides — purely a convenience binding.
pool_state.py ──imports──► utils.py, abc.py, acquire.py, metrics.py
balancer_policy/base.py ──imports──► pool_state.py (PoolStateProvider only)
pool_manager.py ──imports──► pool_state.py, balancer_policy/, health.py, metrics.py, abc.py, acquire.py
health.py ──TYPE_CHECKING import──► pool_manager.py (BasePoolManager)
health.py ──imports──► pool_state.py (PoolState)
acquire.py ──TYPE_CHECKING import──► balancer_policy/base.py, pool_state.py
```

Expand Down
43 changes: 31 additions & 12 deletions docs/metrics-dashboard-example.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,23 @@ Datadog, New Relic, etc.) — only the query syntax differs.

## Metric reference

All gauges are registered by `register_hasql_metrics()` from
`example/otlp/common.py`. Driver-specific extras require an additional
`register_extra_gauges()` call.
Use `async with observe_hasql_metrics(pool, sample_interval=1.0, extra_keys=())`
from `example/otlp/common.py` around the workload, after `pool.ready()`.
One initial sample precedes registration; a background task samples once per
interval on the owning event loop, independently of the workload. Export
(default: 10 seconds) is independent of sampling (default: 1 second); both
intervals must be positive and finite. Worker callbacks and manual collection
read only the latest detached immutable snapshot, never live manager data.

Each callback retains one sample, not an atomic batch across instruments.
A blocked event loop delays sampling. Sampling errors are logged and clear
observations until sampling recovers. Missing lag/extra keys are omitted.
The context stops and awaits sampling before pool close; retain a pool cleanup
`finally` before readiness and a provider cleanup `finally` before pool creation.
Use `await asyncio.to_thread(provider.shutdown)` even if pool close fails.
Final SDK collection can export the last detached sample after sampling stops.
The example provider bounds export calls to 10 seconds; SDK shutdown defaults
to 30 seconds. See the README quick start for the complete cleanup structure.

| Gauge name | Labels | Description |
|---|---|---|
Expand Down Expand Up @@ -44,7 +58,7 @@ gauges. `timedelta` values from the `time` lag key are converted to seconds.

### Driver-specific extra keys

**psycopg3** (`register_extra_gauges(pool, [...])`):
**psycopg3** (`observe_hasql_metrics(pool, extra_keys=[...])`):

| Key | Description |
|---|---|
Expand All @@ -58,7 +72,7 @@ gauges. `timedelta` values from the `time` lag key are converted to seconds.
| `returns_bad` | Connections returned in bad state |
| `usage_ms` | Total connection usage time (ms) |

**SQLAlchemy** (`register_extra_gauges(pool, ["overflow"])`):
**SQLAlchemy** (`observe_hasql_metrics(pool, extra_keys=("overflow",))`):

| Key | Description |
|---|---|
Expand Down Expand Up @@ -97,8 +111,11 @@ Thresholds:
- 0 → red
```

When replicas drop to 0, read traffic either falls back to the master
(if `fallback_master=True`) or fails entirely.
When fresh replicas drop to 0, reads use an available master if
`fallback_master=True`; otherwise they can use an available known stale
replica. With no candidates, acquisition waits up to its timeout. Zero fresh
replicas does not mean reads fail entirely. Separately,
`master_as_replica_weight` can include a master alongside fresh replicas.

#### Panel 1.3: Host Health Map (Table)

Expand Down Expand Up @@ -230,7 +247,7 @@ Right Y-axis: db.pool.connections.used{host="replica-1"}

### Row 4 — Driver-Specific Panels

These panels require `register_extra_gauges()` and are only relevant
These panels require `extra_keys` in `observe_hasql_metrics()` and are relevant
for drivers that expose extra pool internals.

#### Panel 4.1: Queue Depth — psycopg3 (Time Series)
Expand Down Expand Up @@ -284,7 +301,7 @@ up to `max_overflow`. When overflow is consistently > 0, the base
| Rule | Severity | Condition | Meaning |
|---|---|---|---|
| No master | Critical | `db.pool.masters == 0` for 30s | All writes will fail |
| No replicas | Warning | `db.pool.replicas == 0` for 1m | Reads fall back to master or fail |
| No replicas | Warning | `db.pool.replicas == 0` for 1m | Reads use master/stale fallback or wait up to acquire timeout |
| Pool near exhaustion | Warning | `used / max > 0.9` for 1m | Pool running out of connections |
| Host unhealthy | Warning | `db.pool.healthy == 0` for 1m | Host lost its detected role |
| High health-check latency | Warning | `health_check.duration > 0.5s` for 2m | Host may be degrading |
Expand Down Expand Up @@ -410,12 +427,14 @@ See the per-driver OTLP scripts in `example/otlp/`:
| `aiopg_sa.py` | aiopg + SQLAlchemy | — |
| `asyncsqlalchemy.py` | SQLAlchemy async | `overflow` |

Each script creates a `PoolManager`, registers metrics, and runs a
simple workload loop so you can see data flowing into your collector.
Each script creates a `PoolManager`, samples metrics within the observation
context, and runs a workload loop. Sampling stops before pool cleanup; provider
shutdown runs off-loop. Run from the repository checkout as a module to avoid
shadowing driver packages with script filenames.

```bash
# Start the collector (e.g. Grafana Alloy, OTel Collector, etc.)
# Then run any example:
OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 \
python example/otlp/asyncpg.py --dsn postgresql://u:p@db1,db2/mydb
python -m example.otlp.asyncpg --dsn postgresql://u:p@db1,db2/mydb
```
42 changes: 22 additions & 20 deletions example/otlp/aiopg.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

Usage:
OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 \
python example/otlp/aiopg.py --dsn postgresql://u:p@db1,db2/mydb
python -m example.otlp.aiopg --dsn postgresql://u:p@db1,db2/mydb

Dependencies: hasql, aiopg, opentelemetry-sdk,
opentelemetry-exporter-otlp-proto-grpc
Expand All @@ -13,7 +13,7 @@

from hasql.driver.aiopg import PoolManager

from common import register_hasql_metrics, setup_meter_provider
from .common import observe_hasql_metrics, setup_meter_provider

parser = argparse.ArgumentParser()
parser.add_argument("--dsn", required=True, help="Multi-host PostgreSQL DSN")
Expand All @@ -27,26 +27,28 @@ async def main():

provider = setup_meter_provider(export_interval_ms=args.interval * 1000)

pool = PoolManager(
args.dsn,
fallback_master=True,
pool_factory_kwargs={"minsize": 2, "maxsize": 10},
)
await pool.ready()
register_hasql_metrics(pool)

print(f"Exporting metrics every {args.interval}s. Press Ctrl+C to stop.")
try:
while True:
async with pool.acquire_master() as conn:
async with conn.cursor() as cur:
await cur.execute("SELECT 1")
await asyncio.sleep(1)
except KeyboardInterrupt:
pass
pool = PoolManager(
args.dsn,
fallback_master=True,
pool_factory_kwargs={"minsize": 2, "maxsize": 10},
)
try:
await pool.ready()
async with observe_hasql_metrics(pool):
print(
f"Exporting metrics every {args.interval}s. "
"Press Ctrl+C to stop.",
)
while True:
async with pool.acquire_master() as conn:
async with conn.cursor() as cur:
await cur.execute("SELECT 1")
await asyncio.sleep(1)
finally:
await pool.close()
finally:
await pool.close()
provider.shutdown()
await asyncio.to_thread(provider.shutdown)


if __name__ == "__main__":
Expand Down
38 changes: 20 additions & 18 deletions example/otlp/aiopg_sa.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

Usage:
OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 \
python example/otlp/aiopg_sa.py --dsn postgresql://u:p@db1,db2/mydb
python -m example.otlp.aiopg_sa --dsn postgresql://u:p@db1,db2/mydb

Dependencies: hasql, aiopg, opentelemetry-sdk,
opentelemetry-exporter-otlp-proto-grpc
Expand All @@ -13,7 +13,7 @@

from hasql.driver.aiopg_sa import PoolManager

from common import register_hasql_metrics, setup_meter_provider
from .common import observe_hasql_metrics, setup_meter_provider

parser = argparse.ArgumentParser()
parser.add_argument("--dsn", required=True, help="Multi-host PostgreSQL DSN")
Expand All @@ -27,24 +27,26 @@ async def main():

provider = setup_meter_provider(export_interval_ms=args.interval * 1000)

pool = PoolManager(
args.dsn,
fallback_master=True,
)
await pool.ready()
register_hasql_metrics(pool)

print(f"Exporting metrics every {args.interval}s. Press Ctrl+C to stop.")
try:
while True:
async with pool.acquire_master() as conn:
await conn.execute("SELECT 1")
await asyncio.sleep(1)
except KeyboardInterrupt:
pass
pool = PoolManager(
args.dsn,
fallback_master=True,
)
try:
await pool.ready()
async with observe_hasql_metrics(pool):
print(
f"Exporting metrics every {args.interval}s. "
"Press Ctrl+C to stop.",
)
while True:
async with pool.acquire_master() as conn:
await conn.execute("SELECT 1")
await asyncio.sleep(1)
finally:
await pool.close()
finally:
await pool.close()
provider.shutdown()
await asyncio.to_thread(provider.shutdown)


if __name__ == "__main__":
Expand Down
Loading
Loading