Metadata-Version: 2.5
Name: policyengine-observability
Version: 3.0.2
Summary: Shared PolicyEngine observability runtime for logs, timings, metrics, and OpenTelemetry.
Author-email: PolicyEngine <hello@policyengine.org>
License-File: LICENSE
Requires-Python: >=3.11
Provides-Extra: all
Requires-Dist: fastapi; extra == 'all'
Requires-Dist: flask>=2.2; extra == 'all'
Requires-Dist: google-auth>=2.38.0; extra == 'all'
Requires-Dist: google-cloud-logging>=3.15.0; extra == 'all'
Requires-Dist: httpx; extra == 'all'
Requires-Dist: opentelemetry-api>=1.43.0; extra == 'all'
Requires-Dist: opentelemetry-exporter-otlp-proto-grpc>=1.43.0; extra == 'all'
Requires-Dist: opentelemetry-exporter-otlp-proto-http>=1.43.0; extra == 'all'
Requires-Dist: opentelemetry-sdk>=1.43.0; extra == 'all'
Provides-Extra: dev
Requires-Dist: build; extra == 'dev'
Requires-Dist: coverage; extra == 'dev'
Requires-Dist: pyright>=1.1.405; extra == 'dev'
Requires-Dist: pytest; extra == 'dev'
Requires-Dist: ruff>=0.9.0; extra == 'dev'
Requires-Dist: towncrier>=24.8.0; extra == 'dev'
Provides-Extra: fastapi
Requires-Dist: fastapi; extra == 'fastapi'
Provides-Extra: flask
Requires-Dist: flask>=2.2; extra == 'flask'
Provides-Extra: google
Requires-Dist: google-auth>=2.38.0; extra == 'google'
Requires-Dist: google-cloud-logging>=3.15.0; extra == 'google'
Provides-Extra: httpx
Requires-Dist: httpx; extra == 'httpx'
Provides-Extra: otel
Requires-Dist: opentelemetry-api>=1.43.0; extra == 'otel'
Requires-Dist: opentelemetry-sdk>=1.43.0; extra == 'otel'
Provides-Extra: otlp-grpc
Requires-Dist: opentelemetry-api>=1.43.0; extra == 'otlp-grpc'
Requires-Dist: opentelemetry-exporter-otlp-proto-grpc>=1.43.0; extra == 'otlp-grpc'
Requires-Dist: opentelemetry-sdk>=1.43.0; extra == 'otlp-grpc'
Provides-Extra: otlp-http
Requires-Dist: opentelemetry-api>=1.43.0; extra == 'otlp-http'
Requires-Dist: opentelemetry-exporter-otlp-proto-http>=1.43.0; extra == 'otlp-http'
Requires-Dist: opentelemetry-sdk>=1.43.0; extra == 'otlp-http'
Description-Content-Type: text/markdown

# policyengine-observability

`policyengine-observability` provides a shared runtime for structured logs,
OpenTelemetry traces and metrics, request propagation, and framework request
instrumentation. Each application owns its service identity, attribute policy,
logging destinations, OTLP endpoints, and credentials.

Version 2 uses an explicitly owned runtime. `configure(config)` returns that
runtime, and every adapter or manual operation receives it. Configuration does
not select a destination from the deployment platform or contact a remote
service.

## Install

Install only the integrations used by the service:

```bash
pip install "policyengine-observability[otel,otlp-grpc]"
pip install "policyengine-observability[flask,httpx,google]"
pip install "policyengine-observability[fastapi,httpx,google]"
```

The base package has no required dependencies. OpenTelemetry, OTLP exporters,
Google authentication and logging, and web frameworks are optional extras and
are imported only when configured or called.

## Configure a runtime

Service and deployment identity are explicit. Logging and OTel use separate
configuration so they can write to different stores.

```python
from policyengine_observability import (
    DeploymentIdentity,
    LoggingConfig,
    ObservabilityConfig,
    OTelConfig,
    OTLPExporterConfig,
    ServiceIdentity,
    StdoutLogDestination,
    configure,
)

config = ObservabilityConfig(
    service=ServiceIdentity(
        name="example-api",
        namespace="policyengine.example",
        version="1.2.3",
        role="api",
    ),
    deployment=DeploymentIdentity(
        environment="production",
        platform="google_cloud_run",
        region="us-central1",
    ),
    logging=LoggingConfig(
        destinations=(StdoutLogDestination(),),
    ),
    otel=OTelConfig(
        traces=OTLPExporterConfig(endpoint="collector:4317"),
        metrics=OTLPExporterConfig(endpoint="collector:4317"),
    ),
    dispatch_attribute_keys=frozenset({"job_id", "run_id"}),
)
runtime = configure(config)
```

Local logs and spans accept explicitly supplied safe scalar attributes by
default. Set `application_attribute_keys` to a `frozenset` only when a consumer
needs a strict local attribute allowlist. Asynchronous context still transports
only `dispatch_attribute_keys`, and metrics still use their separate
low-cardinality allowlist.

`configure` validates the complete configuration before it creates workers,
exporters, or logging handlers. Invalid values raise `ConfigurationError` with
the fields that must be corrected. Unavailable credentials or destinations
after successful validation remain nonfatal runtime failures.

The default logging destination is one-line JSON on standard output. An OTel
runtime without an exporter still creates local trace context for log
correlation. It does not send traces or metrics remotely.

## Logging destinations

Application code emits a provider-neutral record. `LoggingConfig` selects one
or more destination strategies when the runtime starts.

```python
from policyengine_observability import (
    GoogleCloudLogDestination,
    GoogleCloudLogFormatter,
)

logging = LoggingConfig(
    destinations=(
        StdoutLogDestination(
            formatter=GoogleCloudLogFormatter("trace-project"),
        ),
        GoogleCloudLogDestination(
            project_id="logging-project",
            log_name="example-api",
            queue_capacity=1_000,
            batch_size=100,
            write_timeout_seconds=5,
        ),
    ),
)
```

The formatter adds Cloud Logging trace-correlation fields only to the output
it formats. The canonical record retains the portable `trace_id`, `span_id`,
and `trace_sampled` fields. Each destination and formatter receives a deep
copy of the canonical record, so mutations to nested values remain local to
that destination.

Configured sensitive values are replaced in messages, exception details, and
allowlisted string attributes before length limits are applied and before the
record reaches a logging destination, span, or metric exporter. This also
applies to exception events on spans and local internal diagnostics.

Every queued destination has its own bounded queue and worker. A blocked or
failing destination cannot delay another destination or the application
operation that emitted the record. Queue saturation drops the newest record
and records a local diagnostic.

### Custom destinations

Use `CustomLogDestination` for a writer owned by an application or another
package:

```python
from policyengine_observability import CustomLogDestination

logging = LoggingConfig(
    destinations=(
        CustomLogDestination(
            name="internal-log-store",
            writer_factory=lambda: InternalLogWriter(),
            delivery="queued",
        ),
    ),
)
```

The writer implements `write(record)`. It may optionally implement
`write_many(records)` and `close()`. A separate integration package can also
provide a class implementing `LogDestinationStrategy`. Network writers should
always use `delivery="queued"`. Cleanup for inline writers runs on daemon
threads, and orderly shutdown waits for it only within the configured logging
shutdown timeout.

## OTLP destinations and authentication

Traces and metrics have independent exporter configurations:

```python
otel = OTelConfig(
    traces=OTLPExporterConfig(
        endpoint="https://trace-collector.example",
        protocol="http/protobuf",
        headers=(("x-api-key", "trace-key"),),
    ),
    metrics=OTLPExporterConfig(
        endpoint="metrics-collector.example:4317",
        protocol="grpc",
        headers=(("x-api-key", "metric-key"),),
    ),
)
```

For OTLP over HTTP, the default `endpoint_mode="base"` appends
`/v1/traces` or `/v1/metrics` to the configured endpoint. Set
`endpoint_mode="signal"` when the endpoint already identifies the exact
signal route.

For a Google ID-token protected collector, select the authentication strategy
explicitly:

```python
from policyengine_observability import GoogleIdTokenAuth

traces = OTLPExporterConfig(
    endpoint="collector.example:443",
    auth=GoogleIdTokenAuth("https://collector.example"),
)
```

`ObservabilityConfig.from_env` reads the standard common and signal-specific
OTel settings, including:

```bash
OTEL_EXPORTER_OTLP_ENDPOINT=collector.example:4317
OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=trace-collector.example:4317
OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=metric-collector.example:4317
OTEL_EXPORTER_OTLP_PROTOCOL=grpc
OTEL_EXPORTER_OTLP_TRACES_HEADERS=x-api-key=trace-key
OTEL_EXPORTER_OTLP_METRICS_HEADERS=x-api-key=metric-key
```

The common HTTP endpoint is treated as a base URL. Signal-specific HTTP
endpoint variables are passed to the exporter exactly as configured, following
the OpenTelemetry environment-variable contract.

`POLICYENGINE_OTEL_GOOGLE_AUDIENCE` selects `GoogleIdTokenAuth` for both
signals. `POLICYENGINE_OTEL_TRACES_GOOGLE_AUDIENCE` and
`POLICYENGINE_OTEL_METRICS_GOOGLE_AUDIENCE` override it per signal.

Google authentication constructed from `GCP_CREDENTIALS_JSON`,
`MODAL_IDENTITY_TOKEN`, or `OBSERVABILITY_GOOGLE_OIDC_TOKEN` stays in process
memory. The package passes credential data and subject tokens directly to
Google Auth and does not create temporary credential files. A path explicitly
provided through `GOOGLE_APPLICATION_CREDENTIALS` remains supported.

Applications can share a collector by configuring the same endpoint and
credentials. Another application can use a separate collector or any
OTLP-compatible service without changing this package.

## Flask and FastAPI

Install a framework adapter once while constructing the application:

```python
from flask import Flask
from policyengine_observability import instrument_flask

app = Flask(__name__)
instrument_flask(app, runtime)
```

```python
from fastapi import FastAPI
from policyengine_observability import instrument_fastapi

app = FastAPI()
instrument_fastapi(app, runtime)
```

Both adapters extract W3C trace context and
`X-PolicyEngine-Request-Id`, create one server span, emit one completion
record, record request metrics, and clear context-local state. Repeated calls
reuse the first runtime associated with the application. Install the Flask
adapter before serving the first request. If Flask rejects callback
registration, the adapter restores the prior callback registries, leaves the
application unmarked, reports a local diagnostic, and returns without raising.

## Application instrumentation

The package instruments framework and transport mechanics. Applications own
the names and boundaries of their domain operations.

```python
with runtime.operation(
    "simulation.run",
    attributes={"backend": "modal"},
):
    with runtime.span("simulation.build"):
        simulation = build_simulation()
    result = calculate(simulation)
```

Operations and spans support synchronous and asynchronous context management
and decoration:

```python
@runtime.span("simulation.calculate")
async def calculate(simulation):
    return await simulation.calculate()
```

Function arguments and return values are never captured. The runtime accepts
only application-configured attribute keys and scalar values.

```python
runtime.set_context(auth_result="accepted", simulation_id="sim-123")
runtime.event("simulation.dispatched", attributes={"backend": "modal"})

try:
    load_result()
except ValueError as error:
    runtime.record_exception(error, handled=True)
```

## Python logging and HTTP propagation

Instrument a specific Python logger or configure root capture explicitly:

```python
import logging
from policyengine_observability import instrument_logging

logger = logging.getLogger("policyengine.example")
instrument_logging(logger, runtime)
```

Only the supplied HTTPX client is modified:

```python
import httpx
from policyengine_observability import instrument_httpx

client = httpx.AsyncClient()
instrument_httpx(client, runtime)
```

The request hook injects active W3C context and the PolicyEngine request ID.
Other clients in the process remain unchanged.

For asynchronous dispatch, send the bounded observability context as transport
metadata beside the application payload and restore it around the worker
operation:

```python
observability_context = runtime.capture_context()
worker.spawn(
    payload,
    observability_context=observability_context,
)

def worker(payload, *, observability_context=None):
    with runtime.operation(
        "simulation.run",
        remote_context=observability_context,
    ):
        return run_simulation(payload)
```

`capture_context()` includes W3C trace context, its capture time, the active
PolicyEngine request ID, and scalar attributes named by
`dispatch_attribute_keys`. Starting the remote operation restores only those
configured dispatch attributes. They remain available to nested
`capture_context()` calls and are attached to logs, nested operations, and
nested spans inside the operation. They are never added to metric labels
unless separately included in `metric_attribute_keys`.

Keep this context separate from the application payload. Invalid or stale
trace context can reduce correlation, but it does not prevent the observed
application code from running. A recent direct dispatch continues the trace;
delayed, retry, and aggregate work starts a trace linked to the dispatch span.

## Process identity

OpenTelemetry resource attributes require `service.instance.id` to identify
one telemetry-producing process. Deployment revisions identify code shared by
multiple containers and workers, so they must not be used alone as the process
identity. Construct the deployment identity after the application process has
started:

```python
import os

from policyengine_observability import DeploymentIdentity, process_instance_id

service_name = "example-api"
deployment = DeploymentIdentity(
    environment="production",
    platform="google_cloud_run",
    region="us-central1",
    instance_id=process_instance_id(service_name, os.getenv("K_REVISION")),
)
```

The helper returns one stable value for a service within the current process.
It combines the optional platform identifier with a process ID and random UUID,
so separate workers and containers cannot publish cumulative metrics under the
same resource identity. A child process receives a new value on its first call.

## Process restoration and shutdown

After a process image or memory snapshot is restored, rebuild process-local
locks, context, queues, threads, credentials, and exporters before accepting
work:

```python
@modal.enter(snap=False)
def restore_process_state(self):
    self.runtime.restart_after_snapshot()
```

Call `runtime.shutdown()` during orderly process shutdown. Logging and OTel
each use their own timeout. After configuration validation succeeds, missing
optional dependencies, credential failures, unavailable destinations, queue
saturation, exporter errors, and shutdown timeouts produce rate-limited
diagnostics on standard error. These runtime failures do not change application
responses, return values, or exceptions.

## Release workflow

Changes include a Towncrier fragment in `changelog.d/`. Pull requests run Ruff,
tests, type checks, and coverage checks. The release workflow builds
distributions and publishes through PyPI trusted publishing.

## License

Code in this repository is released under the [MIT License](LICENSE). Original
text and figures are released under
[CC BY 4.0](https://creativecommons.org/licenses/by/4.0/) with attribution to
PolicyEngine. Third-party materials retain their terms.
