Files
mcp-nextcloud/tests/unit/test_processor_metering.py
T
Chris CoutinhoandClaude Opus 4.8 9676bb3106 feat(ingest): per-tier escalation via procrastinate queue-hop
Split external (procrastinate) document processing into per-tier queues so a
document is attempted at most once per tier and requeued to the next tier's
queue on a low-quality parse, using procrastinate's native retry.

- escalation.py: TIER_LADDER (fast->structured->ocr) + EscalateError signal
- registry: process_tier (one tier) + evaluate_escalation post-parse gate
  (reuses classify_from_text) + next_available_tier; shared _classify_result
  and _oversize_result with the inline pipeline
- processor: process_document(tier=...) runs one tier and raises EscalateError
  before embed (junk text never indexed); inline memory path unchanged
- queue/procrastinate: ingest-fast|structured|ocr queues; TieredEscalationStrategy
  (queue-hop on EscalateError, bounded same-tier transient retry); queue-aware
  task; producer defers to ingest-fast; per-queue counts + all-queue reclaim
- cli: worker --tier {fast,structured,ocr}
- billing: pages_ocr usage event + pipeline_tier metadata (paid OCR billed apart)
- observability: astrolabe_ingest_queue_depth{queue,status} gauge + per-queue
  counts in nc_get_vector_sync_status / management status endpoint
- config: INGEST_ESCALATION_ENABLED (default true), INGEST_TRANSIENT_MAX_ATTEMPTS

INGEST_ESCALATION_ENABLED=false and INGEST_QUEUE=memory preserve prior behaviour.

Deck #323.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-13 13:22:18 +02:00

229 lines
7.2 KiB
Python

"""Unit tests for the indexing-path usage-metering helper (Deck #67).
``record_indexing_usage`` records the billable events after a document's chunks
are embedded: ``tokens_embedded`` for every document, and ``pages_embedded``
only for parsed files (real ``page_count``). Text content (no ``page_count``)
meters tokens only — ``pages_embedded`` is a charge for parsing, not content
size (card #282). These cover the value mapping, the flag/zero-chunk no-ops, the
text-only path, and the best-effort failure path without standing up the full
document pipeline.
"""
from unittest.mock import AsyncMock, MagicMock
import pytest
from nextcloud_mcp_server.vector import processor
@pytest.fixture
def store_spy(monkeypatch):
"""Patch UsageEventStore.shared() to return a spy store."""
store = MagicMock()
store.record_usage_event = AsyncMock()
monkeypatch.setattr(
processor.UsageEventStore, "shared", AsyncMock(return_value=store)
)
return store
@pytest.mark.unit
async def test_parsed_file_records_pages_and_tokens(store_spy):
"""A parsed PDF fires both events: pages_embedded = real page count."""
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=110,
token_count=4242,
total_chars=170826,
page_count=12,
)
calls = store_spy.record_usage_event.await_args_list
by_metric = {c.kwargs["metric"]: c.kwargs["value"] for c in calls}
# pages_embedded is the real parsed-page count, NOT the chunk count.
assert by_metric == {"pages_embedded": 12, "tokens_embedded": 4242}
# Intentional ordering: tokens (recorded for every doc) before pages (the
# conditional parsing cost). Asserted so a refactor can't silently reverse
# it — a comment alone is easier to delete than a failing test.
assert calls[0].kwargs["metric"] == "tokens_embedded"
assert calls[1].kwargs["metric"] == "pages_embedded"
for c in calls:
# Hot-path fast-gate + tenant-local attribution metadata.
assert c.kwargs["enabled"] is True
assert c.kwargs["metadata"]["provider"] == "mistral"
assert c.kwargs["metadata"]["model"] == "mistral-embed"
assert c.kwargs["metadata"]["user_id"] == "alice"
assert c.kwargs["metadata"]["doc_type"] == "file"
@pytest.mark.unit
async def test_text_doc_records_tokens_only(store_spy):
"""Unparsed text content (no page_count) meters tokens, never pages."""
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="note",
user_id="alice",
chunk_count=4,
token_count=512,
total_chars=7000,
page_count=None,
)
calls = store_spy.record_usage_event.await_args_list
by_metric = {c.kwargs["metric"]: c.kwargs["value"] for c in calls}
assert by_metric == {"tokens_embedded": 512}
assert "pages_embedded" not in by_metric
@pytest.mark.unit
async def test_zero_pages_skips_pages(store_spy):
"""page_count=0 (e.g. an empty/corrupt PDF) records tokens but no pages."""
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=4,
token_count=99,
total_chars=1000,
page_count=0,
)
calls = store_spy.record_usage_event.await_args_list
by_metric = {c.kwargs["metric"]: c.kwargs["value"] for c in calls}
assert by_metric == {"tokens_embedded": 99}
@pytest.mark.unit
async def test_negative_pages_skips_pages(store_spy):
"""A malformed negative page_count meters as 'no pages' (tokens only)."""
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=4,
token_count=99,
total_chars=1000,
page_count=-1,
)
calls = store_spy.record_usage_event.await_args_list
by_metric = {c.kwargs["metric"]: c.kwargs["value"] for c in calls}
assert by_metric == {"tokens_embedded": 99}
@pytest.mark.unit
async def test_disabled_is_noop(store_spy):
"""Flag off → no store access, no events."""
await processor.record_indexing_usage(
enabled=False,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=10,
token_count=20,
total_chars=5,
page_count=3,
)
store_spy.record_usage_event.assert_not_awaited()
@pytest.mark.unit
async def test_zero_chunks_is_noop(store_spy):
"""A document with no chunks records nothing (no zero-value rows)."""
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=0,
token_count=0,
total_chars=0,
page_count=3,
)
store_spy.record_usage_event.assert_not_awaited()
@pytest.mark.unit
async def test_store_failure_is_swallowed(monkeypatch):
"""A store-construction failure is logged, never raised into indexing."""
monkeypatch.setattr(
processor.UsageEventStore,
"shared",
AsyncMock(side_effect=RuntimeError("boom")),
)
# Must not raise.
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=3,
token_count=7,
total_chars=9,
page_count=2,
)
@pytest.mark.unit
async def test_ocr_tier_records_pages_ocr(store_spy):
"""OCR-tier pages are metered as a separate pages_ocr line (Deck #323)."""
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=20,
token_count=900,
total_chars=40000,
page_count=8,
pipeline_tier="ocr",
)
by_metric = {
c.kwargs["metric"]: c.kwargs["value"]
for c in store_spy.record_usage_event.await_args_list
}
# pages_ocr fires IN ADDITION to pages_embedded for OCR-tier pages.
assert by_metric == {
"tokens_embedded": 900,
"pages_embedded": 8,
"pages_ocr": 8,
}
# pipeline_tier is threaded into the billing metadata for CP attribution.
for c in store_spy.record_usage_event.await_args_list:
assert c.kwargs["metadata"]["pipeline_tier"] == "ocr"
@pytest.mark.unit
async def test_fast_tier_does_not_record_pages_ocr(store_spy):
"""A CPU-cheap fast-tier parse must NOT incur the paid pages_ocr line."""
await processor.record_indexing_usage(
enabled=True,
provider="mistral",
model="mistral-embed",
doc_type="file",
user_id="alice",
chunk_count=10,
token_count=500,
total_chars=20000,
page_count=4,
pipeline_tier="fast",
)
metrics = {c.kwargs["metric"] for c in store_spy.record_usage_event.await_args_list}
assert "pages_ocr" not in metrics
assert metrics == {"tokens_embedded", "pages_embedded"}