feat(usage): record per-tenant usage events into the app DB

Deck #67 data-plane slice: tenant Pods record billable operations
(embedding queries, pages/chunks embedded) into an app-DB usage_events
table that the control plane later pulls read-only into the billing
ledger and syncs to Stripe Meter Events.

- migration 007: usage_events table (Postgres TIMESTAMPTZ/JSONB/UUID
  with portable SQLite fallbacks), indexed (occurred_at, metric) for the
  CP rollup's per-day range scan + GROUP BY metric.
- UsageEventStore: best-effort, flag-gated writer reusing the shared
  RefreshTokenStorage engine; ON CONFLICT (event_id) DO NOTHING for
  idempotent retries; dialect-branched occurred_at bind. All work
  (incl. metadata JSON encode) is swallowed so a metering failure never
  surfaces to the user op.
- USAGE_METERING_ENABLED flag (default off) wired through Settings +
  env map; off-path touches no storage, so OSS self-hosters get an empty
  table and zero write overhead.
- two recording hooks: embeddings_queries (per nc_semantic_search, which
  nc_semantic_search_answer reuses) and pages_chunks (after dense
  embedding succeeds, covering both in-process and procrastinate paths).
- storage.acquire()/.dialect public seams so the sibling store doesn't
  reach into the underscored internal.
- tests parametrized over SQLite + Postgres: flag-off no-op, roundtrip,
  ON CONFLICT dedup, JSON/NULL metadata, and the best-effort swallow of
  both DB errors and unserializable metadata.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Chris Coutinho
2026-06-07 15:13:14 +02:00
co-authored by Claude Opus 4.8
parent 42505f8f87
commit 1c6b1a84ea
9 changed files with 491 additions and 1 deletions
@@ -0,0 +1,78 @@
"""Add usage_events table for per-tenant usage metering.
Deck #67 / control-plane usage-metering.md (pull model). Each tenant Pod
records billable usage (embedding queries, pages/chunks embedded) into its own
app DB; the control plane later pulls this table read-only into the billing
ledger and syncs to Stripe Meter Events. Writes are gated by
``USAGE_METERING_ENABLED`` (default off) — the recording hook is a no-op when
the flag is off, so an OSS self-hoster gets an empty table and zero write
overhead.
Unlike the rest of this schema (unix-epoch ``BigInteger`` timestamps, JSON as
``Text``), this table uses real Postgres ``TIMESTAMPTZ``/``JSONB``/``UUID``
because the control-plane rollup runs ``date_trunc('day', occurred_at AT TIME
ZONE 'UTC')`` and ``GROUP BY day, metric`` directly against Postgres, which
requires a genuine timestamptz column. SQLite (OSS/tests) uses portable
fallbacks (``TEXT``/``TIMESTAMP``); the control plane never queries SQLite.
Revision ID: 007
Revises: 006
Create Date: 2026-06-10 12:00:00.000000
"""
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql
from alembic import op
revision = "007"
down_revision = "006"
branch_labels = None
depends_on = None
def upgrade() -> None:
is_pg = op.get_bind().dialect.name == "postgresql"
op.create_table(
"usage_events",
# Pod-generated idempotency key. UUID on Postgres; TEXT on SQLite,
# which has no native UUID type. Stored/bound as a plain string in
# both backends (see usage/store.py).
sa.Column(
"event_id",
postgresql.UUID(as_uuid=False) if is_pg else sa.Text,
primary_key=True,
),
# Operation completion time (UTC). Real TIMESTAMPTZ on Postgres so the
# CP rollup's date_trunc(... AT TIME ZONE 'UTC') works; portable
# TIMESTAMP on SQLite (stored as ISO text, queryable in tests).
sa.Column(
"occurred_at",
postgresql.TIMESTAMP(timezone=True) if is_pg else sa.TIMESTAMP,
nullable=False,
),
# Catalog metric: 'embeddings_queries' or 'pages_chunks'.
sa.Column("metric", sa.Text, nullable=False),
sa.Column("value", sa.BigInteger, nullable=False),
# Rawest unit per request (provider, model, tokens, doc_type, ...).
# JSONB on Postgres so the CP can slice on dimensions later; TEXT
# (json.dumps) on SQLite.
sa.Column(
"metadata",
postgresql.JSONB if is_pg else sa.Text,
nullable=True,
),
)
# Serves the CP rollup's per-day range scan (occurred_at >= / <) plus the
# GROUP BY metric; leading occurred_at makes the range filter index-usable.
op.create_index(
"idx_usage_events_occurred_metric",
"usage_events",
["occurred_at", "metric"],
)
def downgrade() -> None:
op.drop_index("idx_usage_events_occurred_metric", table_name="usage_events")
op.drop_table("usage_events")
+16 -1
View File
@@ -39,7 +39,7 @@ import os
import socket import socket
import sqlite3 import sqlite3
import time import time
from contextlib import asynccontextmanager from contextlib import AbstractAsyncContextManager, asynccontextmanager
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
@@ -666,6 +666,21 @@ class RefreshTokenStorage:
async with self.engine.connect() as conn: async with self.engine.connect() as conn:
yield _DBConn(conn) yield _DBConn(conn)
def acquire(self) -> AbstractAsyncContextManager["_DBConn"]:
"""Public alias for :meth:`_db`: a backend-agnostic connection cm.
Lets sibling stores (e.g. :class:`UsageEventStore`) reuse this
instance's engine, NullPool, and ``_DBConn`` shim without reaching
into the underscored internal. Use as ``async with storage.acquire()
as db:``.
"""
return self._db()
@property
def dialect(self) -> str:
"""Backend dialect name ("sqlite" / "postgresql"), or "unknown" pre-init."""
return self._dialect
async def store_refresh_token( async def store_refresh_token(
self, self,
user_id: str, user_id: str,
+10
View File
@@ -244,6 +244,11 @@ _DEFAULTS: dict[str, Any] = {
# this before a real ACL backfill would silently drop legacy results. # this before a real ACL backfill would silently drop legacy results.
# verify-on-read remains the correctness backstop regardless. # verify-on-read remains the correctness backstop regardless.
"acl_prefilter_enabled": False, "acl_prefilter_enabled": False,
# Usage metering (Deck #67, control-plane usage-metering.md). OFF by
# default so OSS self-hosters don't accrue a metering table or write
# overhead; Astrolabe Cloud provisioning sets it true. When on, billable
# ops record rows into the app-DB usage_events table (best-effort).
"usage_metering_enabled": False,
} }
@@ -829,6 +834,10 @@ class Settings:
embedding_gateway_scope: str | None = None embedding_gateway_scope: str | None = None
tenant_id: str | None = None # per-tenant identity (UUID form) tenant_id: str | None = None # per-tenant identity (UUID form)
acl_prefilter_enabled: bool = False # query-side ACL pre-filter (§11); OFF acl_prefilter_enabled: bool = False # query-side ACL pre-filter (§11); OFF
# Usage metering (Deck #67); OFF by default. When true, billable ops
# record best-effort rows into the app-DB usage_events table for the
# control plane to pull. See nextcloud_mcp_server/usage/store.py.
usage_metering_enabled: bool = False
def __post_init__(self): def __post_init__(self):
"""Validate configuration and set defaults.""" """Validate configuration and set defaults."""
@@ -1441,6 +1450,7 @@ def get_settings() -> Settings:
"embedding_gateway_scope": "EMBEDDING_GATEWAY_SCOPE", "embedding_gateway_scope": "EMBEDDING_GATEWAY_SCOPE",
"tenant_id": "TENANT_ID", "tenant_id": "TENANT_ID",
"acl_prefilter_enabled": "ACL_PREFILTER_ENABLED", "acl_prefilter_enabled": "ACL_PREFILTER_ENABLED",
"usage_metering_enabled": "USAGE_METERING_ENABLED",
} }
# Only pass values that dynaconf actually has; omit unset keys so # Only pass values that dynaconf actually has; omit unset keys so
+24
View File
@@ -39,6 +39,7 @@ from nextcloud_mcp_server.search.access_filter import (
from nextcloud_mcp_server.search.bm25_hybrid import BM25HybridSearchAlgorithm from nextcloud_mcp_server.search.bm25_hybrid import BM25HybridSearchAlgorithm
from nextcloud_mcp_server.search.context import get_chunk_with_context from nextcloud_mcp_server.search.context import get_chunk_with_context
from nextcloud_mcp_server.search.verification import verify_search_results from nextcloud_mcp_server.search.verification import verify_search_results
from nextcloud_mcp_server.usage import UsageEventStore
from nextcloud_mcp_server.utils.validation import parse_modified_timestamp from nextcloud_mcp_server.utils.validation import parse_modified_timestamp
from nextcloud_mcp_server.vector.metrics_publisher import count_indexed from nextcloud_mcp_server.vector.metrics_publisher import count_indexed
from nextcloud_mcp_server.vector.qdrant_client import get_qdrant_client from nextcloud_mcp_server.vector.qdrant_client import get_qdrant_client
@@ -517,6 +518,29 @@ def configure_semantic_tools(mcp: FastMCP):
logger.info("Returning %d results from BM25 hybrid search", len(results)) logger.info("Returning %d results from BM25 hybrid search", len(results))
# Usage metering (Deck #67): one billable 'embeddings_queries'
# event per successful search (the query embedding is the metered
# cost). Best-effort and gated on the flag so the off-path touches
# no storage. nc_semantic_search_answer reuses this tool, so it
# records here too — do not add a second hook there.
if settings.usage_metering_enabled:
try:
store = await UsageEventStore.shared()
await store.record_usage_event(
metric="embeddings_queries",
value=1,
metadata={
"user_id": username,
"fusion": fusion,
"doc_types": doc_types,
},
)
except Exception:
logger.debug(
"usage metering hook (embeddings_queries) skipped",
exc_info=True,
)
return SemanticSearchResponse( return SemanticSearchResponse(
results=results, results=results,
query=query, query=query,
+10
View File
@@ -0,0 +1,10 @@
"""Per-tenant usage metering (data plane).
Records billable operations into the app-DB ``usage_events`` table for the
control plane to pull. Gated by ``USAGE_METERING_ENABLED`` (default off). See
Deck #67 and control-plane ``usage-metering.md``.
"""
from nextcloud_mcp_server.usage.store import UsageEventStore
__all__ = ["UsageEventStore"]
+127
View File
@@ -0,0 +1,127 @@
"""Best-effort usage-event recording for per-tenant metering (Deck #67).
A tenant Pod records billable operations (embedding queries, pages/chunks
embedded) into the app-DB ``usage_events`` table; the control plane later pulls
that table read-only into the billing ledger and syncs to Stripe Meter Events
(see control-plane ``usage-metering.md``). This module owns only the data-plane
recording side.
Design contract:
- **Flag-gated.** Writes are a no-op unless ``USAGE_METERING_ENABLED`` is true,
so OSS self-hosters and unmetered deployments do zero DB work.
- **Best-effort.** A metering-write failure is logged and dropped, never raised
into the user-facing operation. ``ON CONFLICT (event_id) DO NOTHING`` makes a
retried write a no-op.
- **Engine reuse.** Rather than opening its own engine, this store borrows the
process-wide :class:`RefreshTokenStorage` singleton (``get_shared_storage()``)
— same app DB, NullPool, dialect handling, and ``_DBConn`` shim. The shared
storage guarantees Alembic migrations (incl. ``usage_events``) already ran.
"""
import json
import logging
import time
import uuid
from datetime import datetime, timezone
from typing import Any
from nextcloud_mcp_server.auth.storage import RefreshTokenStorage, get_shared_storage
from nextcloud_mcp_server.config import get_settings
from nextcloud_mcp_server.observability.metrics import record_db_operation
logger = logging.getLogger(__name__)
# Parameters bind untyped through the ``sa.text(...)`` shim; asyncpg infers
# each placeholder's type from its target column. For ``occurred_at``
# (TIMESTAMPTZ) it wants a real ``datetime`` (a string is rejected even with a
# CAST), so we bind the aware datetime object on Postgres; SQLite's sqlite3
# driver can't bind a ``datetime`` on Python 3.12+, so we bind an ISO string
# there. ``metadata`` (JSONB) takes a JSON string on both — asyncpg's jsonb
# codec accepts ``str`` directly, so no cast is needed. Same SQL both ways;
# only the ``occurred_at`` bind value differs by dialect.
_INSERT_SQL = (
"INSERT INTO usage_events (event_id, occurred_at, metric, value, metadata) "
"VALUES (?, ?, ?, ?, ?) "
"ON CONFLICT (event_id) DO NOTHING"
)
class UsageEventStore:
"""Append-only writer for the app-DB ``usage_events`` table."""
def __init__(self, storage: RefreshTokenStorage) -> None:
self._storage = storage
@classmethod
async def shared(cls) -> "UsageEventStore":
"""Build a store backed by the process-wide storage singleton.
``get_shared_storage()`` runs ``initialize()`` (and thus Alembic
migrations) on first access, so the ``usage_events`` table is present
by the time any event is recorded.
"""
return cls(await get_shared_storage())
async def record_usage_event(
self,
*,
metric: str,
value: int,
occurred_at: datetime | None = None,
metadata: dict[str, Any] | None = None,
event_id: str | None = None,
) -> None:
"""Record one billable usage event (best-effort, flag-gated).
Does nothing unless ``USAGE_METERING_ENABLED`` is true. Any failure is
logged and swallowed — this must never break the caller's operation.
Args:
metric: Catalog metric, e.g. ``"embeddings_queries"`` or
``"pages_chunks"``.
value: Count/quantity for this event.
occurred_at: Operation completion time; defaults to now (UTC).
metadata: Optional rawest-unit context (provider, model, tokens,
doc_type, ...). Stored as JSONB (Postgres) / JSON text (SQLite).
event_id: Optional idempotency key; defaults to a fresh UUID4.
"""
if not get_settings().usage_metering_enabled:
return
start = time.time()
try:
event_id = event_id or str(uuid.uuid4())
when = occurred_at or datetime.now(timezone.utc)
# asyncpg takes the datetime object directly; sqlite3 needs a string.
when_bind = (
when if self._storage.dialect == "postgresql" else when.isoformat()
)
# json.dumps lives inside the best-effort try: a non-serializable
# metadata dict must be swallowed like any other write failure, not
# raised into the caller's operation (see the contract above).
params = (
event_id,
when_bind,
metric,
value,
json.dumps(metadata, sort_keys=True) if metadata is not None else None,
)
async with self._storage.acquire() as db:
await db.execute(_INSERT_SQL, params)
await db.commit()
record_db_operation(
self._storage.dialect, "insert", time.time() - start, "success"
)
except Exception:
# Best-effort: never surface a metering failure to the user op.
record_db_operation(
self._storage.dialect, "insert", time.time() - start, "error"
)
logger.warning(
"usage metering write dropped (metric=%s, value=%s)",
metric,
value,
exc_info=True,
)
+23
View File
@@ -29,6 +29,7 @@ from nextcloud_mcp_server.observability.metrics import (
) )
from nextcloud_mcp_server.observability.tracing import trace_operation from nextcloud_mcp_server.observability.tracing import trace_operation
from nextcloud_mcp_server.search.pdf_highlighter import PDFHighlighter from nextcloud_mcp_server.search.pdf_highlighter import PDFHighlighter
from nextcloud_mcp_server.usage import UsageEventStore
from nextcloud_mcp_server.vector import payload_keys from nextcloud_mcp_server.vector import payload_keys
from nextcloud_mcp_server.vector.document_chunker import ( from nextcloud_mcp_server.vector.document_chunker import (
DocumentChunker, DocumentChunker,
@@ -821,6 +822,28 @@ async def _index_document(
chunks=len(chunk_texts), chunks=len(chunk_texts),
chars=total_chars, chars=total_chars,
) )
# Usage metering (Deck #67): record chunks embedded as a billable
# 'pages_chunks' event. Best-effort and gated on the flag so the
# off-path (OSS default) touches no storage; placed after the
# embedding succeeds so it can never affect the indexing path.
if settings.usage_metering_enabled:
try:
store = await UsageEventStore.shared()
await store.record_usage_event(
metric="pages_chunks",
value=len(chunk_texts),
metadata={
"provider": provider,
"model": settings.get_embedding_model_name(),
"doc_type": doc_task.doc_type,
"user_id": doc_task.user_id,
"total_chars": total_chars,
},
)
except Exception:
logger.debug(
"usage metering hook (pages_chunks) skipped", exc_info=True
)
async def generate_sparse_embeddings(): async def generate_sparse_embeddings():
"""Generate sparse embeddings (BM25 for keyword matching).""" """Generate sparse embeddings (BM25 for keyword matching)."""
+12
View File
@@ -157,6 +157,18 @@ class TestGetSettings:
assert settings.vector_sync_processor_workers == 5 assert settings.vector_sync_processor_workers == 5
assert settings.vector_sync_queue_max_size == 5000 assert settings.vector_sync_queue_max_size == 5000
@patch.dict(os.environ, {}, clear=True)
def test_usage_metering_disabled_by_default(self):
"""USAGE_METERING_ENABLED defaults to False (OSS doesn't self-monitor)."""
_reload_config()
assert get_settings().usage_metering_enabled is False
@patch.dict(os.environ, {"USAGE_METERING_ENABLED": "true"}, clear=True)
def test_usage_metering_enabled_via_env(self):
"""USAGE_METERING_ENABLED=true maps to settings.usage_metering_enabled."""
_reload_config()
assert get_settings().usage_metering_enabled is True
class TestChunkConfigValidation: class TestChunkConfigValidation:
"""Test document chunking configuration validation.""" """Test document chunking configuration validation."""
+191
View File
@@ -0,0 +1,191 @@
"""Unit tests for ``UsageEventStore`` (Deck #67 usage metering, data plane).
Parametrized over both supported backends via the shared ``storage_backend``
fixture: SQLite (default, always runs) and Postgres (opt-in, gated on
``TEST_DATABASE_URL`` — bring up ``docker compose --profile postgres up -d
postgres-test`` and export
``TEST_DATABASE_URL=postgresql+asyncpg://mcp:mcp@localhost:5433/mcp``).
Covers the recording contract: flag-gated no-op, insert roundtrip, ON CONFLICT
dedup, JSON metadata roundtrip, NULL metadata, and the best-effort guarantee
that a DB failure is swallowed instead of surfacing to the caller.
"""
import json
import tempfile
import uuid
from pathlib import Path
import pytest
from cryptography.fernet import Fernet
import nextcloud_mcp_server.usage.store as store_module
from nextcloud_mcp_server.auth.storage import RefreshTokenStorage
from nextcloud_mcp_server.usage.store import UsageEventStore
pytestmark = pytest.mark.unit
@pytest.fixture
async def storage(storage_backend):
"""Initialized RefreshTokenStorage backed by SQLite or Postgres."""
key = Fernet.generate_key()
if storage_backend["kind"] == "sqlite":
with tempfile.TemporaryDirectory() as tmpdir:
db_path = Path(tmpdir) / "usage.db"
s = RefreshTokenStorage(db_path=str(db_path), encryption_key=key)
await s.initialize()
yield s
else:
s = RefreshTokenStorage(database_url=storage_backend["url"], encryption_key=key)
await s.initialize()
try:
yield s
finally:
await storage_backend["reset"]()
def _set_metering(monkeypatch, enabled: bool) -> None:
"""Force the metering flag without mutating global dynaconf state.
``record_usage_event`` calls ``get_settings()`` (imported into the store
module's namespace), and ``get_settings()`` builds a fresh Settings per
call — so patching the symbol the store sees is the clean seam.
"""
class _Settings:
usage_metering_enabled = enabled
monkeypatch.setattr(store_module, "get_settings", lambda: _Settings())
async def _count(storage: RefreshTokenStorage) -> int:
async with storage.acquire() as db:
cursor = await db.execute("SELECT COUNT(*) FROM usage_events")
row = await cursor.fetchone()
return row[0]
async def _fetch(storage: RefreshTokenStorage, event_id: str):
async with storage.acquire() as db:
cursor = await db.execute(
"SELECT event_id, occurred_at, metric, value, metadata "
"FROM usage_events WHERE event_id = ?",
(event_id,),
)
return await cursor.fetchone()
async def test_flag_off_is_noop(storage, monkeypatch):
"""With metering disabled, nothing is written (zero DB work)."""
_set_metering(monkeypatch, False)
store = UsageEventStore(storage)
await store.record_usage_event(metric="pages_chunks", value=5)
assert await _count(storage) == 0
async def test_insert_roundtrip(storage, monkeypatch):
"""A recorded event lands and reads back with the right fields."""
_set_metering(monkeypatch, True)
store = UsageEventStore(storage)
eid = str(uuid.uuid4())
await store.record_usage_event(
metric="pages_chunks",
value=7,
event_id=eid,
metadata={"provider": "gateway"},
)
row = await _fetch(storage, eid)
assert row is not None
# Postgres returns event_id as a uuid.UUID; normalize to str for compare.
assert str(row[0]) == eid
assert row[2] == "pages_chunks"
assert row[3] == 7
async def test_on_conflict_dedup(storage, monkeypatch):
"""A duplicate event_id is a no-op; the first write is retained."""
_set_metering(monkeypatch, True)
store = UsageEventStore(storage)
eid = str(uuid.uuid4())
await store.record_usage_event(metric="pages_chunks", value=1, event_id=eid)
await store.record_usage_event(metric="embeddings_queries", value=99, event_id=eid)
assert await _count(storage) == 1
row = await _fetch(storage, eid)
assert row[2] == "pages_chunks" # DO NOTHING, not DO UPDATE
assert row[3] == 1
async def test_metadata_json_roundtrip(storage, monkeypatch):
"""Nested metadata round-trips as JSON on both backends."""
_set_metering(monkeypatch, True)
store = UsageEventStore(storage)
eid = str(uuid.uuid4())
meta = {"provider": "gateway", "model": "titan", "nested": {"chunks": 3}}
await store.record_usage_event(
metric="pages_chunks", value=3, event_id=eid, metadata=meta
)
row = await _fetch(storage, eid)
raw = row[4]
# asyncpg returns JSONB as a JSON string (no codec); SQLite stores TEXT.
loaded = raw if isinstance(raw, dict) else json.loads(raw)
assert loaded == meta
async def test_metadata_none_is_null(storage, monkeypatch):
"""Omitting metadata stores SQL NULL, not the string 'null'."""
_set_metering(monkeypatch, True)
store = UsageEventStore(storage)
eid = str(uuid.uuid4())
await store.record_usage_event(
metric="embeddings_queries", value=1, event_id=eid, metadata=None
)
row = await _fetch(storage, eid)
assert row[4] is None
async def test_best_effort_swallows_db_errors(storage, monkeypatch):
"""A DB failure is logged + dropped, never raised into the caller."""
_set_metering(monkeypatch, True)
store = UsageEventStore(storage)
recorded: list[tuple] = []
monkeypatch.setattr(
store_module,
"record_db_operation",
lambda *args, **kwargs: recorded.append(args),
)
def _boom():
raise RuntimeError("db down")
# ``acquire`` raises on call — the store must catch and continue.
monkeypatch.setattr(storage, "acquire", _boom)
# Must not raise.
await store.record_usage_event(metric="pages_chunks", value=1)
assert recorded, "record_db_operation should be called on the error path"
assert recorded[-1][3] == "error"
async def test_best_effort_swallows_unserializable_metadata(storage, monkeypatch):
"""Non-serializable metadata is swallowed, not raised into the caller.
json.dumps runs inside the best-effort try, so a metadata value the JSON
encoder can't handle must drop the event like any other write failure
rather than surfacing to the user op.
"""
_set_metering(monkeypatch, True)
store = UsageEventStore(storage)
# An arbitrary object is not JSON-serializable; json.dumps raises TypeError.
bad_metadata = {"obj": object()}
# Must not raise.
await store.record_usage_event(
metric="pages_chunks", value=1, metadata=bad_metadata
)
# Nothing was written — the encode failed before the insert.
assert await _count(storage) == 0