diff --git a/docs/design.md b/docs/design.md index fea18970..f81d2cd8 100644 --- a/docs/design.md +++ b/docs/design.md @@ -106,7 +106,7 @@ from microsoft_agents_a365.observability.core import ( configure( service_name="my-agent", service_namespace="my-namespace", - token_resolver=lambda agent_id, tenant_id: get_auth_token(), + token_resolver=lambda agent_id, tenant_id: get_app_only_obs_token(agent_id, tenant_id), cluster_category="prod" ) @@ -388,8 +388,8 @@ BatchSpanProcessor ← Accumulate spans ▼ Agent365Exporter.export() ← Send to backend ├── Partition by (tenant_id, agent_id) - ├── Resolve endpoint via PowerPlatformApiDiscovery - └── POST to /maven/agent365/agents/{agentId}/traces + ├── Resolve app-only OBS token via token_resolver(agent_id, tenant_id) + └── POST to /observabilityService/tenants/{tenantId}/otlp/agents/{agentId}/traces ``` ### MCP Tool Discovery Flow @@ -434,7 +434,7 @@ from microsoft_agents_a365.observability.core import configure configure( service_name="my-agent", service_namespace="my-namespace", - token_resolver=lambda agent_id, tenant_id: get_token(), + token_resolver=lambda agent_id, tenant_id: get_app_only_obs_token(agent_id, tenant_id), cluster_category="prod", # or "ppe", "test" # Advanced options via exporter_options parameter: # max_queue_size=2048, diff --git a/docs/integrating-with-existing-opentelemetry.md b/docs/integrating-with-existing-opentelemetry.md index 9e3fb123..e72498a3 100644 --- a/docs/integrating-with-existing-opentelemetry.md +++ b/docs/integrating-with-existing-opentelemetry.md @@ -4,7 +4,7 @@ This guide is for developers whose application **already** initializes OpenTelem ## The integration rule -> **Initialize your existing OpenTelemetry stack first, then call Agent 365's `configure()`.** The SDK detects the existing `TracerProvider` and adds its processors to it. Your existing backend receives every span; the Agent 365 backend also receives spans when `ENABLE_A365_OBSERVABILITY_EXPORTER=true` and a `token_resolver` is provided (otherwise `configure()` falls back to `ConsoleSpanExporter`). +> **Initialize your existing OpenTelemetry stack first, then call Agent 365's `configure()`.** The SDK detects the existing `TracerProvider` and adds its processors to it. Your existing backend receives every span; the Agent 365 backend also receives spans when `ENABLE_A365_OBSERVABILITY_EXPORTER=true` and an app-only OBS `token_resolver` is provided (otherwise `configure()` falls back to `ConsoleSpanExporter`). If the Agent 365 exporter is enabled without a resolver, the console fallback is kept and nothing is sent to Agent 365. The detection happens in [`config.py`](../libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/config.py): if a real (non-no-op) `TracerProvider` is already set (detected via a non-None `resource` attribute), `configure()` adds an `_EnrichingBatchSpanProcessor` (wrapping the configured exporter) and a custom `SpanProcessor` to that provider rather than creating a new one. @@ -25,7 +25,7 @@ configure_azure_monitor(connection_string=os.environ["APPLICATIONINSIGHTS_CONNEC configure( service_name="my-agent", service_namespace="my-namespace", - token_resolver=my_token_resolver, + token_resolver=my_app_only_obs_token_resolver, ) ``` @@ -55,12 +55,19 @@ trace.set_tracer_provider(provider) configure( service_name="my-agent", service_namespace="my-namespace", - token_resolver=my_token_resolver, + token_resolver=my_app_only_obs_token_resolver, ) ``` → Runnable version: [`observability-with-otlp`](https://github.com/microsoft/Agent365-Samples/tree/main/python/observability-with-otlp) sample (defaults to `ConsoleSpanExporter` for zero setup). +The Agent 365 backend exporter always uses +`/observabilityService/tenants/{tenantId}/otlp/agents/{agentId}/traces?api-version=1`. +The resolver must return an app-only OBS token for the exporting agent identity; +delegated workload/OBO tokens are not read from request context and are rejected +by the S2S service. Cache the resolver result and refresh near expiry because it +is invoked for each export batch/identity group. + ## Auto-instrumentation vs. manual instrumentation The OTel **backend** (where spans go) and the **instrumentation style** (how spans are produced) are independent axes. You can mix them freely. diff --git a/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md b/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md index dcf0f6a9..1159082d 100644 --- a/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md +++ b/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md @@ -2,6 +2,33 @@ All notable changes to this package will be documented in this file. +## [Unreleased] + +### Breaking Changes + +- **OBS exports always use `/observabilityService`** — The `use_s2s_endpoint` + option is deprecated and ignored, even when `False`. Exports no longer select + or fall back to `/observability`. Provide an app-only OBS token independently + of your agent's workload auth; the S2S service rejects delegated `scp` tokens. +- **OBS export requires the configured app-only resolver** — The exporter uses + only `token_resolver(agent_id, tenant_id)` for Agent 365 authentication. If + the Agent 365 exporter is enabled without a resolver, the existing console + fallback is kept and nothing is sent to Agent 365. Empty tokens or acquisition + failures fail export without an HTTP request or delegated fallback. The + exporter invokes the resolver on every export batch/identity group, so + resolvers must cache the acquired token and refresh only near expiry. + +### Migration + +- Request the OBS resource `/.default` scope + (`api://9b975845-388f-4429-889e-eab1ef63949c/.default`) for the exporting + agent identity. Do not use delegated `Agent365.Observability.OtelWrite` tokens + for OBS export. +- Validate resolver tokens before returning them: reject any `scp` claim and any + `idtyp` other than `app`; when `idtyp` is absent, accept only a non-empty + `roles` array or a non-empty `oid` equal to `sub`; verify `aud` is the OBS + resource and the token is not expired. + ## [0.3.0] ### Breaking Changes diff --git a/libraries/microsoft-agents-a365-observability-core/README.md b/libraries/microsoft-agents-a365-observability-core/README.md index 529d74a9..51a0300c 100644 --- a/libraries/microsoft-agents-a365-observability-core/README.md +++ b/libraries/microsoft-agents-a365-observability-core/README.md @@ -17,6 +17,35 @@ pip install microsoft-agents-a365-observability-core For usage examples and detailed documentation, see the [Observability documentation](https://learn.microsoft.com/microsoft-agent-365/developer/observability?tabs=python) on Microsoft Learn. +### Agent 365 OBS export authentication + +When `ENABLE_A365_OBSERVABILITY_EXPORTER` is enabled, exports always use the +S2S OTLP route: + +```text +/observabilityService/tenants/{tenantId}/otlp/agents/{agentId}/traces?api-version=1 +``` + +The deprecated `use_s2s_endpoint` option is ignored, even when set to `False`; +domain overrides change only the host. The exporter never falls back to +`/observability` and never reads delegated request-context tokens. + +Provide an app-only OBS `token_resolver(agent_id, tenant_id)` for the exporting +agent identity. The resolver is invoked for each export batch and identity +group, so it should cache the acquired token and refresh near expiry. Empty +tokens or resolver failures fail that export batch without sending an HTTP +request or retrying on a delegated route. If the Agent 365 exporter is enabled +without a resolver, the existing `ConsoleSpanExporter` fallback is kept and +nothing is sent to Agent 365. + +Resolvers should request the OBS resource `/.default` scope +(`api://9b975845-388f-4429-889e-eab1ef63949c/.default`) and validate the token +before returning it: reject any `scp` claim and any `idtyp` other than `app`; +when `idtyp` is absent, accept only a non-empty `roles` array or a non-empty +`oid` equal to `sub`; also verify `aud` is the OBS resource and the token is not +expired. Workload authentication for MCP, Microsoft Graph, and other OBO calls +is separate and unchanged. + ## Support For issues, questions, or feedback: @@ -33,4 +62,3 @@ For issues, questions, or feedback: Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT License - see the [LICENSE](../../LICENSE.md) file for details. - diff --git a/libraries/microsoft-agents-a365-observability-core/docs/design.md b/libraries/microsoft-agents-a365-observability-core/docs/design.md index cc1944ab..570eeabd 100644 --- a/libraries/microsoft-agents-a365-observability-core/docs/design.md +++ b/libraries/microsoft-agents-a365-observability-core/docs/design.md @@ -44,7 +44,7 @@ from microsoft_agents_a365.observability.core import configure configure( service_name="my-agent", service_namespace="my-namespace", - token_resolver=lambda agent_id, tenant_id: get_token(), + token_resolver=lambda agent_id, tenant_id: get_app_only_obs_token(agent_id, tenant_id), cluster_category="prod" ) ``` @@ -219,12 +219,27 @@ This ensures that context values set via `BaggageBuilder` are recorded as span a **Export flow:** 1. Partition spans by `(tenant_id, agent_id)` tuple 2. For each partition: - - Resolve endpoint via `PowerPlatformApiDiscovery` - - Resolve auth token via `token_resolver(agent_id, tenant_id)` + - Resolve endpoint host via the configured domain override or production endpoint + - Build the S2S OTLP URL `/observabilityService/tenants/{tenant_id}/otlp/agents/{agent_id}/traces?api-version=1` + - Resolve an app-only OBS token via `token_resolver(agent_id, tenant_id)` - Build OTLP-like JSON payload - - POST to `/maven/agent365/agents/{agentId}/traces` + - POST to the S2S OTLP route 3. Retry transient failures (408, 429, 5xx) up to 3 times with exponential backoff +The exporter always uses the S2S `/observabilityService` route. The +`use_s2s_endpoint` option is deprecated and ignored, including when `False`. +There is no delegated `/observability` fallback for batch export, request export, +401/403/404 responses, missing tokens, or token acquisition failures. + +`token_resolver` must return an app-only OBS token for the exporting agent and +tenant. It is invoked on every export batch/identity group, so resolvers should +cache and refresh tokens near expiry. Resolvers should request the OBS resource +`/.default` scope (`api://9b975845-388f-4429-889e-eab1ef63949c/.default`) and +validate the token before returning it: reject any `scp` claim and any `idtyp` +other than `app`; when `idtyp` is absent, accept only a non-empty `roles` array +or a non-empty `oid` equal to `sub`; verify `aud` is the OBS resource and the +token is not expired. + **Configuration via `Agent365ExporterOptions`:** ```python from microsoft_agents_a365.observability.core.exporters import Agent365ExporterOptions @@ -232,7 +247,7 @@ from microsoft_agents_a365.observability.core.exporters import Agent365ExporterO options = Agent365ExporterOptions( cluster_category="prod", token_resolver=my_token_resolver, - use_s2s_endpoint=False, + use_s2s_endpoint=False, # Deprecated and ignored; export still uses S2S. max_queue_size=2048, scheduled_delay_ms=5000, exporter_timeout_ms=30000, diff --git a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/config.py b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/config.py index e1757127..3d124a05 100644 --- a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/config.py +++ b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/config.py @@ -4,7 +4,6 @@ import os import logging import threading -from collections.abc import Callable from typing import Any, Optional from opentelemetry import trace @@ -17,7 +16,7 @@ from opentelemetry.sdk.trace.export import ConsoleSpanExporter from .exporters.agent365_exporter import _Agent365Exporter -from .exporters.agent365_exporter_options import Agent365ExporterOptions +from .exporters.agent365_exporter_options import Agent365ExporterOptions, TokenResolver from .exporters.enriching_span_processor import ( _EnrichingBatchSpanProcessor, ) @@ -60,7 +59,7 @@ def configure( service_name: str, service_namespace: str, logger_name: str = DEFAULT_LOGGER_NAME, - token_resolver: Callable[[str, str], str | None] | None = None, + token_resolver: TokenResolver | None = None, cluster_category: str = "prod", exporter_options: Agent365ExporterOptions | SpectraExporterOptions | None = None, suppress_invoke_agent_input: bool = False, @@ -72,8 +71,8 @@ def configure( :param service_name: The name of the service. :param service_namespace: The namespace of the service. :param logger_name: The name of the logger to collect telemetry from. - :param token_resolver: (Deprecated) Callable that returns an auth token for a given agent + tenant. - Use exporter_options instead. + :param token_resolver: (Deprecated) Callable that returns an app-only OBS token + for a given agent + tenant. Use exporter_options instead. :param cluster_category: (Deprecated) Environment / cluster category (e.g. "prod"). Use exporter_options instead. :param exporter_options: Exporter configuration. Pass Agent365ExporterOptions for A365 API @@ -103,7 +102,7 @@ def _configure_internal( service_name: str, service_namespace: str, logger_name: str, - token_resolver: Callable[[str, str], str | None] | None = None, + token_resolver: TokenResolver | None = None, cluster_category: str = "prod", exporter_options: Agent365ExporterOptions | SpectraExporterOptions | None = None, suppress_invoke_agent_input: bool = False, @@ -270,7 +269,7 @@ def configure( service_name: str, service_namespace: str, logger_name: str = DEFAULT_LOGGER_NAME, - token_resolver: Callable[[str, str], str | None] | None = None, + token_resolver: TokenResolver | None = None, cluster_category: str = "prod", exporter_options: Agent365ExporterOptions | SpectraExporterOptions | None = None, suppress_invoke_agent_input: bool = False, @@ -282,8 +281,8 @@ def configure( :param service_name: The name of the service. :param service_namespace: The namespace of the service. :param logger_name: The name of the logger to collect telemetry from. - :param token_resolver: (Deprecated) Callable that returns an auth token for a given agent + tenant. - Use exporter_options instead. + :param token_resolver: (Deprecated) Callable that returns an app-only OBS token for a given + agent + tenant. Use exporter_options instead. :param cluster_category: (Deprecated) Environment / cluster category (e.g. "prod"). Use exporter_options instead. :param exporter_options: Exporter configuration. Pass Agent365ExporterOptions for A365 API diff --git a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/__init__.py b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/__init__.py index 9b0ab065..69cdb993 100644 --- a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/__init__.py +++ b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/__init__.py @@ -1,9 +1,9 @@ # Copyright (c) Microsoft Corporation. # Licensed under the MIT License. -from .agent365_exporter_options import Agent365ExporterOptions +from .agent365_exporter_options import Agent365ExporterOptions, TokenResolver from .spectra_exporter_options import SpectraExporterOptions # Agent365Exporter is not exported intentionally. # It should only be used internally by the observability core module. -__all__ = ["Agent365ExporterOptions", "SpectraExporterOptions"] +__all__ = ["Agent365ExporterOptions", "SpectraExporterOptions", "TokenResolver"] diff --git a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter.py b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter.py index 1acbee61..0313c023 100644 --- a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter.py +++ b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter.py @@ -5,18 +5,21 @@ from __future__ import annotations +import asyncio +import inspect import json import logging import threading import time -from collections.abc import Callable, Sequence -from typing import Any, final +from collections.abc import Awaitable, Sequence +from typing import Any, cast, final import requests from opentelemetry.sdk.trace import ReadableSpan from opentelemetry.sdk.trace.export import SpanExporter, SpanExportResult from opentelemetry.trace import StatusCode +from .agent365_exporter_options import TokenResolver from .utils import ( DEFAULT_MAX_PAYLOAD_BYTES, build_export_url, @@ -43,20 +46,23 @@ logger = logging.getLogger(__name__) +async def _await_token(awaitable: Awaitable[str | None]) -> str | None: + return await awaitable + + @final class _Agent365Exporter(SpanExporter): """ Agent 365 span exporter for Agent 365: * Partitions spans by (tenantId, agentId) * Builds OTLP-like JSON: resourceSpans -> scopeSpans -> spans - * POSTs per group to https://{endpoint}/observability/tenants/{tenantId}/otlp/agents/{agentId}/traces?api-version=1 - * or, when use_s2s_endpoint is True, https://{endpoint}/observabilityService/tenants/{tenantId}/otlp/agents/{agentId}/traces?api-version=1 - * Adds Bearer token via token_resolver(agentId, tenantId) + * POSTs per group to the S2S /observabilityService OTLP route. + * Adds an app-only authorization token via token_resolver(agentId, tenantId) """ def __init__( self, - token_resolver: Callable[[str, str], str | None], + token_resolver: TokenResolver, cluster_category: str = "prod", use_s2s_endpoint: bool = False, max_payload_bytes: int = DEFAULT_MAX_PAYLOAD_BYTES, @@ -126,8 +132,8 @@ def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: headers = {"content-type": "application/json"} try: - token = self._token_resolver(agent_id, tenant_id) - if token: + token = self._resolve_token(agent_id, tenant_id) + if token is not None and token.strip(): # Warn if sending bearer token over non-HTTPS connection if not url.lower().startswith("https://"): logger.warning( @@ -135,13 +141,21 @@ def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: "This may expose credentials in transit." ) headers["authorization"] = f"Bearer {token}" - logger.debug(f"Token resolved successfully for agent {agent_id}") + logger.debug( + f"App-only OBS token resolved successfully for agent {agent_id}" + ) else: - logger.debug(f"No token returned for agent {agent_id}") + logger.error( + f"No app-only OBS token returned for agent {agent_id}, " + f"tenant {tenant_id}; export request will not be sent" + ) + any_failure = True + continue except Exception as e: # If token resolution fails, treat as failure for this group logger.error( - f"Token resolution failed for agent {agent_id}, tenant {tenant_id}: {e}" + f"App-only OBS token resolution failed for agent {agent_id}, " + f"tenant {tenant_id}: {e}" ) any_failure = True continue @@ -276,6 +290,24 @@ def _post_with_retries(self, url: str, body: str, headers: dict[str, str]) -> bo return False return False + def _resolve_token(self, agent_id: str, tenant_id: str) -> str | None: + token = self._token_resolver(agent_id, tenant_id) + if inspect.isawaitable(token): + try: + asyncio.get_running_loop() + except RuntimeError: + return asyncio.run(_await_token(cast(Awaitable[str | None], token))) + if inspect.iscoroutine(token): + token.close() + elif isinstance(token, asyncio.Future): + token.cancel() + raise RuntimeError( + "Agent365Exporter cannot await an async token_resolver while running " + "inside an active event loop; use a synchronous cached resolver or refresh " + "the app-only OBS token before export." + ) + return token + # ------------- Payload mapping ------------------ def _map_and_truncate_spans( diff --git a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter_options.py b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter_options.py index a3981df8..e951968a 100644 --- a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter_options.py +++ b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter_options.py @@ -5,6 +5,8 @@ from .utils import DEFAULT_MAX_PAYLOAD_BYTES +TokenResolver = Callable[[str, str], str | Awaitable[str | None] | None] + class Agent365ExporterOptions: """ @@ -15,7 +17,7 @@ class Agent365ExporterOptions: def __init__( self, cluster_category: str = "prod", - token_resolver: Optional[Callable[[str, str], Awaitable[Optional[str]]]] = None, + token_resolver: Optional[TokenResolver] = None, use_s2s_endpoint: bool = False, max_queue_size: int = 2048, scheduled_delay_ms: int = 5000, @@ -26,8 +28,11 @@ def __init__( """ Args: cluster_category: Cluster region argument. Defaults to 'prod'. - token_resolver: Async callable that resolves the auth token (REQUIRED). - use_s2s_endpoint: Use the S2S endpoint instead of standard endpoint. + token_resolver: Callable that resolves an app-only OBS token (REQUIRED when the + Agent 365 exporter is enabled). It is invoked for each export batch/identity + group, so implementations should cache and refresh tokens near expiry. + use_s2s_endpoint: Deprecated compatibility option. Ignored; export always uses + the S2S /observabilityService OTLP route, even when this is False. max_queue_size: Maximum queue size for the batch processor. Default is 2048. scheduled_delay_ms: Delay between export batches (ms). Default is 5000. exporter_timeout_ms: Timeout for the export operation (ms). Default is 30000. diff --git a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/utils.py b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/utils.py index 1e357176..d3f140e2 100644 --- a/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/utils.py +++ b/libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/utils.py @@ -258,16 +258,13 @@ def build_export_url( endpoint: Base endpoint URL or domain. agent_id: The agent identifier to include in the URL path. tenant_id: The tenant identifier to include in the URL path. - use_s2s_endpoint: Whether to use the S2S endpoint path format. + use_s2s_endpoint: Deprecated compatibility option. Ignored; export always uses + the S2S endpoint path format. Returns: The fully constructed export URL with path and query parameters. """ - endpoint_path = ( - f"/observabilityService/tenants/{tenant_id}/otlp/agents/{agent_id}/traces" - if use_s2s_endpoint - else f"/observability/tenants/{tenant_id}/otlp/agents/{agent_id}/traces" - ) + endpoint_path = f"/observabilityService/tenants/{tenant_id}/otlp/agents/{agent_id}/traces" parsed = urlparse(endpoint) if parsed.scheme and "://" in endpoint: diff --git a/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md index 456a642a..ab08f4f8 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md +++ b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md @@ -4,6 +4,27 @@ All notable changes to this package will be documented in this file. ## [Unreleased] +### Breaking Changes + +- **Hosting OBS token cache requires an app-only resolver** — + `refresh_observability_token(agent_id, tenant_id, token_resolver)` acquires + and caches app-only OBS tokens for export. The resolver receives the + configured OBS scopes and must acquire a token for the exporting agent + identity, not its blueprint or the workload's user. Acquisition failures + propagate to the caller; empty tokens fail refresh and clear stale cache state. +- **`get_observability_token(...)` no longer exchanges tokens** — It returns the + cached app-only token acquired by `refresh_observability_token`, or `None`. +- **Delegated OBS registration is a no-op** — `register_observability(...)` + logs one error and stores no token. It never calls `Authorization.exchange_token` + and never acquires a delegated OBS token. + +### Migration + +- Call and return `refresh_observability_token(agent_id, tenant_id, app_only_token_resolver)` + from the exporter `token_resolver`. +- Keep workload MCP/Graph/OBO authentication unchanged; only OBS export + authentication moves to app-only S2S. + ### Changed - **`OutputLoggingMiddleware`** — Updated to use new scope APIs (`Request`, `SpanDetails`, `UserDetails`). Removed `TenantDetails` and `ExecutionType` dependencies. Middleware no longer gates on tenant presence. diff --git a/libraries/microsoft-agents-a365-observability-hosting/README.md b/libraries/microsoft-agents-a365-observability-hosting/README.md index 2b3ed82a..8984b3e5 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/README.md +++ b/libraries/microsoft-agents-a365-observability-hosting/README.md @@ -7,3 +7,62 @@ This library provides hosting components for Agent 365 observability. ```bash pip install microsoft-agents-a365-observability-hosting ``` + +## App-only OBS token cache + +Agent 365 observability export is S2S-only. Use +`AgenticTokenCache.refresh_observability_token(agent_id, tenant_id, token_resolver)` +from the exporter token resolver to acquire and cache an app-only OBS token for +the exporting agent identity. + +```python +from microsoft_agents_a365.observability.hosting.token_cache_helpers import AgenticTokenCache + +cache = AgenticTokenCache() + + +async def acquire_app_only_obs_token(agent_id: str, tenant_id: str, scopes: list[str]) -> str: + # Acquire a final app-only OBS token for agent_id and tenant_id. + ... + + +async def token_resolver(agent_id: str, tenant_id: str) -> str: + return await cache.refresh_observability_token(agent_id, tenant_id, acquire_app_only_obs_token) +``` + +The resolver receives the OBS `/.default` scope and must not perform user_fic or +OBO authentication. For blueprint-backed agents, acquire the final agent +identity app-only token before returning it; do not return the intermediate +blueprint assertion or a delegated workload token. The cache retries transient +acquisition failures, isolates entries by `(agent_id, tenant_id)`, respects JWT +expiry with refresh skew, and uses a fallback TTL for opaque tokens. + +If you use an async exporter resolver, `_Agent365Exporter` runs it with +`asyncio.run` on the thread that performs the export: + +- Scheduled batch exports and the final export during `shutdown()` run on the + BatchSpanProcessor worker thread. `shutdown()` still blocks its caller until + that export finishes. +- `force_flush()` exports on the calling thread. If you call it from inside a + running event loop, the async resolver can't be awaited and that export + fails with a logged error. + +Create async clients inside the resolver; don't reuse `aiohttp` or +`azure.identity.aio` clients bound to your app's event loop. From async code, +call `force_flush()` and `shutdown()` with `await asyncio.to_thread(...)`. A +synchronous, thread-safe cached resolver avoids these event-loop restrictions +and is the pattern the Agent365-Samples Python samples use. +Call `refresh_observability_token` only from the exporter's `token_resolver`; +its per-key locks are `asyncio.Lock`s, so do not also refresh the same cache +instance from your app's own event loop or threads. + +The previous delegated registration shape using `TurnContext` and +`Authorization.exchange_token` is removed for OBS export. `register_observability(...)` +is accepted only as a no-op compatibility shape: the cache logs once, stores no +token, and never calls `exchange_token`. + +Token resolvers should validate returned tokens before caching: reject any `scp` +claim and any `idtyp` other than `app`; when `idtyp` is absent, accept only a +non-empty `roles` array or a non-empty `oid` equal to `sub`; also verify `aud` +is the OBS resource (`9b975845-388f-4429-889e-eab1ef63949c`) and the token is +not expired. diff --git a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/__init__.py b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/__init__.py index 23daf0a3..ea1f5850 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/__init__.py +++ b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/__init__.py @@ -2,6 +2,6 @@ # Licensed under the MIT License. """Token cache helpers for observability.""" -from .agent_token_cache import AgenticTokenCache, AgenticTokenStruct +from .agent_token_cache import AgenticTokenCache, AgenticTokenStruct, ObservabilityTokenResolver -__all__ = ["AgenticTokenCache", "AgenticTokenStruct"] +__all__ = ["AgenticTokenCache", "AgenticTokenStruct", "ObservabilityTokenResolver"] diff --git a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/agent_token_cache.py b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/agent_token_cache.py index 00f6f9d4..d317f304 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/agent_token_cache.py +++ b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/token_cache_helpers/agent_token_cache.py @@ -1,56 +1,86 @@ # Copyright (c) Microsoft Corporation. # Licensed under the MIT License. -""" -Token cache for observability tokens per (agentId, tenantId). -""" +"""App-only token cache for Agent 365 observability export.""" from __future__ import annotations +import asyncio +import base64 +import json import logging -from dataclasses import dataclass +import time +from collections.abc import Awaitable, Callable, Sequence +from dataclasses import dataclass, field +from inspect import isawaitable from threading import Lock from microsoft_agents.hosting.core.app.oauth.authorization import Authorization from microsoft_agents.hosting.core.turn_context import TurnContext +from microsoft_agents_a365.runtime.environment_utils import get_observability_authentication_scope logger = logging.getLogger(__name__) +ObservabilityTokenResolver = Callable[[str, str, list[str]], str | Awaitable[str | None] | None] + @dataclass class AgenticTokenStruct: - """Structure containing the token generation components.""" + """Deprecated delegated OBS token generator shape. + + OBS export is S2S-only. Instances of this type are accepted only for source + compatibility; the cache never calls ``authorization.exchange_token``. + """ authorization: Authorization - """The user authorization object for token exchange.""" + """The user authorization object from the removed delegated OBS flow.""" turn_context: TurnContext - """The turn context for the current conversation.""" + """The turn context from the removed delegated OBS flow.""" auth_handler_name: str | None = "AGENTIC" - """The name of the authentication handler.""" + """The name of the removed delegated authentication handler.""" class AgenticTokenCache: - """ - Caches observability tokens per (agentId, tenantId) using the provided - UserAuthorization and TurnContext. + """Caches app-only OBS tokens per ``(agent_id, tenant_id)``. + + Call :meth:`refresh_observability_token` from the exporter's token resolver + to acquire or refresh an app-only token. The resolver receives the exporting + agent ID, tenant ID, and OBS ``/.default`` scopes. Delegated TurnContext / + Authorization shapes are logged once and ignored. """ @dataclass class _Entry: """Internal entry structure for cache storage.""" - agentic_token_struct: AgenticTokenStruct - """The token generation structure.""" + scopes: tuple[str, ...] + token: str | None = None + expires_on_ms: float | None = None + acquired_on_ms: float | None = None + lock: asyncio.Lock = field(default_factory=asyncio.Lock, repr=False, compare=False) - scopes: list[str] - """The observability scopes for token requests.""" + _default_refresh_skew_ms = 60_000 + _default_max_token_age_ms = 3_600_000 + _max_exp_seconds = 86_400 + _max_cache_size = 10_000 - def __init__(self) -> None: + def __init__(self, observability_scopes: Sequence[str] | None = None) -> None: """Initialize the token cache.""" - self._map: dict[str, AgenticTokenCache._Entry] = {} + if isinstance(observability_scopes, str): + observability_scopes = (observability_scopes,) + self._map: dict[tuple[str, str], AgenticTokenCache._Entry] = {} self._lock = Lock() + self._observability_scopes = ( + None if observability_scopes is None else tuple(observability_scopes) + ) + self._removed_registration_logged = False + + @staticmethod + def _make_key(agent_id: str, tenant_id: str) -> tuple[str, str]: + # A tuple key keeps identities apart even when an ID contains a separator character. + return (agent_id, tenant_id) def register_observability( self, @@ -59,79 +89,224 @@ def register_observability( token_generator: AgenticTokenStruct, observability_scopes: list[str], ) -> None: - """ - Register observability for the specified agent and tenant. + """Deprecated no-op for the removed delegated OBS registration flow.""" + self._log_removed_registration_once() + + async def refresh_observability_token( + self, + agent_id: str, + tenant_id: str, + token_resolver: ObservabilityTokenResolver, + ) -> str: + """Refresh an app-only OBS token for the exporting agent identity. Args: - agent_id: The agent identifier. - tenant_id: The tenant identifier. - token_generator: The token generator structure. - observability_scopes: The observability scopes. + agent_id: The exporting agent instance identifier. + tenant_id: The exporting tenant identifier. + token_resolver: App-only token resolver receiving ``(agent_id, + tenant_id, scopes)``. + + Returns: + The cached app-only OBS token. Raises: - ValueError: If agent_id or tenant_id is empty or None. - TypeError: If token_generator is None. + TypeError: If ``token_resolver`` is not callable. + ValueError: If agent or tenant IDs are empty, or no scopes are configured. + Exception: Propagates resolver failures after retry handling. """ - if not agent_id or not agent_id.strip(): - raise ValueError("agent_id cannot be None or whitespace") + if not callable(token_resolver): + raise TypeError("token_resolver must be callable") - if not tenant_id or not tenant_id.strip(): - raise ValueError("tenant_id cannot be None or whitespace") + if not agent_id or not agent_id.strip() or not tenant_id or not tenant_id.strip(): + raise ValueError("[AgenticTokenCache] Agent and tenant IDs are required") - if token_generator is None: - raise TypeError("token_generator cannot be None") + key = self._make_key(agent_id, tenant_id) + entry = self._get_or_create_entry(key) + async with entry.lock: + token = entry.token + if token is not None and not self._is_expired(entry): + return token - key = f"{agent_id}:{tenant_id}" - - # First registration wins; subsequent calls ignored (idempotent) - with self._lock: - if key not in self._map: - self._map[key] = AgenticTokenCache._Entry( - agentic_token_struct=token_generator, scopes=observability_scopes - ) - logger.debug(f"Registered observability for {key}") - else: - logger.debug(f"Observability already registered for {key}, ignoring") + return await self._acquire_token(agent_id, tenant_id, entry, token_resolver) async def get_observability_token(self, agent_id: str, tenant_id: str) -> str | None: + """Return a non-expired cached app-only OBS token, or ``None``. + + This method is a pure cache read. It never acquires a token and never + calls delegated token exchange. """ - Get the observability token for the specified agent and tenant. + key = self._make_key(agent_id, tenant_id) + with self._lock: + entry = self._map.get(key) - Args: - agent_id: The agent identifier. - tenant_id: The tenant identifier. + if entry is None or entry.token is None: + logger.debug("[AgenticTokenCache] No token cached for %s", key) + return None + if self._is_expired(entry): + logger.debug("[AgenticTokenCache] Token expired for %s", key) + return None + return entry.token - Returns: - The observability token if available; otherwise, None. + def invalidate_token(self, agent_id: str, tenant_id: str) -> None: + """Invalidate one cached token. + + The entry is removed, so a refresh already in flight for this identity + cannot repopulate the cache. """ - key = f"{agent_id}:{tenant_id}" + key = self._make_key(agent_id, tenant_id) + with self._lock: + entry = self._map.pop(key, None) + if entry is not None: + self._clear_token(entry) - logger.debug(f"Cache lookup for {key}") + def invalidate_all(self) -> None: + """Invalidate all cached tokens.""" + with self._lock: + self._map.clear() + def _get_or_create_entry(self, key: tuple[str, str]) -> _Entry: with self._lock: entry = self._map.get(key) + if entry is not None: + if not entry.scopes: + raise ValueError("[AgenticTokenCache] Entry has invalid scopes") + return entry + + scopes = self._get_effective_scopes() + if not scopes: + raise ValueError("[AgenticTokenCache] No valid scopes") + + # Evict the oldest idle entries until there is room. Entries with a refresh in flight keep + # their lock, so the cache can briefly exceed its bound; the next insertion trims it back. + while len(self._map) >= self._max_cache_size: + idle_key = next( + ( + existing_key + for existing_key, existing in self._map.items() + if not existing.lock.locked() + ), + None, + ) + if idle_key is None: + break + del self._map[idle_key] - if entry is None: - logger.debug(f"Cache miss for {key}") - return None + entry = AgenticTokenCache._Entry(scopes=scopes) + self._map[key] = entry + return entry - logger.debug(f"Cache hit for {key}, exchanging token") + def _get_effective_scopes(self) -> tuple[str, ...]: + scopes = self._observability_scopes + if scopes is None: + scopes = tuple(get_observability_authentication_scope()) + return tuple(scope for scope in scopes if scope and scope.strip()) - try: - authorization = entry.agentic_token_struct.authorization - turn_context = entry.agentic_token_struct.turn_context - auth_handler_id = entry.agentic_token_struct.auth_handler_name - - # Exchange the turn token for an observability token - token = await authorization.exchange_token( - context=turn_context, - scopes=entry.scopes, - auth_handler_id=auth_handler_id, + async def _acquire_token( + self, + agent_id: str, + tenant_id: str, + entry: _Entry, + resolver: ObservabilityTokenResolver, + ) -> str: + max_retries = 2 + last_error: BaseException | None = None + for attempt in range(max_retries + 1): + logger.info( + "[AgenticTokenCache] Acquiring app-only token attempt %s/%s", + attempt + 1, + max_retries + 1, ) - - logger.info(f"Successfully exchanged token for {key}") - return token - except Exception as e: - # Return None if token generation fails - logger.error(f"Token exchange failed for {key}: {type(e).__name__}") + try: + result = resolver(agent_id, tenant_id, list(entry.scopes)) + token = await result if isawaitable(result) else result + if token is None or not token.strip(): + raise RuntimeError( + "[AgenticTokenCache] App-only token resolver returned no token" + ) + entry.token = token + entry.acquired_on_ms = time.time() * 1000 + exp = self._decode_exp(token) + if exp is not None: + entry.expires_on_ms = exp * 1000 + else: + entry.expires_on_ms = None + logger.warning("[AgenticTokenCache] No exp claim, fallback TTL") + logger.info("[AgenticTokenCache] Token cached") + return token + except Exception as error: + last_error = error + if self._is_retriable_error(error) and attempt < max_retries: + logger.warning( + "[AgenticTokenCache] Retriable token acquisition failure attempt %s: %s", + attempt + 1, + error, + ) + await asyncio.sleep(0.2 * (attempt + 1)) + continue + logger.error("[AgenticTokenCache] Token acquisition failed: %s", error) + self._clear_token(entry) + raise + + self._clear_token(entry) + raise RuntimeError("[AgenticTokenCache] Token acquisition failed") from last_error + + def _decode_exp(self, jwt: str) -> int | None: + try: + parts = jwt.split(".") + if len(parts) < 2: + return None + payload = parts[1] + padded = payload + "=" * ((4 - (len(payload) % 4)) % 4) + decoded = base64.urlsafe_b64decode(padded.encode("utf-8")) + claims = json.loads(decoded.decode("utf-8")) + if not isinstance(claims, dict): + return None + exp = claims.get("exp") + if not isinstance(exp, int | float): + return None + max_exp = int(time.time()) + self._max_exp_seconds + return min(int(exp), max_exp) + except Exception: return None + + def _is_expired(self, entry: _Entry) -> bool: + now = time.time() * 1000 + if entry.expires_on_ms is not None: + return now >= entry.expires_on_ms - self._default_refresh_skew_ms + if entry.acquired_on_ms is not None: + return now >= entry.acquired_on_ms + self._default_max_token_age_ms + return True + + def _is_retriable_error(self, error: Exception) -> bool: + if isinstance(error, (TimeoutError, ConnectionError)): + return True + + message = str(error).lower() + if "timeout" in message or "econnreset" in message or "network" in message: + return True + + status = self._get_status(error) + return status in (408, 429) or (status is not None and 500 <= status < 600) + + def _get_status(self, error: Exception) -> int | None: + for name in ("status", "status_code"): + value = getattr(error, name, None) + if isinstance(value, int): + return value + return None + + def _clear_token(self, entry: _Entry) -> None: + entry.token = None + entry.expires_on_ms = None + entry.acquired_on_ms = None + + def _log_removed_registration_once(self) -> None: + with self._lock: + if self._removed_registration_logged: + return + self._removed_registration_logged = True + logger.error( + "[AgenticTokenCache] Delegated OBS token registration was removed and does " + "nothing; S2S OBS needs an app-only token. Call " + "refresh_observability_token(agent_id, tenant_id, token_resolver) instead." + ) diff --git a/libraries/microsoft-agents-a365-observability-hosting/pyproject.toml b/libraries/microsoft-agents-a365-observability-hosting/pyproject.toml index 4c4ac387..df25b8a6 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/pyproject.toml +++ b/libraries/microsoft-agents-a365-observability-hosting/pyproject.toml @@ -27,6 +27,7 @@ keywords = ["observability", "telemetry", "tracing", "opentelemetry", "monitorin dependencies = [ "microsoft-agents-hosting-core", "microsoft-agents-a365-observability-core", + "microsoft-agents-a365-runtime", "opentelemetry-api", ] diff --git a/libraries/microsoft-agents-a365-runtime/CHANGELOG.md b/libraries/microsoft-agents-a365-runtime/CHANGELOG.md new file mode 100644 index 00000000..f43df08e --- /dev/null +++ b/libraries/microsoft-agents-a365-runtime/CHANGELOG.md @@ -0,0 +1,14 @@ +# Changelog — microsoft-agents-a365-runtime + +All notable changes to this package will be documented in this file. + +## [Unreleased] + +### Breaking Changes + +- **OBS authentication scope is app-only** — + `get_observability_authentication_scope()` now returns the OBS resource + `/.default` scope (`api://9b975845-388f-4429-889e-eab1ef63949c/.default`) + instead of the delegated `Agent365.Observability.OtelWrite` scope. OBS export + is S2S-only; use this scope with an app-only resolver for the exporting agent + identity. Workload MCP/Graph/OBO authentication is unchanged. diff --git a/libraries/microsoft-agents-a365-runtime/docs/design.md b/libraries/microsoft-agents-a365-runtime/docs/design.md index d87e900f..50efd9f0 100644 --- a/libraries/microsoft-agents-a365-runtime/docs/design.md +++ b/libraries/microsoft-agents-a365-runtime/docs/design.md @@ -167,8 +167,14 @@ from microsoft_agents_a365.runtime import get_observability_authentication_scope # Get authentication scope for observability scope = get_observability_authentication_scope() +# Returns ["api://9b975845-388f-4429-889e-eab1ef63949c/.default"] ``` +Observability export is S2S-only. This helper returns the OBS resource +`/.default` scope for app-only token acquisition by the exporting agent identity; +it must not be used to acquire delegated `Agent365.Observability.OtelWrite` +tokens for OBS export. + ## Type Definitions ### ClusterCategory diff --git a/libraries/microsoft-agents-a365-runtime/microsoft_agents_a365/runtime/environment_utils.py b/libraries/microsoft-agents-a365-runtime/microsoft_agents_a365/runtime/environment_utils.py index bbb240ac..5bab292f 100644 --- a/libraries/microsoft-agents-a365-runtime/microsoft_agents_a365/runtime/environment_utils.py +++ b/libraries/microsoft-agents-a365-runtime/microsoft_agents_a365/runtime/environment_utils.py @@ -7,10 +7,12 @@ import os -# Authentication scopes for different environments -PROD_OBSERVABILITY_SCOPE = ( - "api://9b975845-388f-4429-889e-eab1ef63949c/Agent365.Observability.OtelWrite" -) +# Authentication scopes for different environments. +# +# Agent 365 observability export uses the app-only S2S route. Token resolvers must +# request the OBS resource's /.default scope for the exporting agent identity; they +# must not acquire delegated Agent365.Observability.OtelWrite tokens for OBS export. +PROD_OBSERVABILITY_SCOPE = "api://9b975845-388f-4429-889e-eab1ef63949c/.default" # Cluster categories for different environments PROD_OBSERVABILITY_CLUSTER_CATEGORY = "prod" @@ -25,7 +27,8 @@ def get_observability_authentication_scope() -> list[str]: Returns the scope for authenticating to the observability service based on the current environment. The scope can be overridden via the A365_OBSERVABILITY_SCOPE_OVERRIDE environment variable - to enable testing against pre-production environments. + to enable testing against pre-production environments. The default scope is the OBS + resource ``/.default`` scope for app-only S2S tokens. Returns: list[str]: The authentication scope for the current environment. diff --git a/tests/observability/core/test_agent365.py b/tests/observability/core/test_agent365.py index 5ae441c5..ca91e13b 100644 --- a/tests/observability/core/test_agent365.py +++ b/tests/observability/core/test_agent365.py @@ -83,6 +83,32 @@ def test_configure_with_exporter_options_and_parameter_precedence(self, mock_is_ ) self.assertTrue(result, "configure() should return True with exporter_options") + def test_agent365_exporter_options_keeps_legacy_s2s_false_default(self): + """Deprecated use_s2s_endpoint option keeps legacy False default but routing ignores it.""" + self.assertFalse(Agent365ExporterOptions().use_s2s_endpoint) + + @patch("microsoft_agents_a365.observability.core.config.is_agent365_exporter_enabled") + def test_configure_uses_console_when_exporter_enabled_without_token_resolver( + self, mock_is_enabled + ): + """Existing console fallback remains when exporter is enabled without a resolver.""" + mock_is_enabled.return_value = True + + with patch( + "microsoft_agents_a365.observability.core.config._telemetry_manager._logger" + ) as mock_logger: + result = configure( + service_name="test-service", + service_namespace="test-namespace", + exporter_options=Agent365ExporterOptions(), + ) + + self.assertTrue(result, "configure() should retain console fallback without a resolver") + mock_logger.warning.assert_called_once_with( + "is_agent365_exporter_enabled() not enabled or token_resolver not set." + " Falling back to console exporter." + ) + @patch("microsoft_agents_a365.observability.core.config._Agent365Exporter") @patch("microsoft_agents_a365.observability.core.config._EnrichingBatchSpanProcessor") @patch("microsoft_agents_a365.observability.core.config.is_agent365_exporter_enabled") diff --git a/tests/observability/core/test_agent365_exporter.py b/tests/observability/core/test_agent365_exporter.py index f3d80ebe..60056109 100644 --- a/tests/observability/core/test_agent365_exporter.py +++ b/tests/observability/core/test_agent365_exporter.py @@ -1,6 +1,8 @@ # Copyright (c) Microsoft Corporation. # Licensed under the MIT License. +import asyncio +import inspect import json import os import unittest @@ -133,7 +135,8 @@ def test_export_success(self): self.assertIn(DEFAULT_ENDPOINT_URL, url) self.assertIn( - "/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces", url + "/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces", + url, ) self.assertEqual(headers["authorization"], "Bearer test_token_123") self.assertEqual(headers["content-type"], "application/json") @@ -255,9 +258,9 @@ def test_s2s_endpoint_path_when_enabled(self): self.assertEqual(headers["authorization"], "Bearer test_token_123") self.assertEqual(headers["content-type"], "application/json") - def test_default_endpoint_path_when_s2s_disabled(self): - """Test 5: Test that default endpoint path is used when use_s2s_endpoint is False.""" - # Arrange - Create exporter with S2S endpoint disabled (default behavior) + def test_s2s_endpoint_path_when_legacy_flag_disabled(self): + """Test 5: Test that S2S endpoint path is used when use_s2s_endpoint is False.""" + # Arrange - Create exporter with deprecated S2S flag disabled default_exporter = _Agent365Exporter( token_resolver=self.mock_token_resolver, cluster_category="test", use_s2s_endpoint=False ) @@ -273,21 +276,138 @@ def test_default_endpoint_path_when_s2s_disabled(self): self.assertEqual(result, SpanExportResult.SUCCESS) mock_post.assert_called_once() - # Verify the call arguments - should use default path with default endpoint + # Verify the call arguments - should still use S2S path with default endpoint args, kwargs = mock_post.call_args url, body, headers = args self.assertIn(DEFAULT_ENDPOINT_URL, url) self.assertIn( - "/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces", url - ) - self.assertNotIn( "/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces", url, ) + self.assertNotIn("/observability/tenants/", url) self.assertEqual(headers["authorization"], "Bearer test_token_123") self.assertEqual(headers["content-type"], "application/json") + def test_omitted_false_true_legacy_flag_all_route_to_s2s(self): + """Omitted/False/True use_s2s_endpoint values are ignored; export always uses S2S.""" + for use_s2s_endpoint in (None, False, True): + with self.subTest(use_s2s_endpoint=use_s2s_endpoint): + kwargs = {} + if use_s2s_endpoint is not None: + kwargs["use_s2s_endpoint"] = use_s2s_endpoint + exporter = _Agent365Exporter( + token_resolver=self.mock_token_resolver, + cluster_category="test", + **kwargs, + ) + spans = [self._create_mock_span("legacy_flag_span")] + + with patch.object(exporter, "_post_with_retries", return_value=True) as mock_post: + result = exporter.export(spans) + + self.assertEqual(result, SpanExportResult.SUCCESS) + url = mock_post.call_args[0][0] + self.assertIn("/observabilityService/tenants/", url) + self.assertNotIn("/observability/tenants/", url) + + def test_empty_token_fails_without_sending_request(self): + """Empty app-only tokens fail the export batch without an HTTP request.""" + for token in (None, "", " "): + with self.subTest(token=token): + resolver = Mock(return_value=token) + exporter = _Agent365Exporter(token_resolver=resolver, cluster_category="test") + spans = [self._create_mock_span("empty_token_span")] + + with patch.object(exporter, "_post_with_retries", return_value=True) as mock_post: + result = exporter.export(spans) + + self.assertEqual(result, SpanExportResult.FAILURE) + resolver.assert_called_once_with("test-agent-456", "test-tenant-123") + mock_post.assert_not_called() + + def test_resolver_failure_fails_without_sending_request(self): + """Resolver acquisition failures fail export without delegated fallback or HTTP.""" + resolver = Mock(side_effect=RuntimeError("app-only acquisition failed")) + exporter = _Agent365Exporter(token_resolver=resolver, cluster_category="test") + spans = [self._create_mock_span("resolver_failure_span")] + + with patch.object(exporter, "_post_with_retries", return_value=True) as mock_post: + result = exporter.export(spans) + + self.assertEqual(result, SpanExportResult.FAILURE) + resolver.assert_called_once_with("test-agent-456", "test-tenant-123") + mock_post.assert_not_called() + + def test_async_resolver_is_awaited_for_sync_export(self): + """Async token resolvers are awaited when export runs outside an active event loop.""" + + async def resolver(agent_id, tenant_id): + return f"async-token-{agent_id}-{tenant_id}" + + exporter = _Agent365Exporter(token_resolver=resolver, cluster_category="test") + spans = [self._create_mock_span("async_resolver_span")] + + with patch.object(exporter, "_post_with_retries", return_value=True) as mock_post: + result = exporter.export(spans) + + self.assertEqual(result, SpanExportResult.SUCCESS) + mock_post.assert_called_once() + + def test_awaitable_resolver_in_active_event_loop_fails_and_is_released(self): + """Inside a running loop, export fails and the resolver's coroutine, future or task is released.""" + + async def pending_token(): + await asyncio.sleep(3600) + return "never-used" + + async def run_in_loop(): + coroutine = pending_token() + future = asyncio.get_running_loop().create_future() + task = asyncio.ensure_future(pending_token()) + for awaitable in (coroutine, future, task): + exporter = _Agent365Exporter( + token_resolver=lambda _agent_id, _tenant_id, value=awaitable: value, + cluster_category="test", + ) + spans = [self._create_mock_span("active_loop_span")] + with patch.object(exporter, "_post_with_retries", return_value=True) as mock_post: + self.assertEqual(exporter.export(spans), SpanExportResult.FAILURE) + mock_post.assert_not_called() + + await asyncio.wait([task], timeout=1) + self.assertEqual(inspect.getcoroutinestate(coroutine), inspect.CORO_CLOSED) + self.assertTrue(future.cancelled()) + self.assertTrue(task.cancelled()) + + asyncio.run(run_in_loop()) + + def test_auth_not_found_errors_do_not_fallback_to_delegated_route(self): + """401/403/404 failures do not trigger a delegated route fallback.""" + for status_code in (401, 403, 404): + with self.subTest(status_code=status_code): + exporter = _Agent365Exporter( + token_resolver=self.mock_token_resolver, + cluster_category="test", + use_s2s_endpoint=False, + ) + spans = [self._create_mock_span("auth_failure_span")] + + with patch("requests.Session.post") as mock_post: + mock_response = Mock() + mock_response.status_code = status_code + mock_response.text = "auth failure" + mock_response.headers = {"x-ms-correlation-id": "corr"} + mock_post.return_value = mock_response + + result = exporter.export(spans) + + self.assertEqual(result, SpanExportResult.FAILURE) + mock_post.assert_called_once() + url = mock_post.call_args[0][0] + self.assertIn("/observabilityService/tenants/", url) + self.assertNotIn("/observability/tenants/", url) + @patch("microsoft_agents_a365.observability.core.exporters.agent365_exporter.logger") def test_export_logging(self, mock_logger): """Test that the exporter logs appropriate messages during export.""" @@ -329,11 +449,13 @@ def test_export_logging(self, mock_logger): unittest.mock.call.debug("Found 1 identity groups with 2 total spans to export"), # Should log endpoint being used at DEBUG (default endpoint) unittest.mock.call.debug( - f"Exporting 2 spans to endpoint: {DEFAULT_ENDPOINT_URL}/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1 " + f"Exporting 2 spans to endpoint: {DEFAULT_ENDPOINT_URL}/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1 " "(tenant: test-tenant-123, agent: test-agent-456)" ), # Should log token resolution success at DEBUG - unittest.mock.call.debug("Token resolved successfully for agent test-agent-456"), + unittest.mock.call.debug( + "App-only OBS token resolved successfully for agent test-agent-456" + ), # Should log HTTP success at DEBUG unittest.mock.call.debug( "HTTP 200 success on attempt 1. " @@ -400,7 +522,7 @@ def test_export_uses_domain_override_when_env_var_set(self): args, kwargs = mock_post.call_args url, body, headers = args - expected_url = f"https://{override_domain}/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" + expected_url = f"https://{override_domain}/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" self.assertEqual(url, expected_url) def test_export_uses_default_endpoint_when_no_override(self): @@ -429,7 +551,7 @@ def test_export_uses_default_endpoint_when_no_override(self): args, kwargs = mock_post.call_args url, body, headers = args - expected_url = f"{DEFAULT_ENDPOINT_URL}/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" + expected_url = f"{DEFAULT_ENDPOINT_URL}/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" self.assertEqual(url, expected_url) def test_export_ignores_empty_domain_override(self): @@ -483,7 +605,7 @@ def test_export_uses_valid_url_override_with_https(self): args, kwargs = mock_post.call_args url, body, headers = args - expected_url = "https://override.example.com/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" + expected_url = "https://override.example.com/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" self.assertEqual(url, expected_url) def test_export_uses_valid_url_override_with_http(self): @@ -511,7 +633,7 @@ def test_export_uses_valid_url_override_with_http(self): args, kwargs = mock_post.call_args url, body, headers = args - expected_url = "http://localhost:8080/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" + expected_url = "http://localhost:8080/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" self.assertEqual(url, expected_url) def test_export_uses_valid_domain_override_with_port(self): @@ -539,7 +661,7 @@ def test_export_uses_valid_domain_override_with_port(self): args, kwargs = mock_post.call_args url, body, headers = args - expected_url = "https://example.com:8080/observability/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" + expected_url = "https://example.com:8080/observabilityService/tenants/test-tenant-123/otlp/agents/test-agent-456/traces?api-version=1" self.assertEqual(url, expected_url) def test_export_ignores_invalid_domain_with_protocol(self): diff --git a/tests/observability/core/test_export_config_consistency.py b/tests/observability/core/test_export_config_consistency.py index 443aa5d0..0e057784 100644 --- a/tests/observability/core/test_export_config_consistency.py +++ b/tests/observability/core/test_export_config_consistency.py @@ -34,8 +34,7 @@ class TestExportConfigConsistency(unittest.TestCase): # ---- pinned production values ---- EXPECTED_ENDPOINT = "https://agent365.svc.cloud.microsoft" - EXPECTED_SCOPE = "api://9b975845-388f-4429-889e-eab1ef63949c/Agent365.Observability.OtelWrite" - EXPECTED_STANDARD_PATH = "/observability/tenants/{tid}/otlp/agents/{aid}/traces" + EXPECTED_SCOPE = "api://9b975845-388f-4429-889e-eab1ef63949c/.default" EXPECTED_S2S_PATH = "/observabilityService/tenants/{tid}/otlp/agents/{aid}/traces" # ---- snapshot assertions ---- @@ -58,17 +57,17 @@ def test_prod_observability_scope_value(self): "and build_export_url() path. All three must stay in sync.", ) - def test_export_url_standard_path_structure(self): - """Standard export URL must use the pinned path pattern.""" + def test_export_url_legacy_false_path_structure(self): + """Deprecated false flag must still use the pinned S2S path pattern.""" url = build_export_url(self.EXPECTED_ENDPOINT, "a1", "t1") expected = ( f"{self.EXPECTED_ENDPOINT}" - f"{self.EXPECTED_STANDARD_PATH.format(tid='t1', aid='a1')}?api-version=1" + f"{self.EXPECTED_S2S_PATH.format(tid='t1', aid='a1')}?api-version=1" ) self.assertEqual( url, expected, - "Standard export URL path changed — also review PROD_OBSERVABILITY_SCOPE " + "Export URL path changed — also review PROD_OBSERVABILITY_SCOPE " "and DEFAULT_ENDPOINT_URL. All three must stay in sync.", ) @@ -97,8 +96,8 @@ def test_scope_and_endpoint_are_coherent(self): self.assertEqual(len(scopes), 1) scope = scopes[0] - # Scope should reference the Agent365 Observability permission - self.assertIn("Agent365.Observability", scope) + # Scope should request app-only tokens for the OBS resource. + self.assertTrue(scope.endswith("/.default")) # Endpoint should be the agent365 service self.assertIn("agent365", DEFAULT_ENDPOINT_URL) diff --git a/tests/observability/hosting/token_cache_helpers/test_agent_token_cache.py b/tests/observability/hosting/token_cache_helpers/test_agent_token_cache.py index eb505ce3..6ca2b6fd 100644 --- a/tests/observability/hosting/token_cache_helpers/test_agent_token_cache.py +++ b/tests/observability/hosting/token_cache_helpers/test_agent_token_cache.py @@ -1,8 +1,12 @@ # Copyright (c) Microsoft Corporation. # Licensed under the MIT License. -"""Tests for AgenticTokenCache and AgenticTokenStruct.""" +"""Tests for AgenticTokenCache app-only OBS token handling.""" +import asyncio +import base64 +import json +import time from unittest.mock import AsyncMock, MagicMock import pytest @@ -14,6 +18,30 @@ ) +class StatusError(Exception): + """Exception with an HTTP-like status code.""" + + def __init__(self, status: int) -> None: + super().__init__(f"status {status}") + self.status = status + + +def _encode_segment(value: dict[str, object]) -> str: + encoded = base64.urlsafe_b64encode(json.dumps(value).encode("utf-8")).decode("utf-8") + return encoded.rstrip("=") + + +def make_jwt(exp_seconds_from_now: int, claims: dict[str, object] | None = None) -> str: + """Create an unsigned JWT-like token for cache expiry tests.""" + payload = { + "exp": int(time.time()) + exp_seconds_from_now, + "idtyp": "app", + } + if claims is not None: + payload.update(claims) + return f"{_encode_segment({'alg': 'none'})}.{_encode_segment(payload)}.sig" + + @pytest.fixture def mock_authorization(): """Create a mock Authorization instance.""" @@ -35,100 +63,379 @@ def token_cache(): @pytest.mark.asyncio -async def test_register_and_retrieve_token_success( - token_cache, mock_authorization, mock_turn_context -): - """Test complete flow: create struct, register, and retrieve token successfully.""" - agent_id = "agent-123" - tenant_id = "tenant-456" - expected_token = "mock-token-xyz" - scopes = ["https://example.com/.default"] +async def test_get_observability_token_returns_none_without_entry(token_cache): + """A cache miss returns None.""" + assert await token_cache.get_observability_token("agent", "tenant") is None + + +@pytest.mark.asyncio +async def test_refresh_passes_identity_and_default_app_only_scope(token_cache): + """Refresh calls the resolver with exporting identity and OBS /.default scope.""" + token = make_jwt(300, {"roles": []}) + calls = [] + + def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: + calls.append((agent_id, tenant_id, scopes)) + return token + + refreshed = await token_cache.refresh_observability_token("agent", "tenant", resolver) + + assert refreshed == token + assert await token_cache.get_observability_token("agent", "tenant") == token + assert calls == [("agent", "tenant", ["api://9b975845-388f-4429-889e-eab1ef63949c/.default"])] + + +@pytest.mark.asyncio +async def test_refresh_uses_custom_app_only_scopes(): + """Constructor-provided scopes are passed through to the resolver.""" + token_cache = AgenticTokenCache(observability_scopes=["api://custom-obs/.default"]) + token = make_jwt(300) + resolver = MagicMock(return_value=token) - mock_authorization.exchange_token.return_value = expected_token + await token_cache.refresh_observability_token("agent", "tenant", resolver) + resolver.assert_called_once_with("agent", "tenant", ["api://custom-obs/.default"]) + + +@pytest.mark.asyncio +async def test_single_scope_string_is_not_split_into_characters(): + """A bare scope string is treated as one scope.""" + token_cache = AgenticTokenCache(observability_scopes="api://custom-obs/.default") + resolver = MagicMock(return_value=make_jwt(300)) + + await token_cache.refresh_observability_token("agent", "tenant", resolver) + + resolver.assert_called_once_with("agent", "tenant", ["api://custom-obs/.default"]) + + +@pytest.mark.asyncio +async def test_refresh_accepts_async_resolver(token_cache): + """Async app-only resolvers are awaited by the hosting cache.""" + token = make_jwt(300) + + async def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: + await asyncio.sleep(0) + return f"{agent_id}:{tenant_id}:{scopes[0]}:{token}" + + refreshed = await token_cache.refresh_observability_token("agent", "tenant", resolver) + + assert refreshed == ( + f"agent:tenant:api://9b975845-388f-4429-889e-eab1ef63949c/.default:{token}" + ) + + +@pytest.mark.asyncio +async def test_legacy_register_shape_logs_once_and_does_not_exchange( + token_cache, mock_authorization, mock_turn_context, caplog +): + """Removed delegated registration shape is a safe no-op.""" token_struct = AgenticTokenStruct( authorization=mock_authorization, turn_context=mock_turn_context, ) - assert token_struct.auth_handler_name == "AGENTIC" - token_cache.register_observability( - agent_id=agent_id, - tenant_id=tenant_id, - token_generator=token_struct, - observability_scopes=scopes, + with caplog.at_level("ERROR"): + token_cache.register_observability("agent", "tenant", token_struct, ["scope"]) + token_cache.register_observability("agent", "tenant", token_struct, ["scope"]) + + mock_authorization.exchange_token.assert_not_called() + assert await token_cache.get_observability_token("agent", "tenant") is None + assert sum("Delegated OBS token" in record.message for record in caplog.records) == 1 + + +@pytest.mark.asyncio +async def test_refresh_rejects_non_callable_resolver(token_cache): + """refresh_observability_token requires a callable resolver.""" + with pytest.raises(TypeError, match="token_resolver must be callable"): + await token_cache.refresh_observability_token("agent", "tenant", "not-callable") + + +@pytest.mark.asyncio +@pytest.mark.parametrize("agent_id,tenant_id", [("", "tenant"), ("agent", " ")]) +async def test_refresh_rejects_empty_identity(token_cache, agent_id, tenant_id): + """Agent and tenant IDs are required for app-only refresh.""" + resolver = MagicMock(return_value=make_jwt(300)) + + with pytest.raises(ValueError, match="Agent and tenant IDs"): + await token_cache.refresh_observability_token(agent_id, tenant_id, resolver) + + resolver.assert_not_called() + + +@pytest.mark.asyncio +async def test_refresh_rejects_empty_scope_configuration(): + """An empty OBS scope configuration fails before token acquisition.""" + token_cache = AgenticTokenCache(observability_scopes=[]) + resolver = MagicMock(return_value=make_jwt(300)) + + with pytest.raises(ValueError, match="No valid scopes"): + await token_cache.refresh_observability_token("agent", "tenant", resolver) + + resolver.assert_not_called() + + +@pytest.mark.asyncio +@pytest.mark.parametrize("token", [None, "", " "]) +async def test_refresh_surfaces_empty_resolver_result(token_cache, token): + """Empty resolver results propagate as refresh failures.""" + resolver = MagicMock(return_value=token) + + with pytest.raises(RuntimeError, match="returned no token"): + await token_cache.refresh_observability_token("agent", "tenant", resolver) + + assert await token_cache.get_observability_token("agent", "tenant") is None + + +@pytest.mark.asyncio +async def test_refresh_propagates_permanent_resolver_failure(token_cache): + """Permanent acquisition errors clear stale cached tokens and propagate.""" + await token_cache.refresh_observability_token("agent", "tenant", lambda *_: make_jwt(300)) + error = RuntimeError("permission denied") + resolver = MagicMock(side_effect=error) + token_cache.invalidate_token("agent", "tenant") + + with pytest.raises(RuntimeError, match="permission denied"): + await token_cache.refresh_observability_token("agent", "tenant", resolver) + + assert await token_cache.get_observability_token("agent", "tenant") is None + + +@pytest.mark.asyncio +async def test_refresh_retries_transient_failure_then_caches(token_cache, monkeypatch): + """Transient resolver failures are retried before succeeding.""" + token = make_jwt(300) + resolver = MagicMock(side_effect=[StatusError(503), token]) + + async def no_sleep(delay: float) -> None: + return None + + monkeypatch.setattr(asyncio, "sleep", no_sleep) + + refreshed = await token_cache.refresh_observability_token("agent", "tenant", resolver) + + assert refreshed == token + assert resolver.call_count == 2 + assert await token_cache.get_observability_token("agent", "tenant") == token + + +@pytest.mark.asyncio +@pytest.mark.parametrize("error", [TimeoutError(), ConnectionResetError()]) +async def test_refresh_retries_standard_transient_exceptions(token_cache, monkeypatch, error): + """Bare timeout and connection-reset errors are retried even without a message.""" + token = make_jwt(300) + resolver = MagicMock(side_effect=[error, token]) + + async def no_sleep(delay: float) -> None: + return None + + monkeypatch.setattr(asyncio, "sleep", no_sleep) + + assert await token_cache.refresh_observability_token("agent", "tenant", resolver) == token + assert resolver.call_count == 2 + + +@pytest.mark.asyncio +async def test_refresh_deduplicates_concurrent_same_identity_acquisition(token_cache): + """Concurrent refreshes for the same key share the first acquired token.""" + token = make_jwt(300) + call_count = 0 + + async def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: + nonlocal call_count + call_count += 1 + await asyncio.sleep(0) + return token + + results = await asyncio.gather( + *(token_cache.refresh_observability_token("agent", "tenant", resolver) for _ in range(8)) + ) + + assert results == [token] * 8 + assert call_count == 1 + + +@pytest.mark.asyncio +async def test_refresh_reuses_cached_token_until_expiry_skew(token_cache): + """A non-expired token is reused; near-expiry tokens are not returned.""" + token = make_jwt(120) + resolver = MagicMock(return_value=token) + + await token_cache.refresh_observability_token("agent", "tenant", resolver) + await token_cache.refresh_observability_token("agent", "tenant", resolver) + + assert resolver.call_count == 1 + assert await token_cache.get_observability_token("agent", "tenant") == token + + token_cache.invalidate_token("agent", "tenant") + near_expiry = make_jwt(30) + await token_cache.refresh_observability_token("agent", "tenant", lambda *_: near_expiry) + assert await token_cache.get_observability_token("agent", "tenant") is None + + +@pytest.mark.asyncio +async def test_opaque_token_uses_fresh_fallback_ttl(token_cache): + """Opaque tokens get a fallback TTL after replacing expired JWT metadata.""" + await token_cache.refresh_observability_token("agent", "tenant", lambda *_: make_jwt(120)) + token_cache.invalidate_token("agent", "tenant") + await token_cache.refresh_observability_token("agent", "tenant", lambda *_: "opaque-token") + + assert await token_cache.get_observability_token("agent", "tenant") == "opaque-token" + entry = token_cache._map[AgenticTokenCache._make_key("agent", "tenant")] + entry.acquired_on_ms = (time.time() * 1000) - token_cache._default_max_token_age_ms - 1 + assert await token_cache.get_observability_token("agent", "tenant") is None + + +@pytest.mark.asyncio +async def test_tokens_are_isolated_by_agent_and_tenant(token_cache): + """Agent and tenant are both part of the cache key.""" + + def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: + return f"{agent_id}:{tenant_id}" + + await token_cache.refresh_observability_token("agent-one", "tenant-a", resolver) + await token_cache.refresh_observability_token("agent-two", "tenant-a", resolver) + await token_cache.refresh_observability_token("agent-one", "tenant-b", resolver) + + assert ( + await token_cache.get_observability_token("agent-one", "tenant-a") == "agent-one:tenant-a" + ) + assert ( + await token_cache.get_observability_token("agent-two", "tenant-a") == "agent-two:tenant-a" + ) + assert ( + await token_cache.get_observability_token("agent-one", "tenant-b") == "agent-one:tenant-b" ) - token = await token_cache.get_observability_token(agent_id, tenant_id) - assert token == expected_token +@pytest.mark.asyncio +async def test_invalidate_one_then_all(token_cache): + """Token invalidation is scoped by key or all keys.""" + token = make_jwt(300) + await token_cache.refresh_observability_token("one", "tenant", lambda *_: token) + await token_cache.refresh_observability_token("two", "tenant", lambda *_: token) -@pytest.mark.parametrize( - "agent_id,tenant_id,token_generator,error_type,error_match", - [ - ("", "tenant-456", "valid", ValueError, "agent_id cannot be None or whitespace"), - ("agent-123", "tenant-456", None, TypeError, "token_generator cannot be None"), - ], -) -def test_register_observability_validation( - token_cache, - mock_authorization, - mock_turn_context, - agent_id, - tenant_id, - token_generator, - error_type, - error_match, -): - """Test that registration validates inputs and raises appropriate errors.""" - struct = None - if token_generator == "valid": - struct = AgenticTokenStruct( - authorization=mock_authorization, - turn_context=mock_turn_context, - ) - - with pytest.raises(error_type, match=error_match): - token_cache.register_observability( - agent_id=agent_id, - tenant_id=tenant_id, - token_generator=struct, - observability_scopes=["scope"], - ) - - -def test_thread_safety(token_cache, mock_authorization, mock_turn_context): - """Test that cache is thread-safe with concurrent registrations.""" - import threading - - agent_id = "agent-123" - tenant_id = "tenant-456" - results = [] - - def register_token(scope_suffix): - try: - struct = AgenticTokenStruct( - authorization=mock_authorization, - turn_context=mock_turn_context, - ) - token_cache.register_observability( - agent_id=agent_id, - tenant_id=tenant_id, - token_generator=struct, - observability_scopes=[f"scope-{scope_suffix}"], - ) - results.append(scope_suffix) - except Exception as e: - results.append(f"error: {e}") - - # Create 10 concurrent registrations - threads = [threading.Thread(target=register_token, args=(i,)) for i in range(10)] - for thread in threads: - thread.start() - for thread in threads: - thread.join() - - # All registrations should succeed - assert len(results) == 10 - # Only one entry should exist (idempotent) - assert f"{agent_id}:{tenant_id}" in token_cache._map + token_cache.invalidate_token("one", "tenant") + assert await token_cache.get_observability_token("one", "tenant") is None + assert await token_cache.get_observability_token("two", "tenant") == token + + token_cache.invalidate_all() + assert await token_cache.get_observability_token("two", "tenant") is None + + +@pytest.mark.asyncio +async def test_invalidate_token_is_not_undone_by_in_flight_refresh(token_cache): + """A refresh that completes after invalidate_token() does not repopulate the cache.""" + release = asyncio.Event() + + async def slow_resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: + await release.wait() + return "token-acquired-before-invalidation" + + in_flight = asyncio.create_task( + token_cache.refresh_observability_token("agent", "tenant", slow_resolver) + ) + await asyncio.sleep(0) + token_cache.invalidate_token("agent", "tenant") + release.set() + + assert await in_flight == "token-acquired-before-invalidation" + assert await token_cache.get_observability_token("agent", "tenant") is None + + fresh = MagicMock(return_value="token-acquired-after-invalidation") + refreshed = await token_cache.refresh_observability_token("agent", "tenant", fresh) + assert refreshed == "token-acquired-after-invalidation" + fresh.assert_called_once() + + +@pytest.mark.asyncio +async def test_refresh_returns_the_cached_token_it_checked(token_cache): + """A concurrent invalidation between the usability check and the return never yields None.""" + token = make_jwt(300) + await token_cache.refresh_observability_token("agent", "tenant", lambda *_: token) + check_expiry = token_cache._is_expired + + def invalidate_during_check(entry: AgenticTokenCache._Entry) -> bool: + expired = check_expiry(entry) + token_cache.invalidate_token("agent", "tenant") + return expired + + token_cache._is_expired = invalidate_during_check + resolver = MagicMock(return_value="unused") + + assert await token_cache.refresh_observability_token("agent", "tenant", resolver) == token + resolver.assert_not_called() + + +@pytest.mark.asyncio +async def test_cache_evicts_oldest_entry_when_capacity_is_reached(token_cache): + """Cache size is bounded and evicts the oldest key.""" + token_cache._max_cache_size = 2 + token = make_jwt(300) + await token_cache.refresh_observability_token("one", "tenant", lambda *_: token) + await token_cache.refresh_observability_token("two", "tenant", lambda *_: token) + await token_cache.refresh_observability_token("three", "tenant", lambda *_: token) + + assert await token_cache.get_observability_token("one", "tenant") is None + assert await token_cache.get_observability_token("two", "tenant") == token + assert await token_cache.get_observability_token("three", "tenant") == token + + +@pytest.mark.asyncio +async def test_separator_bearing_ids_do_not_share_a_cache_entry(token_cache): + """IDs that contain a separator character never alias another identity.""" + first = MagicMock(return_value="token-for-first") + second = MagicMock(return_value="token-for-second") + + assert await token_cache.refresh_observability_token("a:b", "c", first) == "token-for-first" + assert await token_cache.refresh_observability_token("a", "b:c", second) == "token-for-second" + + second.assert_called_once() + assert await token_cache.get_observability_token("a:b", "c") == "token-for-first" + assert await token_cache.get_observability_token("a", "b:c") == "token-for-second" + + +@pytest.mark.asyncio +async def test_identity_churn_keeps_every_per_identity_registry_bounded(token_cache): + """Refresh locks live and die with cache entries, so identity churn stays bounded.""" + token_cache._max_cache_size = 2 + token = make_jwt(300) + for agent_id in ("one", "two", "three", "four"): + await token_cache.refresh_observability_token(agent_id, "tenant", lambda *_: token) + + registries = { + name: len(value) for name, value in vars(token_cache).items() if isinstance(value, dict) + } + assert registries == {"_map": 2} + + token_cache.invalidate_all() + registries = { + name: len(value) for name, value in vars(token_cache).items() if isinstance(value, dict) + } + assert registries == {"_map": 0} + + +@pytest.mark.asyncio +async def test_eviction_keeps_in_flight_entry_then_trims_overflow(token_cache): + """Eviction skips a refresh that is still running; the next insertion restores the bound.""" + token_cache._max_cache_size = 1 + release = asyncio.Event() + + async def slow_resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: + await release.wait() + return "token-for-slow" + + in_flight = asyncio.create_task( + token_cache.refresh_observability_token("slow", "tenant", slow_resolver) + ) + await asyncio.sleep(0) + fast = await token_cache.refresh_observability_token( + "fast", "tenant", lambda *_: "token-for-fast" + ) + + release.set() + assert fast == "token-for-fast" + assert await in_flight == "token-for-slow" + assert await token_cache.get_observability_token("slow", "tenant") == "token-for-slow" + + await token_cache.refresh_observability_token("next", "tenant", lambda *_: "token-for-next") + assert len(token_cache._map) == 1 diff --git a/tests/runtime/test_environment_utils.py b/tests/runtime/test_environment_utils.py index a79a1b43..4fdde30f 100644 --- a/tests/runtime/test_environment_utils.py +++ b/tests/runtime/test_environment_utils.py @@ -12,9 +12,10 @@ def test_get_observability_authentication_scope(): - """Test get_observability_authentication_scope returns production scope.""" + """Test get_observability_authentication_scope returns app-only production scope.""" result = get_observability_authentication_scope() assert result == [PROD_OBSERVABILITY_SCOPE] + assert result == ["api://9b975845-388f-4429-889e-eab1ef63949c/.default"] def test_get_observability_authentication_scope_with_override(monkeypatch): diff --git a/uv.lock b/uv.lock index 6db581d4..f9f13b87 100644 --- a/uv.lock +++ b/uv.lock @@ -3200,6 +3200,7 @@ name = "microsoft-agents-a365-observability-hosting" source = { editable = "libraries/microsoft-agents-a365-observability-hosting" } dependencies = [ { name = "microsoft-agents-a365-observability-core" }, + { name = "microsoft-agents-a365-runtime" }, { name = "microsoft-agents-hosting-core" }, { name = "opentelemetry-api" }, ] @@ -3207,6 +3208,7 @@ dependencies = [ [package.metadata] requires-dist = [ { name = "microsoft-agents-a365-observability-core", editable = "libraries/microsoft-agents-a365-observability-core" }, + { name = "microsoft-agents-a365-runtime", editable = "libraries/microsoft-agents-a365-runtime" }, { name = "microsoft-agents-hosting-core" }, { name = "opentelemetry-api" }, ] diff --git a/versioning/TARGET-VERSION b/versioning/TARGET-VERSION index afaf360d..359a5b95 100644 --- a/versioning/TARGET-VERSION +++ b/versioning/TARGET-VERSION @@ -1 +1 @@ -1.0.0 \ No newline at end of file +2.0.0 \ No newline at end of file