From c29ca9a4611f19da759ffeddd8a3a031ffbc0cf9 Mon Sep 17 00:00:00 2001 From: Aakash Date: Tue, 22 Sep 2026 05:19:39 +0530 Subject: [PATCH] feat: add optional OpenTelemetry tracing --- .github/workflows/ci.yml | 2 +- README.md | 10 ++ docs/configuration-and-api.md | 5 + docs/getting-started.md | 1 + docs/guarantees-and-limitations.md | 9 + docs/observability.md | 208 +++++++++++++++++++++++ docs/semantic-and-async.md | 6 + mkdocs.yml | 1 + pyproject.toml | 4 +- src/trimwise/telemetry.py | 111 +++++++++++++ src/trimwise/trimmer.py | 187 +++++++++++++-------- tests/test_telemetry.py | 258 +++++++++++++++++++++++++++++ 12 files changed, 732 insertions(+), 70 deletions(-) create mode 100644 docs/observability.md create mode 100644 src/trimwise/telemetry.py create mode 100644 tests/test_telemetry.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 75bb8ba..a54d6be 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -69,7 +69,7 @@ jobs: wheel-test/bin/python -c "import asyncio, importlib.metadata; from trimwise import ContextSource, ContextSourceResult, ContextTrimResult, Trimmer; - assert importlib.metadata.version('trimwise') == '0.6.0'; + assert importlib.metadata.version('trimwise') == '0.7.0'; sources = [ContextSource('a', '', '')]; sync_result = Trimmer().trim_context(sources, 8, unit='characters'); async_result = asyncio.run(Trimmer().atrim_context(sources, 8, unit='characters')); diff --git a/README.md b/README.md index 3213436..967c5fd 100644 --- a/README.md +++ b/README.md @@ -123,6 +123,16 @@ Unlike token-pruning compressors, it keeps readable source pieces; see the [research comparison](https://trimwise.readthedocs.io/en/latest/research-foundations/#how-trimwise-compares-with-model-based-compression) for the tradeoff. +## Observe trimming in your application + +Trimwise creates OpenTelemetry spans for its public operations when your application configures +an OpenTelemetry SDK. Without an SDK provider, the API stays a no-op and sends nothing. Your +application keeps control of exporters, collector endpoints, credentials, resources, sampling, +batching, and shutdown; none of those settings are added to `TrimConfig`. + +See [Observability](https://trimwise.readthedocs.io/en/latest/observability/) for OTLP/gRPC and +OTLP/HTTP setup, span names and attributes, trace parenting, and privacy guarantees. + A tight limit can leave out evidence needed to answer a question. The limit applies to Trimwise's returned text, not your entire prompt, so leave room for instructions and the model's answer. Read the [guarantees and limitations](https://trimwise.readthedocs.io/en/latest/guarantees-and-limitations/) diff --git a/docs/configuration-and-api.md b/docs/configuration-and-api.md index fdcb5fc..4a6ebd0 100644 --- a/docs/configuration-and-api.md +++ b/docs/configuration-and-api.md @@ -428,6 +428,10 @@ Most applications should keep the defaults. Configuration changes affect every t through that `Trimmer`; per-call choices such as strategy, query, limit, unit, and custom token counter remain method arguments. +OpenTelemetry exporters, endpoints, credentials, resources, and sampling are application runtime +settings, so they are intentionally absent from `TrimConfig`. See [Observability](observability.md) +for optional tracing setup and the emitted span contract. + ### Configuration validation | Setting | Accepted boundary | @@ -667,3 +671,4 @@ style, MMR balance, and the managed semantic backend. - Follow the [Getting Started guide](getting-started.md). - Compare ranking behavior in [Choosing a Strategy](strategies.md). - Configure embeddings with [Semantic Models and Async Usage](semantic-and-async.md). +- Add optional tracing with [Observability](observability.md). diff --git a/docs/getting-started.md b/docs/getting-started.md index de83294..8a0823e 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -436,4 +436,5 @@ not factual truth. - Trim several sources with [one shared limit](multi-source-context.md). - Review the current [strategy guide](https://github.com/tenwritehq/trimwise#which-strategy-should-i-use). - Learn about [embedding callbacks and FastEmbed](https://github.com/tenwritehq/trimwise#semantic-models). +- Add Trimwise to application traces with [OpenTelemetry](observability.md). - See the [public package on PyPI](https://pypi.org/project/trimwise/). diff --git a/docs/guarantees-and-limitations.md b/docs/guarantees-and-limitations.md index 99ce20f..e9f22e4 100644 --- a/docs/guarantees-and-limitations.md +++ b/docs/guarantees-and-limitations.md @@ -191,6 +191,14 @@ callback is awaited on the calling event loop. The underlying synchronous counter, callback, or FastEmbed inference may continue after the awaiting task is cancelled. +### Built-in tracing excludes source content + +Optional OpenTelemetry spans contain bounded operation metadata such as strategy, unit, counts, +and source or batch size. Trimwise does not attach source text, output text, queries, context +wrappers, exception messages, or stack traces. The application still owns any attributes added to +parent spans or spans created inside callbacks. See [Observability](observability.md) for the exact +attribute contract. + ## Best-effort goals ### Structural mode aims for document-wide coverage @@ -439,4 +447,5 @@ starting point; representative downstream evaluation decides whether they work f - Inspect the pipeline in [How Trimwise Works](how-it-works.md). - Configure backends in [Semantic Models and Async Usage](semantic-and-async.md). - Review the public contract in [Configuration and API Reference](configuration-and-api.md). +- Add optional tracing with [Observability](observability.md). - Track deferred work in the [roadmap](https://github.com/tenwritehq/trimwise/blob/main/ROADMAP.md). diff --git a/docs/observability.md b/docs/observability.md new file mode 100644 index 0000000..10673c0 --- /dev/null +++ b/docs/observability.md @@ -0,0 +1,208 @@ +--- +title: Observe Trimwise with OpenTelemetry +description: Add Trimwise spans to your application's existing traces without giving the library control of exporters, endpoints, or credentials. +--- + +# Observe Trimwise with OpenTelemetry + +Trimwise can add its work to your application's distributed traces. Tracing is optional: if your +application does not configure an OpenTelemetry SDK, the installed OpenTelemetry API is a no-op +and Trimwise sends nothing over the network. + +Trimwise creates spans only. Your application chooses whether to collect them and owns the SDK, +exporter, collector endpoint, authentication headers, credentials, TLS, resources, sampling, +batching, flushing, and shutdown. These operational settings do not belong in `TrimConfig`, which +remains limited to trimming behavior. This follows the +[OpenTelemetry guidance for instrumented libraries](https://opentelemetry.io/docs/specs/otel/library-guidelines/): +libraries use the API, while the final application configures the SDK and exporters. + +## What you get + +Each public operation creates one span: + +| Call | Span name | +| --- | --- | +| `trim()` and `atrim()` | `trimwise.trim` | +| `trim_context()` and `atrim_context()` | `trimwise.trim_context` | +| `atrim_many()` | `trimwise.trim_many` | + +Sync and async forms use the same names so a dashboard does not need separate queries. A batch +creates one span for the whole call rather than one span per input. + +When your application already has a current span, the Trimwise span becomes its child. This works +through the worker threads used by `atrim()` and `atrim_context()`, and through an awaited async +embedding callback. + +Successful single-source and context spans use these attributes: + +| Attribute | Meaning | +| --- | --- | +| `trimwise.strategy.requested` | Valid strategy passed by the caller, including `auto` | +| `trimwise.strategy.resolved` | Strategy actually used after resolving `auto` | +| `trimwise.unit` | `tokens`, `words`, or `characters` | +| `trimwise.limit` | Requested output ceiling | +| `trimwise.input.count` | Measured input size in `trimwise.unit` | +| `trimwise.output.count` | Measured output size in `trimwise.unit` | +| `trimwise.trimmed` | Whether any source text was removed | +| `trimwise.source.count` | Number of input sources for a context call | + +A successful `trimwise.trim_many` span records `trimwise.batch.size` and whether any item was +trimmed. Batch inputs may use different units, so Trimwise does not add misleading aggregate input +or output counts. + +Failed and cancelled calls set the span status to error and add `error.type`, such as +`builtins.ValueError`. Trimwise does not attach the exception message or stack trace. + +## Send traces with OTLP/gRPC + +Install the SDK and the gRPC exporter in your **application**: + +```bash +python -m pip install opentelemetry-sdk opentelemetry-exporter-otlp-proto-grpc +export OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 +``` + +Configure OpenTelemetry once when your application starts: + +```python +from opentelemetry import trace +from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import BatchSpanProcessor + + +provider = TracerProvider( + resource=Resource.create( + { + "service.name": "answer-api", + "deployment.environment.name": "production", + } + ) +) +provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter())) +trace.set_tracer_provider(provider) +``` + +The exporter reads `OTEL_EXPORTER_OTLP_ENDPOINT`. Use your collector's TLS and authentication +settings in production. Call `provider.shutdown()` from your application's shutdown hook so the +batch processor can flush pending spans. + +## Send traces with OTLP/HTTP + +Use the HTTP exporter when your collector accepts OTLP over HTTP/protobuf: + +```bash +python -m pip install opentelemetry-sdk opentelemetry-exporter-otlp-proto-http +export OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4318 +``` + +The provider setup is the same except for the exporter import: + +```python +from opentelemetry import trace +from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import BatchSpanProcessor + + +provider = TracerProvider(resource=Resource.create({"service.name": "answer-api"})) +provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter())) +trace.set_tracer_provider(provider) +``` + +With the base endpoint above, the HTTP exporter sends traces to `/v1/traces`. Call +`provider.shutdown()` when the application stops. + +## See Trimwise inside an application trace + +Once the provider is configured, use Trimwise normally: + +```python +from opentelemetry import trace +from trimwise import Trimmer + + +app_tracer = trace.get_tracer("answer-api") +trimmer = Trimmer() + +with app_tracer.start_as_current_span("answer.generate") as span: + span.set_attribute("app.plan", "standard") + result = trimmer.trim( + document, + limit=500, + strategy="lexical", + query="What caused the outage?", + ) +``` + +The trace has this shape: + +```text +answer.generate +└── trimwise.trim +``` + +Put application-specific fields on your `Resource` or parent span as shown above. That keeps +deployment and tenant metadata under application control. Trimwise does not accept an unrestricted +custom-attribute dictionary. + +## Privacy and cardinality + +Trimwise records operation metadata, not content. Its spans never include: + +- Input or output text. +- Queries. +- Context prefixes, suffixes, or separators. +- Embedding passages or vectors. +- Exception messages or stack traces. +- Collector endpoints, headers, credentials, or other configuration. + +Do not put source text, queries, user IDs, request IDs, or other high-cardinality or sensitive +values on parent spans unless your own telemetry policy permits them. OpenTelemetry context flows +into embedding callbacks, so spans created by your callback can become descendants of the +Trimwise span; the callback remains responsible for its own attributes and privacy controls. + +## What Trimwise does not configure + +Trimwise depends only on `opentelemetry-api`. It does not install or configure: + +- An OpenTelemetry SDK. +- OTLP/gRPC or OTLP/HTTP exporters. +- A collector or observability backend. +- Sampling, queues, retries, or batch limits. +- Metrics or log export. +- Global propagators or resource attributes. + +Your collector can derive request counts, error rates, and duration histograms from spans if its +span-metrics connector is enabled. That is a collector choice rather than a Trimwise runtime +feature. + +## Turn tracing off + +Do not configure an OpenTelemetry SDK, or configure your application's sampler to drop these +spans. No Trimwise flag is required. The OpenTelemetry API safely returns no-op spans when no SDK +provider is installed. + +## Troubleshooting + +**No Trimwise spans appear:** verify that the SDK is configured before the first trim call, the +exporter package matches your collector protocol, and the exporter endpoint uses port `4317` for +gRPC or `4318` for HTTP by convention. + +**Spans appear under the wrong service:** set `service.name` on the application's `Resource`. +Trimwise deliberately does not choose a service name. + +**The process exits before spans arrive:** call `provider.shutdown()` during application shutdown. +Do not call it after every trim. + +**You expected metrics or logs:** Trimwise emits trace spans only. Configure application logging +and collector-derived span metrics separately. + +## Continue exploring + +- Follow the [Getting Started guide](getting-started.md). +- Review every call and result in [Configuration and API Reference](configuration-and-api.md). +- Understand worker threads and callbacks in [Semantic Models and Async Usage](semantic-and-async.md). +- Read the [guarantees and limitations](guarantees-and-limitations.md). diff --git a/docs/semantic-and-async.md b/docs/semantic-and-async.md index dd790e3..739af21 100644 --- a/docs/semantic-and-async.md +++ b/docs/semantic-and-async.md @@ -380,6 +380,11 @@ For CPU-only structural or lexical work, async calls can overlap at the worker-t FastEmbed, calls sharing one `Trimmer` still wait on that instance's model lock. `atrim_many()` does not add cross-call background batching or make parallel CPU inference requests. +When the application configures OpenTelemetry, the public operation span remains current across +these worker-thread and async-callback boundaries. Spans created inside an embedding callback can +therefore appear below the Trimwise span. See [Observability](observability.md) for setup, span +names, and the privacy contract. + ## Cancellation Cancellation behavior depends on what `atrim()`, `atrim_context()`, or `atrim_many()` is awaiting: @@ -469,3 +474,4 @@ compression ratios. - Follow the [Getting Started guide](getting-started.md). - Compare all selection modes in [Choosing a Strategy](strategies.md). - Review planned semantic quality work in the [roadmap](https://github.com/tenwritehq/trimwise/blob/main/ROADMAP.md). +- Add optional tracing with [Observability](observability.md). diff --git a/mkdocs.yml b/mkdocs.yml index c5a4b23..e8dfe9a 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -18,6 +18,7 @@ nav: - How it works: how-it-works.md - Strategies: strategies.md - Semantic models and async: semantic-and-async.md + - Observability: observability.md - Configuration and API: configuration-and-api.md - Guarantees and limitations: guarantees-and-limitations.md - Research foundations: research-foundations.md diff --git a/pyproject.toml b/pyproject.toml index 30dd666..00f7900 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "flit_core.buildapi" [project] name = "trimwise" -version = "0.6.0" +version = "0.7.0" description = "High-signal text trimming for better LLM prompts." readme = "README.md" requires-python = ">=3.10" @@ -32,6 +32,7 @@ classifiers = [ dependencies = [ "markdown-it-py>=4.2,<5", "numpy>=1.26,<3", + "opentelemetry-api>=1.44,<2", "tiktoken>=0.13,<1", ] @@ -41,6 +42,7 @@ semantic-gpu = ["fastembed-gpu>=0.8,<1"] dev = [ "build>=1.5,<2", "mypy>=2.3,<3", + "opentelemetry-sdk>=1.44,<2", "pytest>=9.1,<10", "pytest-asyncio>=1.4,<2", "pytest-cov>=7.1,<8", diff --git a/src/trimwise/telemetry.py b/src/trimwise/telemetry.py new file mode 100644 index 0000000..53251da --- /dev/null +++ b/src/trimwise/telemetry.py @@ -0,0 +1,111 @@ +"""Create privacy-safe OpenTelemetry spans around public trimming operations.""" + +from __future__ import annotations + +from collections.abc import Iterator, Mapping, Sequence +from contextlib import contextmanager +from importlib.metadata import PackageNotFoundError, version + +from opentelemetry import trace +from opentelemetry.trace import Span, Status, StatusCode + +from trimwise.models import BudgetUnit, ContextTrimResult, Strategy, TrimResult + +_VERSION: str | None +try: + _VERSION = version("trimwise") +except PackageNotFoundError: + _VERSION = None + +_TRACER = trace.get_tracer("trimwise", _VERSION) +_UNITS = {unit.value for unit in BudgetUnit} +_STRATEGIES = {strategy.value for strategy in Strategy} + + +def _request_attributes( + limit: object, + unit: object, + strategy: object, +) -> dict[str, str | int]: + """Keep only validated, bounded request values for a span. + + Args: + limit: Requested output limit. + unit: Requested measurement unit. + strategy: Requested ranking strategy. + + Returns: + Safe attributes available before the trimming call completes. + """ + attributes: dict[str, str | int] = {} + if isinstance(limit, int) and not isinstance(limit, bool) and limit >= 0: + attributes["trimwise.limit"] = limit + + unit_value = unit.value if isinstance(unit, BudgetUnit) else unit + if isinstance(unit_value, str) and unit_value in _UNITS: + attributes["trimwise.unit"] = unit_value + + strategy_value = strategy.value if isinstance(strategy, Strategy) else strategy + if isinstance(strategy_value, str) and strategy_value in _STRATEGIES: + attributes["trimwise.strategy.requested"] = strategy_value + return attributes + + +@contextmanager +def _span(name: str, attributes: Mapping[str, str | int]) -> Iterator[Span]: + """Create a span that records bounded failure types without exception content. + + Args: + name: Stable public operation name. + attributes: Validated request attributes. + + Yields: + Current OpenTelemetry span for result attributes. + + Raises: + BaseException: Re-raises any trimming failure unchanged. + """ + with _TRACER.start_as_current_span( + name, + attributes=attributes, + record_exception=False, + set_status_on_exception=False, + ) as span: + try: + yield span + except BaseException as error: + error_type = type(error) + span.set_attribute( + "error.type", + f"{error_type.__module__}.{error_type.__qualname__}", + ) + span.set_status(Status(StatusCode.ERROR)) + raise + + +def _record_result(span: Span, result: TrimResult | ContextTrimResult) -> None: + """Add common bounded result values to a completed operation span. + + Args: + span: Span that represents the public operation. + result: Successful single-source or shared-context result. + """ + span.set_attribute("trimwise.strategy.resolved", result.strategy.value) + span.set_attribute("trimwise.unit", result.unit.value) + span.set_attribute("trimwise.limit", result.limit) + span.set_attribute("trimwise.input.count", result.input_count) + span.set_attribute("trimwise.output.count", result.output_count) + span.set_attribute("trimwise.trimmed", result.trimmed) + if isinstance(result, ContextTrimResult): + span.set_attribute("trimwise.source.count", len(result.sources)) + + +def _record_batch_result(span: Span, results: Sequence[TrimResult]) -> None: + """Add bounded aggregate values to a completed batch span. + + Args: + span: Span that represents the batch operation. + results: Ordered successful batch results. + """ + span.set_attribute("trimwise.batch.size", len(results)) + span.set_attribute("trimwise.trimmed", any(result.trimmed for result in results)) diff --git a/src/trimwise/trimmer.py b/src/trimwise/trimmer.py index 1e3bba5..dca04ad 100644 --- a/src/trimwise/trimmer.py +++ b/src/trimwise/trimmer.py @@ -57,6 +57,12 @@ invoke_async_embedding_callback, normalize_callback_output, ) +from trimwise.telemetry import ( + _record_batch_result, + _record_result, + _request_attributes, + _span, +) @dataclass(frozen=True, slots=True) @@ -203,7 +209,10 @@ def trim( ValueError: If an argument value or strategy/query combination is invalid. SemanticBackendError: If an explicitly requested semantic backend fails. """ - return self._trim(_TrimArguments(text, limit, unit, strategy, query, token_counter)) + with _span("trimwise.trim", _request_attributes(limit, unit, strategy)) as span: + result = self._trim(_TrimArguments(text, limit, unit, strategy, query, token_counter)) + _record_result(span, result) + return result async def atrim( self, @@ -236,21 +245,35 @@ async def atrim( ValueError: If an argument value or strategy/query combination is invalid. SemanticBackendError: If an explicitly requested semantic backend fails. """ - arguments = _TrimArguments(text, limit, unit, strategy, query, token_counter) - callback = self._async_embedding_callback - if callback is None: - return await asyncio.to_thread(self._trim, arguments) - - prepared = await asyncio.to_thread(self._prepare, arguments) - if isinstance(prepared, TrimResult): - return prepared - if prepared.request.strategy not in {Strategy.SEMANTIC, Strategy.HYBRID}: - return await asyncio.to_thread(self._complete, prepared) - - query_text = prepared.request.query or "" - passages = await asyncio.to_thread(_contextual_ranking_texts, prepared.segments) - output = await invoke_async_embedding_callback(callback, query_text, passages) - return await asyncio.to_thread(self._complete_with_embedding_output, prepared, output) + with _span("trimwise.trim", _request_attributes(limit, unit, strategy)) as span: + arguments = _TrimArguments(text, limit, unit, strategy, query, token_counter) + callback = self._async_embedding_callback + if callback is None: + result = await asyncio.to_thread(self._trim, arguments) + else: + prepared = await asyncio.to_thread(self._prepare, arguments) + if isinstance(prepared, TrimResult): + result = prepared + elif prepared.request.strategy not in {Strategy.SEMANTIC, Strategy.HYBRID}: + result = await asyncio.to_thread(self._complete, prepared) + else: + query_text = prepared.request.query or "" + passages = await asyncio.to_thread( + _contextual_ranking_texts, + prepared.segments, + ) + output = await invoke_async_embedding_callback( + callback, + query_text, + passages, + ) + result = await asyncio.to_thread( + self._complete_with_embedding_output, + prepared, + output, + ) + _record_result(span, result) + return result def trim_context( self, @@ -285,19 +308,22 @@ def trim_context( ValueError: If an argument value or strategy/query combination is invalid. SemanticBackendError: If an explicitly requested semantic backend fails. """ - source_snapshot, rendering = _snapshot_sources(sources, separator) - _validate_deduplicate(deduplicate) - arguments = _ContextArguments( - source_snapshot, - limit, - unit, - strategy, - query, - token_counter, - deduplicate, - rendering, - ) - return self._trim_context(arguments) + with _span("trimwise.trim_context", _request_attributes(limit, unit, strategy)) as span: + source_snapshot, rendering = _snapshot_sources(sources, separator) + _validate_deduplicate(deduplicate) + arguments = _ContextArguments( + source_snapshot, + limit, + unit, + strategy, + query, + token_counter, + deduplicate, + rendering, + ) + result = self._trim_context(arguments) + _record_result(span, result) + return result async def atrim_context( self, @@ -335,45 +361,51 @@ async def atrim_context( ValueError: If an argument value or strategy/query combination is invalid. SemanticBackendError: If an explicitly requested semantic backend fails. """ - source_snapshot, rendering = _snapshot_sources(sources, separator) - _validate_deduplicate(deduplicate) - arguments = _ContextArguments( - source_snapshot, - limit, - unit, - strategy, - query, - token_counter, - deduplicate, - rendering, - ) - callback = self._async_embedding_callback - if callback is None: - return await asyncio.to_thread(self._trim_context, arguments) - - prepared = await asyncio.to_thread(self._prepare_context, arguments) - if isinstance(prepared, ContextTrimResult): - return prepared - if prepared.request.strategy not in {Strategy.SEMANTIC, Strategy.HYBRID}: - return await asyncio.to_thread(self._complete_context, prepared) - - passages = await asyncio.to_thread(_contextual_ranking_texts, prepared.segments) - batch = await asyncio.to_thread( - _prepare_passage_batch, - passages, - prepared.request.deduplicate, - ) - output = await invoke_async_embedding_callback( - callback, - prepared.request.query or "", - batch.passages, - ) - return await asyncio.to_thread( - self._complete_context_with_embedding_output, - prepared, - batch, - output, - ) + with _span("trimwise.trim_context", _request_attributes(limit, unit, strategy)) as span: + source_snapshot, rendering = _snapshot_sources(sources, separator) + _validate_deduplicate(deduplicate) + arguments = _ContextArguments( + source_snapshot, + limit, + unit, + strategy, + query, + token_counter, + deduplicate, + rendering, + ) + callback = self._async_embedding_callback + if callback is None: + result = await asyncio.to_thread(self._trim_context, arguments) + else: + prepared = await asyncio.to_thread(self._prepare_context, arguments) + if isinstance(prepared, ContextTrimResult): + result = prepared + elif prepared.request.strategy not in {Strategy.SEMANTIC, Strategy.HYBRID}: + result = await asyncio.to_thread(self._complete_context, prepared) + else: + passages = await asyncio.to_thread( + _contextual_ranking_texts, + prepared.segments, + ) + batch = await asyncio.to_thread( + _prepare_passage_batch, + passages, + prepared.request.deduplicate, + ) + output = await invoke_async_embedding_callback( + callback, + prepared.request.query or "", + batch.passages, + ) + result = await asyncio.to_thread( + self._complete_context_with_embedding_output, + prepared, + batch, + output, + ) + _record_result(span, result) + return result async def atrim_many( self, @@ -400,6 +432,25 @@ async def atrim_many( ValueError: If an input has an invalid value or strategy/query combination. SemanticBackendError: If an explicitly requested semantic backend fails. """ + with _span("trimwise.trim_many", {}) as span: + results = await self._atrim_many(inputs, deduplicate) + _record_batch_result(span, results) + return results + + async def _atrim_many( + self, + inputs: Sequence[TrimInput], + deduplicate: bool, + ) -> list[TrimResult]: + """Run the validated batch implementation inside its public operation span. + + Args: + inputs: Independent trim requests returned in the supplied order. + deduplicate: Whether exact contextual passages share one embedding. + + Returns: + One measured extractive result per input, in input order. + """ if not isinstance(deduplicate, bool): raise TypeError("deduplicate must be a bool") arguments = _batch_arguments(inputs) diff --git a/tests/test_telemetry.py b/tests/test_telemetry.py new file mode 100644 index 0000000..aeffa7d --- /dev/null +++ b/tests/test_telemetry.py @@ -0,0 +1,258 @@ +"""Verify optional OpenTelemetry tracing and its privacy contract.""" + +from __future__ import annotations + +import asyncio +from collections.abc import Iterator, Sequence +from importlib.metadata import version + +import pytest +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +from opentelemetry.trace import StatusCode + +from trimwise import ContextSource, TrimInput, Trimmer, telemetry + + +@pytest.fixture +def span_recorder( + monkeypatch: pytest.MonkeyPatch, +) -> Iterator[tuple[TracerProvider, InMemorySpanExporter]]: + """Capture Trimwise spans without changing the process-global provider. + + Args: + monkeypatch: Pytest patch helper. + + Yields: + Isolated provider and in-memory exporter. + """ + exporter = InMemorySpanExporter() + provider = TracerProvider() + provider.add_span_processor(SimpleSpanProcessor(exporter)) + monkeypatch.setattr(telemetry, "_TRACER", provider.get_tracer("trimwise", version("trimwise"))) + yield provider, exporter + provider.shutdown() + + +def test_trim_span_is_parented_and_records_safe_result_attributes( + span_recorder: tuple[TracerProvider, InMemorySpanExporter], +) -> None: + """Describe a successful trim without recording source or query text. + + Args: + span_recorder: Isolated provider and captured-span exporter. + """ + provider, exporter = span_recorder + app_tracer = provider.get_tracer("example-app") + with app_tracer.start_as_current_span("answer.request") as parent: + parent_id = parent.get_span_context().span_id + result = Trimmer().trim( + "Keep this. Remove the rest of this source.", + 10, + unit="characters", + strategy="auto", + query="private question", + ) + + spans = {span.name: span for span in exporter.get_finished_spans()} + span = spans["trimwise.trim"] + assert span.parent is not None and span.parent.span_id == parent_id + assert span.instrumentation_scope is not None + assert span.instrumentation_scope.name == "trimwise" + assert span.instrumentation_scope.version == version("trimwise") + assert span.attributes == { + "trimwise.limit": 10, + "trimwise.unit": "characters", + "trimwise.strategy.requested": "auto", + "trimwise.strategy.resolved": "lexical", + "trimwise.input.count": 42, + "trimwise.output.count": result.output_count, + "trimwise.trimmed": True, + } + assert "private question" not in str(span.attributes) + + +def test_failure_span_omits_exception_content( + span_recorder: tuple[TracerProvider, InMemorySpanExporter], +) -> None: + """Record a bounded error type without exception events or secret values. + + Args: + span_recorder: Isolated provider and captured-span exporter. + """ + _, exporter = span_recorder + secret = "customer-secret-91" + + def failing_counter(_: str) -> int: + """Raise a message that telemetry must never retain.""" + raise RuntimeError(secret) + + with pytest.raises(RuntimeError, match=secret): + Trimmer().trim( + secret, + 2, + strategy="auto", + query=secret, + token_counter=failing_counter, + ) + + (span,) = exporter.get_finished_spans() + assert span.status.status_code is StatusCode.ERROR + assert span.status.description is None + assert span.events == () + assert span.attributes == { + "trimwise.limit": 2, + "trimwise.unit": "tokens", + "trimwise.strategy.requested": "auto", + "error.type": "builtins.RuntimeError", + } + assert secret not in str(span.attributes) + + +def test_validation_failure_drops_invalid_attribute_values( + span_recorder: tuple[TracerProvider, InMemorySpanExporter], +) -> None: + """Reject unsafe enum-like values without copying them into telemetry. + + Args: + span_recorder: Isolated provider and captured-span exporter. + """ + _, exporter = span_recorder + invalid_unit = "private-unit-value" + + with pytest.raises(ValueError, match="unsupported budget unit"): + Trimmer().trim("private source", 2, unit=invalid_unit) + + (span,) = exporter.get_finished_spans() + assert span.attributes == { + "trimwise.limit": 2, + "trimwise.strategy.requested": "auto", + "error.type": "builtins.ValueError", + } + assert invalid_unit not in str(span.attributes) + + +def test_context_span_records_source_count( + span_recorder: tuple[TracerProvider, InMemorySpanExporter], +) -> None: + """Describe shared-context work with one bounded source count. + + Args: + span_recorder: Isolated provider and captured-span exporter. + """ + _, exporter = span_recorder + result = Trimmer().trim_context( + [ContextSource("alpha beta"), ContextSource("gamma delta")], + 3, + unit="words", + ) + + (span,) = exporter.get_finished_spans() + assert span.name == "trimwise.trim_context" + assert span.attributes is not None + assert span.attributes["trimwise.source.count"] == 2 + assert span.attributes["trimwise.input.count"] == result.input_count + assert span.attributes["trimwise.output.count"] == result.output_count + + +@pytest.mark.asyncio +async def test_async_callback_span_keeps_trim_parent( + span_recorder: tuple[TracerProvider, InMemorySpanExporter], +) -> None: + """Keep application callback spans below the asynchronous Trimwise span. + + Args: + span_recorder: Isolated provider and captured-span exporter. + """ + provider, exporter = span_recorder + app_tracer = provider.get_tracer("example-app") + + async def embed( + _: str, + passages: Sequence[str], + ) -> tuple[list[float], list[list[float]]]: + """Create one application-owned child span while returning valid vectors.""" + with app_tracer.start_as_current_span("embedding.request"): + return [1.0], [[1.0] for _ in passages] + + await Trimmer(async_embedding_callback=embed).atrim( + "Alpha evidence. Beta evidence. Gamma evidence.", + 12, + unit="characters", + strategy="semantic", + query="Alpha", + ) + + spans = {span.name: span for span in exporter.get_finished_spans()} + trim_span = spans["trimwise.trim"] + embed_span = spans["embedding.request"] + assert embed_span.parent is not None + assert embed_span.parent.span_id == trim_span.context.span_id + + +@pytest.mark.asyncio +async def test_batch_emits_one_operation_span_without_item_spans( + span_recorder: tuple[TracerProvider, InMemorySpanExporter], +) -> None: + """Keep batch tracing useful without creating one span per item. + + Args: + span_recorder: Isolated provider and captured-span exporter. + """ + _, exporter = span_recorder + results = await Trimmer().atrim_many( + [ + TrimInput("short", 10, unit="characters"), + TrimInput("longer source", 4, unit="characters"), + ] + ) + + (span,) = exporter.get_finished_spans() + assert span.name == "trimwise.trim_many" + assert span.attributes == { + "trimwise.batch.size": 2, + "trimwise.trimmed": any(result.trimmed for result in results), + } + + +@pytest.mark.asyncio +async def test_cancelled_async_call_closes_error_span( + span_recorder: tuple[TracerProvider, InMemorySpanExporter], +) -> None: + """End the operation span when cancellation reaches an async callback. + + Args: + span_recorder: Isolated provider and captured-span exporter. + """ + _, exporter = span_recorder + started = asyncio.Event() + + async def embed( + _: str, + passages: Sequence[str], + ) -> tuple[list[float], list[list[float]]]: + """Wait indefinitely so the test can cancel the public call.""" + started.set() + await asyncio.Event().wait() + return [1.0], [[1.0] for _ in passages] + + task = asyncio.create_task( + Trimmer(async_embedding_callback=embed).atrim( + "Alpha evidence. Beta evidence. Gamma evidence.", + 12, + unit="characters", + strategy="semantic", + query="Alpha", + ) + ) + await started.wait() + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + + (span,) = exporter.get_finished_spans() + assert span.status.status_code is StatusCode.ERROR + assert span.attributes is not None + assert span.attributes["error.type"] == "asyncio.exceptions.CancelledError" + assert span.events == ()