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"}.
+
+
+ { 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.
+
+
+ { 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"