Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 33 additions & 9 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,12 +60,17 @@ config = ObservabilityConfig(
traces=OTLPExporterConfig(endpoint="collector:4317"),
metrics=OTLPExporterConfig(endpoint="collector:4317"),
),
application_attribute_keys=frozenset({"country_id", "backend"}),
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
Expand Down Expand Up @@ -300,19 +305,38 @@ 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, serialize the bounded correlation context with the
job request and restore it around the worker operation:
For asynchronous dispatch, send the bounded observability context as transport
metadata beside the application payload and restore it around the worker
operation:

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

with runtime.operation(
"simulation.run",
remote_context=request.observability_context,
):
return run_simulation(request)
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 and every nested span 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 restoration and shutdown

After a process image or memory snapshot is restored, rebuild process-local
Expand Down
1 change: 1 addition & 0 deletions changelog.d/async-context.fixed.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Preserve configured dispatch attributes across asynchronous operations, include them in correlated structured logs and nested spans, accept safe scalar application attributes by default, and remove the attribute-count limit.
17 changes: 17 additions & 0 deletions docs/engineering/skills/repository-guidance.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,20 @@ uv run --extra dev towncrier check --compare-with origin/main
without breaking the application operation being observed.
- Preserve structured log schemas. Make additive changes when possible; bump
schema versions for breaking payload changes.
- Treat `capture_context()` output as transport metadata beside an application
payload. Pass it to the receiver's outer `operation` through
`remote_context`; do not insert it into business request models.
- Restore only attributes explicitly listed in `dispatch_attribute_keys`.
Those attributes must remain available to nested dispatches and structured
logs and must be attached to nested spans, but must not become metric labels
unless independently allowlisted in `metric_attribute_keys`.
- Accept explicitly supplied safe scalar attributes in local logs and spans by
default. Use `application_attribute_keys` only when a consumer requires a
strict local allowlist. Do not use that optional local policy to decide what
crosses a process boundary or becomes a metric label.
- Let explicitly supplied receiver attributes override matching remote
attributes. Malformed remote context may reduce telemetry but must not stop
the observed operation.
- Keep metric attributes bounded and low-cardinality. Do not put raw paths,
full URLs, request bodies, or unbounded user-provided values into metric
labels.
Expand All @@ -73,6 +87,9 @@ uv run --extra dev towncrier check --compare-with origin/main
Add focused tests for context behavior and failure paths whenever changing
the runtime or its components. The corresponding `tests/test_runtime_*.py`
modules cover operations, requests, spans, log emission, and tracing.
Remote-context tests must exercise two runtime instances and prove that
configured dispatch attributes survive capture, restoration, logs, spans, and
a subsequent capture. Include malformed context and local-override cases.
Adapter changes should include framework-level tests that exercise request
setup, response headers, error paths, and teardown behavior.

Expand Down
52 changes: 31 additions & 21 deletions policyengine_observability/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,6 @@ def __init__(self, errors: tuple[str, ...]) -> None:
super().__init__(f"Invalid observability configuration:\n{details}")


DEFAULT_APPLICATION_ATTRIBUTE_KEYS = frozenset(set())

DEFAULT_DISPATCH_ATTRIBUTE_KEYS = frozenset(set())

DEFAULT_METRIC_ATTRIBUTE_KEYS = frozenset(
Expand Down Expand Up @@ -106,7 +104,6 @@ class OTelConfig:

@dataclass(frozen=True, slots=True)
class TelemetryLimits:
max_attributes: int = 32
max_string_length: int = 1_024
max_error_message_length: int = 2_048
max_stack_length: int = 16_384
Expand All @@ -120,9 +117,7 @@ class ObservabilityConfig:
logging: LoggingConfig = field(default_factory=LoggingConfig)
otel: OTelConfig = field(default_factory=OTelConfig)
limits: TelemetryLimits = field(default_factory=TelemetryLimits)
application_attribute_keys: frozenset[str] = (
DEFAULT_APPLICATION_ATTRIBUTE_KEYS
)
application_attribute_keys: frozenset[str] | None = None
dispatch_attribute_keys: frozenset[str] = DEFAULT_DISPATCH_ATTRIBUTE_KEYS
metric_attribute_keys: frozenset[str] = DEFAULT_METRIC_ATTRIBUTE_KEYS
sensitive_values: tuple[str, ...] = ()
Expand Down Expand Up @@ -224,11 +219,7 @@ def from_env(
),
),
limits=limits or TelemetryLimits(),
application_attribute_keys=(
application_attribute_keys
if application_attribute_keys is not None
else DEFAULT_APPLICATION_ATTRIBUTE_KEYS
),
application_attribute_keys=application_attribute_keys,
dispatch_attribute_keys=(
dispatch_attribute_keys
if dispatch_attribute_keys is not None
Expand Down Expand Up @@ -278,19 +269,18 @@ def validation_errors(self) -> tuple[str, ...]:
"string."
)

if self.application_attribute_keys is not None:
_attribute_key_errors(
errors,
"application_attribute_keys",
self.application_attribute_keys,
)

for name, values in (
("application_attribute_keys", self.application_attribute_keys),
("dispatch_attribute_keys", self.dispatch_attribute_keys),
("metric_attribute_keys", self.metric_attribute_keys),
):
if not isinstance(values, frozenset):
errors.append(
f"{name} must be a frozenset of non-empty strings."
)
continue
for value in values:
if not isinstance(value, str) or not value.strip():
errors.append(f"{name} entries must be non-empty strings.")
_attribute_key_errors(errors, name, values)

_choice_error(
errors,
Expand Down Expand Up @@ -421,7 +411,6 @@ def validation_errors(self) -> tuple[str, ...]:
)

for name, value in {
"limits.max_attributes": self.limits.max_attributes,
"limits.max_string_length": self.limits.max_string_length,
"limits.max_error_message_length": self.limits.max_error_message_length,
"limits.max_stack_length": self.limits.max_stack_length,
Expand All @@ -436,6 +425,14 @@ def validation_errors(self) -> tuple[str, ...]:
)
return tuple(errors)

@property
def local_attribute_keys(self) -> frozenset[str] | None:
"""Return the optional strict allowlist for local logs and spans."""

if self.application_attribute_keys is None:
return None
return self.application_attribute_keys | self.dispatch_attribute_keys

def diagnostics(self) -> tuple[str, ...]:
messages: list[str] = []
if (
Expand Down Expand Up @@ -623,6 +620,19 @@ def _env_int(
return parsed


def _attribute_key_errors(
errors: list[str],
name: str,
values: object,
) -> None:
if not isinstance(values, frozenset):
errors.append(f"{name} must be a frozenset of non-empty strings.")
return
for value in values:
if not isinstance(value, str) or not value.strip():
errors.append(f"{name} entries must be non-empty strings.")


def _choice_error(
errors: list[str], name: str, value: Any, choices: set[str]
) -> None:
Expand Down
53 changes: 41 additions & 12 deletions policyengine_observability/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -317,10 +317,7 @@ def set_context(self, **attributes: Any) -> None:
safe, omitted = normalize_attributes(
attributes,
self.config,
allowed_keys=(
self.config.application_attribute_keys
| self.config.dispatch_attribute_keys
),
allowed_keys=self.config.local_attribute_keys,
)
request = self._request_state.get()
operation = self._operation_state.get()
Expand Down Expand Up @@ -507,10 +504,7 @@ def _start_operation(
safe, omitted = normalize_attributes(
attributes,
self.config,
allowed_keys=(
self.config.application_attribute_keys
| self.config.dispatch_attribute_keys
),
allowed_keys=self.config.local_attribute_keys,
)
if omitted:
self.diagnostics.increment("attributes.omitted", omitted)
Expand All @@ -519,6 +513,20 @@ def _start_operation(
request_id: str | None = None
if remote_context:
request_id = _valid_request_id(remote_context.get("request_id"))
remote_attributes, remote_omitted = normalize_attributes(
{
key: remote_context.get(key)
for key in self.config.dispatch_attribute_keys
if key in remote_context
},
self.config,
allowed_keys=self.config.dispatch_attribute_keys,
)
if remote_omitted:
self.diagnostics.increment(
"attributes.omitted", remote_omitted
)
safe = {**remote_attributes, **safe}
carrier = {
key: str(value)
for key, value in remote_context.items()
Expand Down Expand Up @@ -615,13 +623,11 @@ def _start_child_span(
safe, omitted = normalize_attributes(
attributes,
self.config,
allowed_keys=(
self.config.application_attribute_keys
| self.config.dispatch_attribute_keys
),
allowed_keys=self.config.local_attribute_keys,
)
if omitted:
self.diagnostics.increment("attributes.omitted", omitted)
safe = {**safe, **self._active_dispatch_attributes()}
return _ChildSpanState(
name=name,
start_time=time.perf_counter(),
Expand Down Expand Up @@ -661,6 +667,20 @@ def _active_context_fields(self) -> dict[str, Any]:
fields.update(self._otel.current_correlation())
return fields

def _active_dispatch_attributes(self) -> dict[str, Any]:
attributes: dict[str, Any] = {}
request = self._request_state.get()
operation = self._operation_state.get()
if request is not None:
attributes.update(request.attributes)
if operation is not None:
attributes.update(operation.attributes)
return {
key: attributes[key]
for key in self.config.dispatch_attribute_keys
if key in attributes
}

def _metric_base(self) -> dict[str, Any]:
return {
"service.name": self.config.service.name,
Expand All @@ -673,6 +693,15 @@ def _metric_base(self) -> dict[str, Any]:

def _emit_record(self, **kwargs: Any) -> None:
try:
supplied_attributes = kwargs.get("attributes")
kwargs["attributes"] = {
**(
dict(supplied_attributes)
if supplied_attributes is not None
else {}
),
**self._active_dispatch_attributes(),
}
self._delivery.emit(build_record(self.config, **kwargs))
except Exception as exc:
self.diagnostics.report("record.emit", exc)
Expand Down
4 changes: 1 addition & 3 deletions policyengine_observability/schema.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,6 @@ def normalize_attributes(
key = str(raw_key).strip()
if (
not key
or len(normalized) >= config.limits.max_attributes
or _prohibited_key(key)
or (allowed_keys is not None and key not in allowed_keys)
):
Expand Down Expand Up @@ -100,8 +99,7 @@ def build_record(
safe_attributes, omitted = normalize_attributes(
attributes,
config,
allowed_keys=config.application_attribute_keys
| config.dispatch_attribute_keys,
allowed_keys=config.local_attribute_keys,
)
if safe_attributes:
record["attributes"] = safe_attributes
Expand Down
Loading
Loading