diff --git a/docs/docs/assets/img/faststream-processing-metrics.png b/docs/docs/assets/img/faststream-processing-metrics.png new file mode 100644 index 0000000..788b521 Binary files /dev/null and b/docs/docs/assets/img/faststream-processing-metrics.png differ diff --git a/docs/docs/assets/img/sqlbroker-state-metrics.png b/docs/docs/assets/img/sqlbroker-state-metrics.png new file mode 100644 index 0000000..0b50b3a Binary files /dev/null and b/docs/docs/assets/img/sqlbroker-state-metrics.png differ diff --git a/docs/docs/sqlbroker/tutorial.md b/docs/docs/sqlbroker/tutorial.md index c61c44d..d82cb03 100644 --- a/docs/docs/sqlbroker/tutorial.md +++ b/docs/docs/sqlbroker/tutorial.md @@ -228,3 +228,57 @@ And relay the messages from the database to another broker. ```python linenums="1" {!> docs_src/sqlbroker/transactional_outbox.py [ln:30-51]!} ``` + +## Observability + +FastStream already supplies Prometheus metrics for message publishing and processing rates and latencies through its +[Prometheus middleware](../getting-started/observability/prometheus.md){.external-link target="_blank"}. + +
+ ![Grafana panels from the FastStream Prometheus middleware: publish and process rates, publish and process duration percentiles, messages in process, and received message size](../assets/img/faststream-processing-metrics.png){ width="100%" .on-glb } +
+ +SQLBroker additionally provides metrics derived from the messages persisted in the database: + +- `sqlbroker_messages` — messages in the primary table, labeled by `queue` and + `state`. +- `sqlbroker_most_overdue_message_age_seconds` — how long the most overdue message has + been eligible for processing, labeled by `queue` and `state`. +- `sqlbroker_archived_messages` — messages in the archive table, labeled by + `queue` and `state`. +- `sqlbroker_state_collection_last_success_timestamp_seconds` — Unix timestamp + of the last successful database sample. + +
+ ![Grafana panels from the SQLBroker state sampler: messages by queue and state, oldest message age, and archived messages by queue and state](../assets/img/sqlbroker-state-metrics.png){ width="100%" .on-glb } +
+ +### Standalone sampler + +If the sampler runs in every broker node, each node queries the shared database and reports database-wide values and exports a duplicate copy of the same series. Prefer one standalone sampler per database, using the packaged +`sqlbroker-state-metrics` command: + +```console +pip install "faststream-sqlbroker[cli]" +sqlbroker-state-metrics \ + --host 0.0.0.0 \ + --port 8000 \ + --message-table message \ + --archive-table message_archive \ + --interval 30 \ + --database-url postgresql+asyncpg://user:pass@localhost/mydb # pragma: allowlist secret +``` + +### In-broker sampler + +The sampler can also run as part of the broker. Install the Prometheus dependency and pass the registry exposed by your metrics endpoint to `SqlBrokerStateMetricsConfig`: + +```console +pip install "faststream-sqlbroker[prometheus]" +``` + +```python linenums="1" +{!> docs_src/sqlbroker/observability_in_broker.py !} +``` + +Mount `metrics_app` at `/metrics` in your ASGI application. These database-wide gauges must not be summed across instances. This applies even to nominally single-node deployments because rolling restarts can briefly run the old and new broker nodes at the same time. diff --git a/docs/docs_src/sqlbroker/observability_in_broker.py b/docs/docs_src/sqlbroker/observability_in_broker.py new file mode 100644 index 0000000..17034cb --- /dev/null +++ b/docs/docs_src/sqlbroker/observability_in_broker.py @@ -0,0 +1,19 @@ +from prometheus_client import CollectorRegistry, make_asgi_app +from sqlalchemy.ext.asyncio import create_async_engine + +from faststream_sqlbroker import SqlBroker +from faststream_sqlbroker.sqlbroker.observability import ( + SqlBrokerStateMetricsConfig, +) + +engine = create_async_engine("postgresql+asyncpg://user:pass@localhost/mydb") +registry = CollectorRegistry() + +broker = SqlBroker( + engine=engine, + state_metrics_config=SqlBrokerStateMetricsConfig( + registry=registry, + interval=30, + ), +) +metrics_app = make_asgi_app(registry=registry) diff --git a/faststream_sqlbroker/sqlbroker/broker/broker.py b/faststream_sqlbroker/sqlbroker/broker/broker.py index 5f7749a..3e5cae7 100644 --- a/faststream_sqlbroker/sqlbroker/broker/broker.py +++ b/faststream_sqlbroker/sqlbroker/broker/broker.py @@ -32,6 +32,10 @@ from faststream.specification.schema.extra.tag import Tag, TagDict from faststream_sqlbroker.sqlbroker.client import SqlBrokerBaseClient + from faststream_sqlbroker.sqlbroker.observability import ( + SqlBrokerStateMetricsConfig, + SqlBrokerStateSampler as SqlBrokerStateSamplerType, + ) class SqlBroker( @@ -49,6 +53,7 @@ def __init__( engine: AsyncEngine, schema: SqlBrokerSchemaConfig | None = None, validate_schema_on_start: bool = True, + state_metrics_config: "SqlBrokerStateMetricsConfig | None" = None, # broker base args graceful_timeout: float | None = 15.0, decoder: Optional["CustomCallable"] = None, @@ -123,9 +128,25 @@ def __init__( ), ) + self._state_metrics_sampler: SqlBrokerStateSamplerType | None + if state_metrics_config is not None: + from faststream_sqlbroker.sqlbroker.observability.sampler import ( + SqlBrokerStateSampler, + ) + + self._state_metrics_sampler = SqlBrokerStateSampler( + engine=engine, + schema=config.schema, + config=state_metrics_config, + ) + else: + self._state_metrics_sampler = None + async def start(self) -> None: await self.connect() await super().start() + if self._state_metrics_sampler is not None: + self._state_metrics_sampler.start() async def stop( self, @@ -133,6 +154,8 @@ async def stop( exc_val: BaseException | None = None, exc_tb: Optional["TracebackType"] = None, ) -> None: + if self._state_metrics_sampler is not None: + await self._state_metrics_sampler.stop() await super().stop(exc_type, exc_val, exc_tb) if self.config.broker_config.engine: await self.config.broker_config.engine.dispose(close=True) diff --git a/faststream_sqlbroker/sqlbroker/observability/__init__.py b/faststream_sqlbroker/sqlbroker/observability/__init__.py new file mode 100644 index 0000000..c1d0de4 --- /dev/null +++ b/faststream_sqlbroker/sqlbroker/observability/__init__.py @@ -0,0 +1,23 @@ +from .metrics import SqlBrokerStateMetrics +from .queries import archived_message_state_summary_query, message_state_summary_query +from .sampler import SqlBrokerStateMetricsConfig, SqlBrokerStateSampler +from .snapshot import ( + ArchivedMessageStateSummary, + MessageStateSummary, + SqlBrokerStateSnapshot, + load_state_snapshot, + state_snapshot_from_rows, +) + +__all__ = ( + "ArchivedMessageStateSummary", + "MessageStateSummary", + "SqlBrokerStateMetrics", + "SqlBrokerStateMetricsConfig", + "SqlBrokerStateSampler", + "SqlBrokerStateSnapshot", + "archived_message_state_summary_query", + "load_state_snapshot", + "message_state_summary_query", + "state_snapshot_from_rows", +) diff --git a/faststream_sqlbroker/sqlbroker/observability/cli.py b/faststream_sqlbroker/sqlbroker/observability/cli.py new file mode 100644 index 0000000..c4a48e9 --- /dev/null +++ b/faststream_sqlbroker/sqlbroker/observability/cli.py @@ -0,0 +1,124 @@ +import asyncio +import logging +import signal +from enum import Enum +from typing import Annotated + +import typer +from prometheus_client import CollectorRegistry, start_http_server +from sqlalchemy.ext.asyncio import create_async_engine + +from faststream_sqlbroker.sqlbroker.observability.sampler import ( + SqlBrokerStateMetricsConfig, + SqlBrokerStateSampler, +) +from faststream_sqlbroker.sqlbroker.schema import SqlBrokerSchemaConfig + +cli = typer.Typer(add_completion=False, pretty_exceptions_short=True) + + +class LogLevel(str, Enum): + DEBUG = "DEBUG" + INFO = "INFO" + WARNING = "WARNING" + ERROR = "ERROR" + + +async def serve( + *, + database_url: str, + host: str, + port: int, + interval: float, + namespace: str, + message_table: str, + archive_table: str | None, + queues: list[str] | None, +) -> None: + logger = logging.getLogger(__name__) + engine = create_async_engine(database_url) + registry = CollectorRegistry() + sampler = SqlBrokerStateSampler( + engine=engine, + schema=SqlBrokerSchemaConfig( + message_table_name=message_table, + message_archive_table_name=archive_table, + ), + config=SqlBrokerStateMetricsConfig( + registry=registry, + interval=interval, + namespace=namespace, + queues=queues, + ), + logger=logger, + ) + server, thread = start_http_server(port, addr=host, registry=registry) + stopped = asyncio.Event() + loop = asyncio.get_running_loop() + for signum in (signal.SIGINT, signal.SIGTERM): + loop.add_signal_handler(signum, stopped.set) + + logger.info("Serving SQLBroker state metrics on http://%s:%d/metrics", host, port) + sampler.start() + try: + await stopped.wait() + finally: + await sampler.stop() + server.shutdown() + thread.join() + await engine.dispose() + + +@cli.command() +def state_metrics( + database_url: Annotated[ + str, + typer.Option( + help="SQLAlchemy async URL, for example postgresql+asyncpg://user:pass@host/db.", # pragma: allowlist secret + ), + ], + host: Annotated[str, typer.Option(help="Metrics listen address.")] = "127.0.0.1", + port: Annotated[int, typer.Option(help="Metrics listen port.")] = 8000, + interval: Annotated[ + float, + typer.Option(help="Database sampling interval in seconds."), + ] = 30.0, + namespace: Annotated[ + str, typer.Option(help="Prometheus metric namespace.") + ] = "sqlbroker", + message_table: Annotated[ + str, + typer.Option(help="SQLBroker primary message table name."), + ] = "message", + archive_table: Annotated[ + str, + typer.Option( + help="SQLBroker archive table name; use an empty value when disabled.", + ), + ] = "message_archive", + queue: Annotated[ + list[str] | None, + typer.Option( + help="Queue to collect; repeat to select multiple queues. All queues by default.", + ), + ] = None, + log_level: Annotated[LogLevel, typer.Option(case_sensitive=False)] = LogLevel.INFO, +) -> None: + """Expose persisted FastStream SQLBroker state as Prometheus metrics.""" + logging.basicConfig(level=log_level.value) + asyncio.run( + serve( + database_url=database_url, + host=host, + port=port, + interval=interval, + namespace=namespace, + message_table=message_table, + archive_table=archive_table or None, + queues=queue or None, + ), + ) + + +def main() -> None: + cli(prog_name="sqlbroker-state-metrics") diff --git a/faststream_sqlbroker/sqlbroker/observability/metrics.py b/faststream_sqlbroker/sqlbroker/observability/metrics.py new file mode 100644 index 0000000..04eff44 --- /dev/null +++ b/faststream_sqlbroker/sqlbroker/observability/metrics.py @@ -0,0 +1,79 @@ +from datetime import datetime, timezone + +from prometheus_client import CollectorRegistry, Gauge + +from faststream_sqlbroker.sqlbroker.observability.snapshot import ( + SqlBrokerStateSnapshot, +) + + +class SqlBrokerStateMetrics: + """Cached Prometheus metrics populated from a persisted-state snapshot.""" + + def __init__( + self, + *, + registry: CollectorRegistry, + namespace: str = "sqlbroker", + ) -> None: + self.messages = Gauge( + "messages", + "Current messages persisted in the SQLBroker primary table.", + ("queue", "state"), + namespace=namespace, + registry=registry, + ) + self.most_overdue_message_age_seconds = Gauge( + "most_overdue_message_age_seconds", + "Lag of the most overdue SQLBroker message by queue and state.", + ("queue", "state"), + namespace=namespace, + registry=registry, + ) + self.archived_messages = Gauge( + "archived_messages", + "Current messages persisted in the SQLBroker archive table.", + ("queue", "state"), + namespace=namespace, + registry=registry, + ) + self.last_success_timestamp_seconds = Gauge( + "state_collection_last_success_timestamp_seconds", + "Unix timestamp of the last successful SQLBroker state collection.", + namespace=namespace, + registry=registry, + ) + self._labels: set[tuple[str, str]] = set() + self._archive_labels: set[tuple[str, str]] = set() + + def apply(self, snapshot: SqlBrokerStateSnapshot) -> None: + current: set[tuple[str, str]] = set() + for item in snapshot.messages: + labels = (item.queue, item.state.value) + current.add(labels) + self.messages.labels(*labels).set(item.message_count) + age = max( + 0.0, + (snapshot.collected_at - item.oldest_next_attempt_at).total_seconds(), + ) + self.most_overdue_message_age_seconds.labels(*labels).set(age) + + for labels in self._labels - current: + self.messages.remove(*labels) + self.most_overdue_message_age_seconds.remove(*labels) + self._labels = current + + archive_current: set[tuple[str, str]] = set() + for archived_item in snapshot.archived_messages: + labels = (archived_item.queue, archived_item.state.value) + archive_current.add(labels) + self.archived_messages.labels(*labels).set(archived_item.message_count) + + for labels in self._archive_labels - archive_current: + self.archived_messages.remove(*labels) + self._archive_labels = archive_current + + def record_success(self, *, collected_at: datetime) -> None: + self.last_success_timestamp_seconds.set( + collected_at.replace(tzinfo=timezone.utc).timestamp() + ) diff --git a/faststream_sqlbroker/sqlbroker/observability/queries.py b/faststream_sqlbroker/sqlbroker/observability/queries.py new file mode 100644 index 0000000..4d5c526 --- /dev/null +++ b/faststream_sqlbroker/sqlbroker/observability/queries.py @@ -0,0 +1,41 @@ +from collections.abc import Collection +from typing import Any + +from sqlalchemy import Select, Table, func, select + + +def message_state_summary_query( + message_table: Table, + *, + queues: Collection[str] | None = None, +) -> Select[Any]: + """Build a portable snapshot query for persisted message states.""" + stmt = select( + message_table.c.queue.label("queue"), + message_table.c.state.label("state"), + func.count().label("message_count"), + func.min(message_table.c.next_attempt_at).label("oldest_next_attempt_at"), + ).group_by(message_table.c.queue, message_table.c.state) + + if queues is not None: + stmt = stmt.where(message_table.c.queue.in_(queues)) + + return stmt + + +def archived_message_state_summary_query( + message_archive_table: Table, + *, + queues: Collection[str] | None = None, +) -> Select[Any]: + """Build a portable count query for archived message states.""" + stmt = select( + message_archive_table.c.queue.label("queue"), + message_archive_table.c.state.label("state"), + func.count().label("message_count"), + ).group_by(message_archive_table.c.queue, message_archive_table.c.state) + + if queues is not None: + stmt = stmt.where(message_archive_table.c.queue.in_(queues)) + + return stmt diff --git a/faststream_sqlbroker/sqlbroker/observability/sampler.py b/faststream_sqlbroker/sqlbroker/observability/sampler.py new file mode 100644 index 0000000..5e18e16 --- /dev/null +++ b/faststream_sqlbroker/sqlbroker/observability/sampler.py @@ -0,0 +1,93 @@ +import asyncio +import logging +from collections.abc import Collection +from contextlib import suppress +from dataclasses import dataclass +from typing import TYPE_CHECKING + +from prometheus_client import CollectorRegistry + +from faststream_sqlbroker.sqlbroker.observability.metrics import SqlBrokerStateMetrics +from faststream_sqlbroker.sqlbroker.observability.snapshot import load_state_snapshot +from faststream_sqlbroker.sqlbroker.schema import ( + SqlBrokerSchemaConfig, + define_sqlbroker_schema, +) + +if TYPE_CHECKING: + from sqlalchemy.ext.asyncio import AsyncEngine + + +@dataclass(frozen=True, kw_only=True) +class SqlBrokerStateMetricsConfig: + registry: CollectorRegistry + interval: float = 30.0 + namespace: str = "sqlbroker" + queues: Collection[str] | None = None + + def __post_init__(self) -> None: + if self.interval <= 0: + msg = "interval must be greater than zero" + raise ValueError(msg) + + +class SqlBrokerStateSampler: + """Periodically query persisted state and populate cached Prometheus metrics.""" + + def __init__( + self, + *, + engine: "AsyncEngine", + schema: SqlBrokerSchemaConfig | None = None, + config: SqlBrokerStateMetricsConfig, + logger: logging.Logger | None = None, + ) -> None: + self._engine = engine + self._schema = define_sqlbroker_schema(config=schema or SqlBrokerSchemaConfig()) + self._interval = config.interval + self._queues = tuple(config.queues) if config.queues is not None else None + self._metrics = SqlBrokerStateMetrics( + registry=config.registry, + namespace=config.namespace, + ) + self._logger = logger or logging.getLogger(__name__) + self._stop = asyncio.Event() + self._task: asyncio.Task[None] | None = None + + @property + def metrics(self) -> SqlBrokerStateMetrics: + return self._metrics + + async def collect(self) -> None: + async with self._engine.connect() as connection: + snapshot = await load_state_snapshot( + connection, + schema=self._schema, + queues=self._queues, + ) + self._metrics.apply(snapshot) + self._metrics.record_success(collected_at=snapshot.collected_at) + + async def _run(self) -> None: + while not self._stop.is_set(): + try: + await self.collect() + except Exception: + self._logger.exception("SQLBroker state collection failed") + + with suppress(asyncio.TimeoutError): + await asyncio.wait_for(self._stop.wait(), timeout=self._interval) + + def start(self) -> None: + if self._task is not None and not self._task.done(): + return + self._stop.clear() + self._task = asyncio.create_task(self._run()) + + async def stop(self) -> None: + self._stop.set() + if self._task is not None: + self._task.cancel() + with suppress(asyncio.CancelledError): + await self._task + self._task = None diff --git a/faststream_sqlbroker/sqlbroker/observability/snapshot.py b/faststream_sqlbroker/sqlbroker/observability/snapshot.py new file mode 100644 index 0000000..9011480 --- /dev/null +++ b/faststream_sqlbroker/sqlbroker/observability/snapshot.py @@ -0,0 +1,98 @@ +from dataclasses import dataclass +from datetime import datetime, timezone +from typing import Any + +from sqlalchemy.engine import RowMapping +from sqlalchemy.ext.asyncio import AsyncConnection + +from faststream_sqlbroker.sqlbroker.message import SqlBrokerMessageState +from faststream_sqlbroker.sqlbroker.observability.queries import ( + archived_message_state_summary_query, + message_state_summary_query, +) +from faststream_sqlbroker.sqlbroker.schema import ( + SqlBrokerSchemaDefinition, + SqlBrokerSchemaType, +) + + +@dataclass(frozen=True) +class MessageStateSummary: + queue: str + state: SqlBrokerMessageState + message_count: int + oldest_next_attempt_at: datetime + + +@dataclass(frozen=True) +class ArchivedMessageStateSummary: + queue: str + state: SqlBrokerMessageState + message_count: int + + +@dataclass(frozen=True) +class SqlBrokerStateSnapshot: + collected_at: datetime + messages: tuple[MessageStateSummary, ...] + archived_messages: tuple[ArchivedMessageStateSummary, ...] = () + + +def _state(value: Any) -> SqlBrokerMessageState: + if isinstance(value, SqlBrokerMessageState): + return value + try: + return SqlBrokerMessageState(value) + except ValueError: + return SqlBrokerMessageState[value] + + +def state_snapshot_from_rows( + rows: list[RowMapping], + *, + collected_at: datetime | None = None, +) -> SqlBrokerStateSnapshot: + now = collected_at or datetime.now(timezone.utc).replace(tzinfo=None) + return SqlBrokerStateSnapshot( + collected_at=now, + messages=tuple( + MessageStateSummary( + queue=row["queue"], + state=_state(row["state"]), + message_count=row["message_count"], + oldest_next_attempt_at=row["oldest_next_attempt_at"], + ) + for row in rows + ), + ) + + +async def load_state_snapshot( + connection: AsyncConnection, + *, + schema: SqlBrokerSchemaDefinition, + queues: tuple[str, ...] | None = None, +) -> SqlBrokerStateSnapshot: + table = schema.tables[SqlBrokerSchemaType.MESSAGE] + result = await connection.execute(message_state_summary_query(table, queues=queues)) + snapshot = state_snapshot_from_rows(list(result.mappings())) + + archive_table = schema.get_table(SqlBrokerSchemaType.MESSAGE_ARCHIVE) + if archive_table is None: + return snapshot + + archive_result = await connection.execute( + archived_message_state_summary_query(archive_table, queues=queues) + ) + return SqlBrokerStateSnapshot( + collected_at=snapshot.collected_at, + messages=snapshot.messages, + archived_messages=tuple( + ArchivedMessageStateSummary( + queue=row["queue"], + state=_state(row["state"]), + message_count=row["message_count"], + ) + for row in archive_result.mappings() + ), + ) diff --git a/pyproject.toml b/pyproject.toml index a76aade..d0d0433 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -10,13 +10,23 @@ authors = [ { name = "Arseniy Popov", email = "arseniypopov@gmail.com" }, ] requires-python = ">=3.10" -version = "0.1.0a6" +version = "0.1.0a7" dependencies = [ "faststream>=0.7.0rc1", "sqlalchemy>=2.0.44", ] +[project.optional-dependencies] +prometheus = ["prometheus-client>=0.20"] +cli = [ + "prometheus-client>=0.20", + "typer>=0.12.1", +] + +[project.scripts] +sqlbroker-state-metrics = "faststream_sqlbroker.sqlbroker.observability.cli:main" + [tool.uv.build-backend] module-root = "." module-name = "faststream_sqlbroker" @@ -39,6 +49,8 @@ test = [ "asyncmy>=0.2.10", "aiosqlite>=0.22.1", "cryptography>=43.0.0", + "prometheus-client>=0.20", + "typer>=0.12.1", ] docs = [ diff --git a/tests/docs/sqlbroker/test_observability.py b/tests/docs/sqlbroker/test_observability.py new file mode 100644 index 0000000..8812395 --- /dev/null +++ b/tests/docs/sqlbroker/test_observability.py @@ -0,0 +1,13 @@ +from faststream_sqlbroker import SqlBroker + + +def test_in_broker_observability() -> None: + from docs.docs_src.sqlbroker.observability_in_broker import ( + broker, + engine, + metrics_app, + ) + + assert isinstance(broker, SqlBroker) + assert engine.dialect.name == "postgresql" + assert callable(metrics_app) diff --git a/tests/test_observability.py b/tests/test_observability.py new file mode 100644 index 0000000..1cfe37e --- /dev/null +++ b/tests/test_observability.py @@ -0,0 +1,216 @@ +import asyncio +from datetime import datetime, timedelta, timezone + +import pytest +from prometheus_client import CollectorRegistry +from sqlalchemy import insert +from sqlalchemy.ext.asyncio import AsyncEngine + +from faststream_sqlbroker.sqlbroker.broker.broker import SqlBroker +from faststream_sqlbroker.sqlbroker.message import SqlBrokerMessageState +from faststream_sqlbroker.sqlbroker.observability import ( + ArchivedMessageStateSummary, + MessageStateSummary, + SqlBrokerStateMetrics, + SqlBrokerStateMetricsConfig, + SqlBrokerStateSampler, + SqlBrokerStateSnapshot, +) +from faststream_sqlbroker.sqlbroker.schema import ( + SqlBrokerSchemaConfig, + SqlBrokerSchemaType, + define_sqlbroker_schema, +) + + +def _sample(registry: CollectorRegistry, name: str, labels: dict[str, str]) -> float: + value = registry.get_sample_value(name, labels) + assert value is not None + return value + + +async def _wait_for_sample( + registry: CollectorRegistry, + name: str, + labels: dict[str, str], + expected: float, +) -> None: + for _ in range(200): + if registry.get_sample_value(name, labels) == expected: + return + await asyncio.sleep(0.01) + assert registry.get_sample_value(name, labels) == expected + + +@pytest.mark.asyncio() +async def test_snapshot_and_sampler_populate_persisted_state( + engine: AsyncEngine, recreate_tables: None +) -> None: + schema_config = SqlBrokerSchemaConfig() + schema = define_sqlbroker_schema(config=schema_config) + table = schema.tables[SqlBrokerSchemaType.MESSAGE] + archive_table = schema.tables[SqlBrokerSchemaType.MESSAGE_ARCHIVE] + now = datetime.now(timezone.utc).replace(tzinfo=None) + created_at = now - timedelta(seconds=60) + next_attempt_at = now - timedelta(seconds=10) + async with engine.begin() as connection: + await connection.execute( + insert(table), + [ + { + "queue": "orders", + "payload": b"1", + "state": SqlBrokerMessageState.PENDING, + "created_at": created_at, + "next_attempt_at": next_attempt_at, + "attempts_count": 0, + "deliveries_count": 0, + }, + { + "queue": "orders", + "payload": b"2", + "state": SqlBrokerMessageState.PROCESSING, + "created_at": created_at, + "next_attempt_at": next_attempt_at, + "attempts_count": 0, + "deliveries_count": 1, + }, + ], + ) + await connection.execute( + insert(archive_table), + { + "queue": "orders", + "payload": b"3", + "state": SqlBrokerMessageState.COMPLETED, + "created_at": created_at, + "attempts_count": 1, + "deliveries_count": 1, + }, + ) + + registry = CollectorRegistry() + sampler = SqlBrokerStateSampler( + engine=engine, + schema=schema_config, + config=SqlBrokerStateMetricsConfig(registry=registry), + ) + await sampler.collect() + + assert ( + _sample(registry, "sqlbroker_messages", {"queue": "orders", "state": "pending"}) + == 1 + ) + assert ( + _sample( + registry, "sqlbroker_messages", {"queue": "orders", "state": "processing"} + ) + == 1 + ) + assert ( + _sample( + registry, + "sqlbroker_archived_messages", + {"queue": "orders", "state": "completed"}, + ) + == 1 + ) + age = _sample( + registry, + "sqlbroker_most_overdue_message_age_seconds", + {"queue": "orders", "state": "pending"}, + ) + assert 9 <= age < 20 + assert _sample( + registry, "sqlbroker_state_collection_last_success_timestamp_seconds", {} + ) + + +@pytest.mark.asyncio() +async def test_running_consuming_broker_updates_state_metrics( + engine: AsyncEngine, recreate_tables: None +) -> None: + registry = CollectorRegistry() + broker = SqlBroker( + engine=engine, + state_metrics_config=SqlBrokerStateMetricsConfig( + registry=registry, interval=0.01 + ), + ) + handler_started = asyncio.Event() + release_handler = asyncio.Event() + + @broker.subscriber( + queues=["orders"], + max_workers=1, + max_fetch_interval=0.01, + min_fetch_interval=0.01, + fetch_batch_size=1, + flush_interval=0.01, + ) + async def handler() -> None: + handler_started.set() + await release_handler.wait() + + await broker.connect() + await broker.publish({"order_id": 1}, queue="orders") + await broker.start() + try: + await asyncio.wait_for(handler_started.wait(), timeout=2) + await _wait_for_sample( + registry, + "sqlbroker_messages", + {"queue": "orders", "state": "processing"}, + 1, + ) + + release_handler.set() + await _wait_for_sample( + registry, + "sqlbroker_archived_messages", + {"queue": "orders", "state": "completed"}, + 1, + ) + finally: + await broker.stop() + + +def test_metrics_remove_stale_state_series() -> None: + registry = CollectorRegistry() + metrics = SqlBrokerStateMetrics(registry=registry) + + now = datetime.now(timezone.utc).replace(tzinfo=None) + metrics.apply( + SqlBrokerStateSnapshot( + collected_at=now, + messages=( + MessageStateSummary( + queue="orders", + state=SqlBrokerMessageState.PENDING, + message_count=1, + oldest_next_attempt_at=now, + ), + ), + archived_messages=( + ArchivedMessageStateSummary( + queue="orders", + state=SqlBrokerMessageState.COMPLETED, + message_count=1, + ), + ), + ) + ) + metrics.apply(SqlBrokerStateSnapshot(collected_at=now, messages=())) + + assert ( + registry.get_sample_value( + "sqlbroker_messages", {"queue": "orders", "state": "pending"} + ) + is None + ) + assert ( + registry.get_sample_value( + "sqlbroker_archived_messages", {"queue": "orders", "state": "completed"} + ) + is None + ) diff --git a/tests/test_observability_cli.py b/tests/test_observability_cli.py new file mode 100644 index 0000000..fd54dde --- /dev/null +++ b/tests/test_observability_cli.py @@ -0,0 +1,193 @@ +import asyncio +import signal +import socket +import sys +import urllib.error +import urllib.request +from datetime import datetime, timezone +from typing import Any + +import pytest +from sqlalchemy import insert +from sqlalchemy.ext.asyncio import AsyncEngine +from typer.testing import CliRunner + +from faststream_sqlbroker.sqlbroker.message import SqlBrokerMessageState +from faststream_sqlbroker.sqlbroker.observability import cli as cli_module +from faststream_sqlbroker.sqlbroker.schema import ( + SqlBrokerSchemaConfig, + SqlBrokerSchemaType, + define_sqlbroker_schema, +) + +runner = CliRunner() + + +@pytest.fixture() +def served(monkeypatch: pytest.MonkeyPatch) -> list[dict[str, Any]]: + """Capture the arguments the CLI would serve with, without starting the sampler.""" + calls: list[dict[str, Any]] = [] + + async def fake_serve(**kwargs: Any) -> None: + calls.append(kwargs) + + monkeypatch.setattr(cli_module, "serve", fake_serve) + return calls + + +def test_cli_defaults(served: list[dict[str, Any]]) -> None: + result = runner.invoke( + cli_module.cli, + ["--database-url", "sqlite+aiosqlite:///broker.db"], + ) + + assert result.exit_code == 0, result.output + assert served == [ + { + "database_url": "sqlite+aiosqlite:///broker.db", + "host": "127.0.0.1", + "port": 8000, + "interval": 30.0, + "namespace": "sqlbroker", + "message_table": "message", + "archive_table": "message_archive", + "queues": None, + }, + ] + + +def test_cli_queue_filters_are_repeatable(served: list[dict[str, Any]]) -> None: + result = runner.invoke( + cli_module.cli, + [ + "--database-url", + "sqlite+aiosqlite:///broker.db", + "--queue", + "orders", + "--queue", + "email", + ], + ) + + assert result.exit_code == 0, result.output + assert served[0]["queues"] == ["orders", "email"] + + +def test_cli_empty_archive_table_disables_archive_metrics( + served: list[dict[str, Any]], +) -> None: + result = runner.invoke( + cli_module.cli, + ["--database-url", "sqlite+aiosqlite:///broker.db", "--archive-table", ""], + ) + + assert result.exit_code == 0, result.output + assert served[0]["archive_table"] is None + + +def test_cli_requires_database_url(served: list[dict[str, Any]]) -> None: + result = runner.invoke(cli_module.cli, []) + + assert result.exit_code != 0 + assert not served + + +def _free_tcp_port() -> int: + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: + sock.bind(("127.0.0.1", 0)) + return sock.getsockname()[1] + + +def _fetch(url: str) -> str: + with urllib.request.urlopen(url, timeout=1) as response: + return response.read().decode() + + +async def _wait_for_metrics( + url: str, *, timeout: float, expected: str | None = None +) -> str: + """Poll the metrics endpoint until it is reachable and (optionally) contains `expected`. + + The CLI starts the HTTP server before its first background sample completes, so an + early scrape can observe an endpoint that is reachable but not yet populated. + """ + loop = asyncio.get_running_loop() + deadline = loop.time() + timeout + last_error: Exception | None = None + body = "" + while loop.time() < deadline: + try: + body = await asyncio.to_thread(_fetch, url) + if expected is None or expected in body: + return body + except (urllib.error.URLError, ConnectionError) as exc: + last_error = exc + await asyncio.sleep(0.05) + if last_error is not None and not body: + msg = f"metrics endpoint at {url} never became ready: {last_error}" + raise AssertionError(msg) + return body + + +@pytest.mark.asyncio() +async def test_cli_serves_real_metrics_over_http_and_shuts_down_gracefully( + engine: AsyncEngine, recreate_tables: None +) -> None: + """Run the CLI as a real subprocess against a real database end to end.""" + schema = define_sqlbroker_schema(config=SqlBrokerSchemaConfig()) + table = schema.tables[SqlBrokerSchemaType.MESSAGE] + now = datetime.now(timezone.utc).replace(tzinfo=None) + async with engine.begin() as connection: + await connection.execute( + insert(table), + { + "queue": "orders", + "payload": b"1", + "state": SqlBrokerMessageState.PENDING, + "created_at": now, + "next_attempt_at": now, + "attempts_count": 0, + "deliveries_count": 0, + }, + ) + + database_url = engine.url.render_as_string(hide_password=False) + port = _free_tcp_port() + + process = await asyncio.create_subprocess_exec( + sys.executable, + "-c", + "from faststream_sqlbroker.sqlbroker.observability.cli import cli; cli()", + "--database-url", + database_url, + "--host", + "127.0.0.1", + "--port", + str(port), + "--interval", + "0.05", + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.STDOUT, + ) + try: + expected = 'sqlbroker_messages{queue="orders",state="pending"} 1.0' + body = await _wait_for_metrics( + f"http://127.0.0.1:{port}/metrics", timeout=10, expected=expected + ) + assert expected in body + + process.send_signal(signal.SIGINT) + try: + await asyncio.wait_for(process.wait(), timeout=5) + except asyncio.TimeoutError as exc: + process.kill() + await process.wait() + msg = "CLI did not shut down gracefully after SIGINT" + raise AssertionError(msg) from exc + finally: + if process.returncode is None: + process.kill() + await process.wait() + + output = (await process.stdout.read()).decode() if process.stdout else "" + assert process.returncode == 0, output diff --git a/uv.lock b/uv.lock index 91c53ef..4b3fa5f 100644 --- a/uv.lock +++ b/uv.lock @@ -54,6 +54,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/00/b7/e3bf5133d697a08128598c8d0abc5e16377b51465a33756de24fa7dee953/aiosqlite-0.22.1-py3-none-any.whl", hash = "sha256:21c002eb13823fad740196c5a2e9d8e62f6243bd9e7e4a1f87fb5e44ecb4fceb", size = 17405, upload-time = "2025-12-23T19:25:42.139Z" }, ] +[[package]] +name = "annotated-doc" +version = "0.0.4" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/57/ba/046ceea27344560984e26a590f90bc7f4a75b06701f653222458922b558c/annotated_doc-0.0.4.tar.gz", hash = "sha256:fbcda96e87e9c92ad167c2e53839e57503ecfda18804ea28102353485033faa4", size = 7288, upload-time = "2025-11-10T22:07:42.062Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/1e/d3/26bf1008eb3d2daa8ef4cacc7f3bfdc11818d111f7e2d0201bc6e3b49d45/annotated_doc-0.0.4-py3-none-any.whl", hash = "sha256:571ac1dc6991c450b25a9c2d84a3705e2ae7a53467b5d111c24fa8baabbed320", size = 5303, upload-time = "2025-11-10T22:07:40.673Z" }, +] + [[package]] name = "annotated-types" version = "0.7.0" @@ -708,13 +717,22 @@ wheels = [ [[package]] name = "faststream-sqlbroker" -version = "0.1.0a6" +version = "0.1.0a7" source = { editable = "." } dependencies = [ { name = "faststream" }, { name = "sqlalchemy" }, ] +[package.optional-dependencies] +cli = [ + { name = "prometheus-client" }, + { name = "typer" }, +] +prometheus = [ + { name = "prometheus-client" }, +] + [package.dev-dependencies] dev = [ { name = "aiokafka" }, @@ -730,6 +748,7 @@ dev = [ { name = "mkdocstrings", extra = ["python"] }, { name = "mypy" }, { name = "pre-commit" }, + { name = "prometheus-client" }, { name = "pymdown-extensions" }, { name = "pytest" }, { name = "pytest-asyncio" }, @@ -737,6 +756,7 @@ dev = [ { name = "pytest-timeout" }, { name = "pytest-xdist" }, { name = "ruff" }, + { name = "typer" }, ] docs = [ { name = "mdx-include" }, @@ -757,18 +777,24 @@ test = [ { name = "asyncpg" }, { name = "cryptography" }, { name = "dirty-equals" }, + { name = "prometheus-client" }, { name = "pytest" }, { name = "pytest-asyncio" }, { name = "pytest-cov" }, { name = "pytest-timeout" }, { name = "pytest-xdist" }, + { name = "typer" }, ] [package.metadata] requires-dist = [ { name = "faststream", specifier = ">=0.7.0rc1" }, + { name = "prometheus-client", marker = "extra == 'cli'", specifier = ">=0.20" }, + { name = "prometheus-client", marker = "extra == 'prometheus'", specifier = ">=0.20" }, { name = "sqlalchemy", specifier = ">=2.0.44" }, + { name = "typer", marker = "extra == 'cli'", specifier = ">=0.12.1" }, ] +provides-extras = ["prometheus", "cli"] [package.metadata.requires-dev] dev = [ @@ -785,6 +811,7 @@ dev = [ { name = "mkdocstrings", extras = ["python"], specifier = ">=0.27.0" }, { name = "mypy", specifier = "==1.19.1" }, { name = "pre-commit", specifier = "==4.5.1" }, + { name = "prometheus-client", specifier = ">=0.20" }, { name = "pymdown-extensions", specifier = ">=10.0" }, { name = "pytest", specifier = "==9.0.3" }, { name = "pytest-asyncio", specifier = "==1.3.0" }, @@ -792,6 +819,7 @@ dev = [ { name = "pytest-timeout", specifier = ">=2.4.0" }, { name = "pytest-xdist", specifier = ">=3.8.0" }, { name = "ruff", specifier = "==0.14.14" }, + { name = "typer", specifier = ">=0.12.1" }, ] docs = [ { name = "mdx-include", specifier = ">=1.4.2" }, @@ -812,11 +840,13 @@ test = [ { name = "asyncpg", specifier = ">=0.31.0" }, { name = "cryptography", specifier = ">=43.0.0" }, { name = "dirty-equals", specifier = "==0.11" }, + { name = "prometheus-client", specifier = ">=0.20" }, { name = "pytest", specifier = "==9.0.3" }, { name = "pytest-asyncio", specifier = "==1.3.0" }, { name = "pytest-cov", specifier = ">=5.0.0" }, { name = "pytest-timeout", specifier = ">=2.4.0" }, { name = "pytest-xdist", specifier = ">=3.8.0" }, + { name = "typer", specifier = ">=0.12.1" }, ] [[package]] @@ -1066,6 +1096,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/de/1f/77fa3081e4f66ca3576c896ae5d31c3002ac6607f9747d2e3aa49227e464/markdown-3.10.2-py3-none-any.whl", hash = "sha256:e91464b71ae3ee7afd3017d9f358ef0baf158fd9a298db92f1d4761133824c36", size = 108180, upload-time = "2026-02-09T14:57:25.787Z" }, ] +[[package]] +name = "markdown-it-py" +version = "4.2.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "mdurl" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/06/ff/7841249c247aa650a76b9ee4bbaeae59370dc8bfd2f6c01f3630c35eb134/markdown_it_py-4.2.0.tar.gz", hash = "sha256:04a21681d6fbb623de53f6f364d352309d4094dd4194040a10fd51833e418d49", size = 82454, upload-time = "2026-05-07T12:08:28.36Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/b3/81/4da04ced5a082363ecfa159c010d200ecbd959ae410c10c0264a38cac0f5/markdown_it_py-4.2.0-py3-none-any.whl", hash = "sha256:9f7ebbcd14fe59494226453aed97c1070d83f8d24b6fc3a3bcf9a38092641c4a", size = 91687, upload-time = "2026-05-07T12:08:27.182Z" }, +] + [[package]] name = "markupsafe" version = "3.0.3" @@ -1151,6 +1193,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/70/bc/6f1c2f612465f5fa89b95bead1f44dcb607670fd42891d8fdcd5d039f4f4/markupsafe-3.0.3-cp314-cp314t-win_arm64.whl", hash = "sha256:32001d6a8fc98c8cb5c947787c5d08b0a50663d139f1305bac5885d98d9b40fa", size = 14146, upload-time = "2025-09-27T18:37:28.327Z" }, ] +[[package]] +name = "mdurl" +version = "0.1.2" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/d6/54/cfe61301667036ec958cb99bd3efefba235e65cdeb9c84d24a8293ba1d90/mdurl-0.1.2.tar.gz", hash = "sha256:bb413d29f5eea38f31dd4754dd7377d4465116fb207585f97bf925588687c1ba", size = 8729, upload-time = "2022-08-14T12:40:10.846Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/b3/38/89ba8ad64ae25be8de66a6d463314cf1eb366222074cfda9ee839c56a4b4/mdurl-0.1.2-py3-none-any.whl", hash = "sha256:84008a41e51615a49fc9966191ff91509e3c40b939176e643fd50a5c2196b8f8", size = 9979, upload-time = "2022-08-14T12:40:09.779Z" }, +] + [[package]] name = "mdx-include" version = "1.4.2" @@ -1446,6 +1497,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/5d/19/fd3ef348460c80af7bb4669ea7926651d1f95c23ff2df18b9d24bab4f3fa/pre_commit-4.5.1-py2.py3-none-any.whl", hash = "sha256:3b3afd891e97337708c1674210f8eba659b52a38ea5f822ff142d10786221f77", size = 226437, upload-time = "2025-12-16T21:14:32.409Z" }, ] +[[package]] +name = "prometheus-client" +version = "0.25.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/1b/fb/d9aa83ffe43ce1f19e557c0971d04b90561b0cfd50762aafb01968285553/prometheus_client-0.25.0.tar.gz", hash = "sha256:5e373b75c31afb3c86f1a52fa1ad470c9aace18082d39ec0d2f918d11cc9ba28", size = 86035, upload-time = "2026-04-09T19:53:42.359Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/8d/9b/d4b1e644385499c8346fa9b622a3f030dce14cd6ef8a1871c221a17a67e7/prometheus_client-0.25.0-py3-none-any.whl", hash = "sha256:d5aec89e349a6ec230805d0df882f3807f74fd6c1a2fa86864e3c2279059fed1", size = 64154, upload-time = "2026-04-09T19:53:41.324Z" }, +] + [[package]] name = "pycparser" version = "3.0" @@ -1804,6 +1864,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/a0/f4/c67b0b3f1b9245e8d266f0f112c500d50e5b4e83cb6f3b71b6528104182a/requests-2.34.2-py3-none-any.whl", hash = "sha256:2a0d60c172f83ac6ab31e4554906c0f3b3588d37b5cb939b1c061f4907e278e0", size = 73075, upload-time = "2026-05-14T19:25:26.443Z" }, ] +[[package]] +name = "rich" +version = "15.0.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "markdown-it-py" }, + { name = "pygments" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/c0/8f/0722ca900cc807c13a6a0c696dacf35430f72e0ec571c4275d2371fca3e9/rich-15.0.0.tar.gz", hash = "sha256:edd07a4824c6b40189fb7ac9bc4c52536e9780fbbfbddf6f1e2502c31b068c36", size = 230680, upload-time = "2026-04-12T08:24:00.75Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/82/3b/64d4899d73f91ba49a8c18a8ff3f0ea8f1c1d75481760df8c68ef5235bf5/rich-15.0.0-py3-none-any.whl", hash = "sha256:33bd4ef74232fb73fe9279a257718407f169c09b78a87ad3d296f548e27de0bb", size = 310654, upload-time = "2026-04-12T08:24:02.83Z" }, +] + [[package]] name = "ruff" version = "0.14.14" @@ -1892,6 +1965,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/1b/02/e511474facfd324529509626097f981986b7adc935fb4eb436c0812c540f/selectolax-0.4.10-cp314-cp314t-win_arm64.whl", hash = "sha256:68e1ef717b47f5cdcd1b151b2176d7184c38cb3772f8964509979840575203ef", size = 1944297, upload-time = "2026-05-26T15:42:50.906Z" }, ] +[[package]] +name = "shellingham" +version = "1.5.4" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/58/15/8b3609fd3830ef7b27b655beb4b4e9c62313a4e8da8c676e142cc210d58e/shellingham-1.5.4.tar.gz", hash = "sha256:8dbca0739d487e5bd35ab3ca4b36e11c4078f3a234bfce294b0a0291363404de", size = 10310, upload-time = "2023-10-24T04:13:40.426Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/e0/f9/0595336914c5619e5f28a1fb793285925a8cd4b432c9da0a987836c7f822/shellingham-1.5.4-py2.py3-none-any.whl", hash = "sha256:7ecfff8f2fd72616f7481040475a65b2bf8af90a56c89140852d1120324e8686", size = 9755, upload-time = "2023-10-24T04:13:38.866Z" }, +] + [[package]] name = "six" version = "1.17.0" @@ -2015,6 +2097,21 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/7b/61/cceae43728b7de99d9b847560c262873a1f6c98202171fd5ed62640b494b/tomli-2.4.1-py3-none-any.whl", hash = "sha256:0d85819802132122da43cb86656f8d1f8c6587d54ae7dcaf30e90533028b49fe", size = 14583, upload-time = "2026-03-25T20:22:03.012Z" }, ] +[[package]] +name = "typer" +version = "0.26.8" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "annotated-doc" }, + { name = "colorama", marker = "sys_platform == 'win32'" }, + { name = "rich" }, + { name = "shellingham" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/7c/f7/68adc395201b20b872d68e975386832e8005ffeacedd43a1d837a32815be/typer-0.26.8.tar.gz", hash = "sha256:c244a6bd558886fe3f8780efb6bdd28bb9aff005a94eedebaa5cb32926fe2f7e", size = 202097, upload-time = "2026-06-26T09:22:45.705Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/80/87/b9fd69c92c6102a066e1b86a35243f53e70bd4c709f2a26d9f4fee4f4dc0/typer-0.26.8-py3-none-any.whl", hash = "sha256:3512ca79ac5c11113414b36e80281b872884477722440691c89d1112e321a49c", size = 122564, upload-time = "2026-06-26T09:22:44.72Z" }, +] + [[package]] name = "typing-extensions" version = "4.15.0"