From b89d3502d4d8318783b8deee07ba9b4980554283 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:41:46 +0100 Subject: [PATCH 01/13] fix(observability): use S2S app-only export Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../observability/core/config.py | 28 ++-- .../observability/core/exporters/__init__.py | 4 +- .../core/exporters/agent365_exporter.py | 52 ++++++-- .../exporters/agent365_exporter_options.py | 13 +- .../observability/core/exporters/utils.py | 9 +- .../runtime/environment_utils.py | 13 +- tests/observability/core/test_agent365.py | 17 +++ .../core/test_agent365_exporter.py | 122 +++++++++++++++--- .../core/test_export_config_consistency.py | 15 +-- tests/runtime/test_environment_utils.py | 3 +- versioning/TARGET-VERSION | 2 +- 11 files changed, 213 insertions(+), 65 deletions(-) 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..0a9d2cc7 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, @@ -175,7 +174,13 @@ def _configure_internal( endpoint=exporter_options.endpoint, ) - elif is_agent365_exporter_enabled() and exporter_options.token_resolver is not None: + elif is_agent365_exporter_enabled(): + if exporter_options.token_resolver is None: + raise ValueError( + "Agent365Exporter requires an app-only OBS token_resolver when " + "ENABLE_A365_OBSERVABILITY_EXPORTER is enabled. Delegated/context " + "tokens are not used for OBS export." + ) exporter = _Agent365Exporter( token_resolver=exporter_options.token_resolver, cluster_category=exporter_options.cluster_category, @@ -186,8 +191,7 @@ def _configure_internal( else: exporter = ConsoleSpanExporter() self._logger.warning( - "is_agent365_exporter_enabled() not enabled or token_resolver not set." - " Falling back to console exporter." + "is_agent365_exporter_enabled() not enabled. Falling back to console exporter." ) # Add span processors @@ -270,7 +274,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 +286,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..8ad3ead4 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,22 @@ 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() + 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..f96ddb14 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,8 +17,8 @@ class Agent365ExporterOptions: def __init__( self, cluster_category: str = "prod", - token_resolver: Optional[Callable[[str, str], Awaitable[Optional[str]]]] = None, - use_s2s_endpoint: bool = False, + token_resolver: Optional[TokenResolver] = None, + use_s2s_endpoint: bool = True, max_queue_size: int = 2048, scheduled_delay_ms: int = 5000, exporter_timeout_ms: int = 30000, @@ -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-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..cd42414d 100644 --- a/tests/observability/core/test_agent365.py +++ b/tests/observability/core/test_agent365.py @@ -83,6 +83,23 @@ 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_defaults_to_legacy_s2s_true(self): + """Deprecated use_s2s_endpoint option defaults to True but routing ignores the value.""" + self.assertTrue(Agent365ExporterOptions().use_s2s_endpoint) + + @patch("microsoft_agents_a365.observability.core.config.is_agent365_exporter_enabled") + def test_configure_fails_when_exporter_enabled_without_token_resolver(self, mock_is_enabled): + """Enabled Agent 365 exporter requires an explicit app-only OBS resolver.""" + mock_is_enabled.return_value = True + + result = configure( + service_name="test-service", + service_namespace="test-namespace", + exporter_options=Agent365ExporterOptions(), + ) + + self.assertFalse(result, "configure() should fail without an app-only token resolver") + @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..f9d72228 100644 --- a/tests/observability/core/test_agent365_exporter.py +++ b/tests/observability/core/test_agent365_exporter.py @@ -133,7 +133,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 +256,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 +274,110 @@ 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_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 +419,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 +492,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 +521,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 +575,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 +603,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 +631,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/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/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 From ff8cb02eb2fff46f64ec453ff4c5a88d752e9899 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:41:53 +0100 Subject: [PATCH 02/13] fix(hosting): require app-only OBS token resolver Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../observability/hosting/__init__.py | 4 + .../hosting/token_cache_helpers/__init__.py | 4 +- .../token_cache_helpers/agent_token_cache.py | 319 ++++++++++++---- .../test_agent_token_cache.py | 354 ++++++++++++++---- 4 files changed, 526 insertions(+), 155 deletions(-) diff --git a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py index 6b41ec88..bca2ba08 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py +++ b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py @@ -14,6 +14,7 @@ A365_PARENT_TRACEPARENT_KEY, OutputLoggingMiddleware, ) +from .token_cache_helpers import AgenticTokenCache, AgenticTokenStruct, ObservabilityTokenResolver __all__ = [ "BaggageMiddleware", @@ -21,4 +22,7 @@ "A365_PARENT_TRACEPARENT_KEY", "ObservabilityHostingManager", "ObservabilityHostingOptions", + "AgenticTokenCache", + "AgenticTokenStruct", + "ObservabilityTokenResolver", ] 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..25dd631b 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,85 @@ # 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 +import time +from collections.abc import Awaitable, Callable, Sequence from dataclasses import dataclass +from inspect import isawaitable from threading import Lock +from typing import cast 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, Sequence[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 - 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] = {} + self._key_locks: dict[str, asyncio.Lock] = {} self._lock = Lock() + self._observability_scopes = ( + None if observability_scopes is None else tuple(observability_scopes) + ) + self._removed_overload_logged = False + + @staticmethod + def make_key(agent_id: str, tenant_id: str) -> str: + """Create a cache key for an agent and tenant.""" + return f"{agent_id}:{tenant_id}" def register_observability( self, @@ -59,79 +88,229 @@ 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_overload_once() + + async def RefreshObservabilityToken( + self, + agent_id: str, + tenant_id: str, + token_resolver: ObservabilityTokenResolver | object, + *removed_overload_args: object, + ) -> str | None: + """Compatibility alias for :meth:`refresh_observability_token`.""" + return await self.refresh_observability_token( + agent_id, + tenant_id, + token_resolver, + *removed_overload_args, + ) + + async def refresh_observability_token( + self, + agent_id: str, + tenant_id: str, + token_resolver: ObservabilityTokenResolver | object, + *removed_overload_args: object, + ) -> str | None: + """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)``. + removed_overload_args: Present only for the removed delegated + TurnContext/Authorization overload; ignored after a one-time log. + + Returns: + The cached app-only OBS token, or ``None`` for the removed overload. Raises: - ValueError: If agent_id or tenant_id is empty or None. - TypeError: If token_generator is None. + 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 removed_overload_args or not callable(token_resolver): + self._log_removed_overload_once() + return None - 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) + lock = self._get_key_lock(key) + async with lock: + entry = self._get_or_create_entry(key) + if entry.token is not None and not self._is_expired(entry): + return entry.token - key = f"{agent_id}:{tenant_id}" + resolver = cast(ObservabilityTokenResolver, token_resolver) + return await self._acquire_token(agent_id, tenant_id, entry, resolver) - # First registration wins; subsequent calls ignored (idempotent) + def get_observability_token(self, agent_id: str, tenant_id: str) -> str | None: + """Get a non-expired cached app-only OBS token.""" + key = self.make_key(agent_id, tenant_id) 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") - - async def get_observability_token(self, agent_id: str, tenant_id: str) -> str | None: - """ - Get the observability token for the specified agent and tenant. + 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. - """ - key = f"{agent_id}:{tenant_id}" + def invalidate_token(self, agent_id: str, tenant_id: str) -> None: + """Invalidate one cached token.""" + key = self.make_key(agent_id, tenant_id) + with self._lock: + entry = self._map.get(key) + 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: 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 - if entry is None: - logger.debug(f"Cache miss for {key}") - return None + scopes = self._get_effective_scopes() + if not scopes: + raise ValueError("[AgenticTokenCache] No valid scopes") - logger.debug(f"Cache hit for {key}, exchanging token") + if len(self._map) >= self._max_cache_size: + oldest_key = next(iter(self._map), None) + if oldest_key is not None: + del self._map[oldest_key] - 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, - ) + entry = AgenticTokenCache._Entry(scopes=scopes) + self._map[key] = entry + return entry - 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__}") + def _get_key_lock(self, key: str) -> asyncio.Lock: + with self._lock: + lock = self._key_locks.get(key) + if lock is None: + lock = asyncio.Lock() + self._key_locks[key] = lock + return lock + + 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()) + + 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, + ) + try: + result = resolver(agent_id, tenant_id, list(entry.scopes)) + token = await cast(Awaitable[str | None], 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: + 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_overload_once(self) -> None: + if self._removed_overload_logged: + return + self._removed_overload_logged = True + logger.error( + "[AgenticTokenCache] Delegated OBS token registration/refresh 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/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..9277b9da 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.""" @@ -34,101 +62,261 @@ def token_cache(): return AgenticTokenCache() +def test_get_observability_token_returns_none_without_entry(token_cache): + """A cache miss returns None.""" + assert 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 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) + + 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_register_and_retrieve_token_success( - token_cache, mock_authorization, mock_turn_context +async def test_legacy_refresh_shape_logs_once_and_does_not_exchange( + token_cache, mock_authorization, mock_turn_context, caplog ): - """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"] + """Removed TurnContext/Authorization refresh shape is a safe no-op.""" + with caplog.at_level("ERROR"): + result1 = await token_cache.RefreshObservabilityToken( + "agent", + "tenant", + mock_turn_context, + mock_authorization, + ["api://old-scope/.default"], + "agentic", + ) + result2 = await token_cache.RefreshObservabilityToken( + "agent", + "tenant", + mock_turn_context, + mock_authorization, + ) + + assert result1 is None + assert result2 is None + mock_authorization.exchange_token.assert_not_called() + assert token_cache.get_observability_token("agent", "tenant") is None + assert sum("Delegated OBS token" in record.message for record in caplog.records) == 1 - mock_authorization.exchange_token.return_value = expected_token +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 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 +@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 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 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 token_cache.get_observability_token("agent", "tenant") == token + + +@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)) ) - token = await token_cache.get_observability_token(agent_id, tenant_id) - assert token == expected_token + assert results == [token] * 8 + assert call_count == 1 -@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, - ) +@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) - 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"], - ) + await token_cache.refresh_observability_token("agent", "tenant", resolver) + await token_cache.refresh_observability_token("agent", "tenant", resolver) + + assert resolver.call_count == 1 + assert 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 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 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 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 token_cache.get_observability_token("agent-one", "tenant-a") == "agent-one:tenant-a" + assert token_cache.get_observability_token("agent-two", "tenant-a") == "agent-two:tenant-a" + assert token_cache.get_observability_token("agent-one", "tenant-b") == "agent-one:tenant-b" + + +@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) + + token_cache.invalidate_token("one", "tenant") + assert token_cache.get_observability_token("one", "tenant") is None + assert token_cache.get_observability_token("two", "tenant") == token + + token_cache.invalidate_all() + assert token_cache.get_observability_token("two", "tenant") is None + + +@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) -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 + assert token_cache.get_observability_token("one", "tenant") is None + assert token_cache.get_observability_token("two", "tenant") == token + assert token_cache.get_observability_token("three", "tenant") == token From 908bd9b9b1baee731be46effd94a9c6da148d0ef Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:42:01 +0100 Subject: [PATCH 03/13] docs(observability): document S2S-only migration Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- docs/design.md | 8 ++-- ...integrating-with-existing-opentelemetry.md | 13 ++++-- .../CHANGELOG.md | 27 ++++++++++++ .../README.md | 29 ++++++++++++- .../docs/design.md | 25 ++++++++--- .../CHANGELOG.md | 22 ++++++++++ .../README.md | 42 +++++++++++++++++++ .../CHANGELOG.md | 14 +++++++ .../docs/design.md | 6 +++ 9 files changed, 173 insertions(+), 13 deletions(-) create mode 100644 libraries/microsoft-agents-a365-runtime/CHANGELOG.md 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..41615e53 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, configuration fails. 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..d30fe184 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 authentication. Missing + resolvers fail configuration when the Agent 365 exporter is enabled; 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..6d861ff7 100644 --- a/libraries/microsoft-agents-a365-observability-core/README.md +++ b/libraries/microsoft-agents-a365-observability-core/README.md @@ -17,6 +17,34 @@ 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, configuration fails instead of silently using another token. + +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 +61,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-hosting/CHANGELOG.md b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md index 456a642a..b44208ce 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md +++ b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md @@ -4,6 +4,28 @@ 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)` (and the + compatibility alias `RefreshObservabilityToken`) replaces delegated + TurnContext/Authorization token exchange for OBS 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. +- **Delegated OBS registration/refresh shapes are no-ops** — + `register_observability(...)` and untyped refresh calls with TurnContext / + Authorization arguments log one error and return no token. They never call + `Authorization.exchange_token` and never acquire a delegated OBS token. + +### Migration + +- Call `refresh_observability_token(agent_id, tenant_id, app_only_token_resolver)` + from the exporter `token_resolver`, then return + `get_observability_token(agent_id, tenant_id)`. +- 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..f9fcdf49 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/README.md +++ b/libraries/microsoft-agents-a365-observability-hosting/README.md @@ -7,3 +7,45 @@ 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. `RefreshObservabilityToken(...)` is also available +as a compatibility alias. + +```python +from microsoft_agents_a365.observability.hosting import AgenticTokenCache + +cache = AgenticTokenCache() + + +async def refresh_app_only_obs(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 resolver(agent_id: str, tenant_id: str) -> str | None: + await cache.refresh_observability_token(agent_id, tenant_id, refresh_app_only_obs) + return cache.get_observability_token(agent_id, tenant_id) +``` + +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. + +The previous delegated registration/refresh shapes using `TurnContext` and +`Authorization.exchange_token` are removed for OBS export. They are accepted only +as no-op compatibility shapes: the cache logs once, returns `None`, 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-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 From ddc6b7005969fb2bfcb906bf623ee91219c36abc Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 14:49:05 +0100 Subject: [PATCH 04/13] fix(observability): trim Python S2S port compatibility Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- ...integrating-with-existing-opentelemetry.md | 2 +- .../CHANGELOG.md | 12 ++--- .../README.md | 3 +- .../observability/core/config.py | 11 ++-- .../exporters/agent365_exporter_options.py | 2 +- .../CHANGELOG.md | 14 +++--- .../README.md | 22 +++++--- .../observability/hosting/__init__.py | 4 -- .../token_cache_helpers/agent_token_cache.py | 50 ++++++------------- tests/observability/core/test_agent365.py | 31 ++++++++---- .../test_agent_token_cache.py | 35 +++---------- 11 files changed, 76 insertions(+), 110 deletions(-) diff --git a/docs/integrating-with-existing-opentelemetry.md b/docs/integrating-with-existing-opentelemetry.md index 41615e53..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 an app-only OBS `token_resolver` is provided (otherwise `configure()` falls back to `ConsoleSpanExporter`). If the Agent 365 exporter is enabled without a resolver, configuration fails. +> **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. diff --git a/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md b/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md index d30fe184..1159082d 100644 --- a/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md +++ b/libraries/microsoft-agents-a365-observability-core/CHANGELOG.md @@ -11,12 +11,12 @@ All notable changes to this package will be documented in this file. 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 authentication. Missing - resolvers fail configuration when the Agent 365 exporter is enabled; 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. + 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 diff --git a/libraries/microsoft-agents-a365-observability-core/README.md b/libraries/microsoft-agents-a365-observability-core/README.md index 6d861ff7..51a0300c 100644 --- a/libraries/microsoft-agents-a365-observability-core/README.md +++ b/libraries/microsoft-agents-a365-observability-core/README.md @@ -35,7 +35,8 @@ 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, configuration fails instead of silently using another token. +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 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 0a9d2cc7..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 @@ -174,13 +174,7 @@ def _configure_internal( endpoint=exporter_options.endpoint, ) - elif is_agent365_exporter_enabled(): - if exporter_options.token_resolver is None: - raise ValueError( - "Agent365Exporter requires an app-only OBS token_resolver when " - "ENABLE_A365_OBSERVABILITY_EXPORTER is enabled. Delegated/context " - "tokens are not used for OBS export." - ) + elif is_agent365_exporter_enabled() and exporter_options.token_resolver is not None: exporter = _Agent365Exporter( token_resolver=exporter_options.token_resolver, cluster_category=exporter_options.cluster_category, @@ -191,7 +185,8 @@ def _configure_internal( else: exporter = ConsoleSpanExporter() self._logger.warning( - "is_agent365_exporter_enabled() not enabled. Falling back to console exporter." + "is_agent365_exporter_enabled() not enabled or token_resolver not set." + " Falling back to console exporter." ) # Add span processors 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 f96ddb14..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 @@ -18,7 +18,7 @@ def __init__( self, cluster_category: str = "prod", token_resolver: Optional[TokenResolver] = None, - use_s2s_endpoint: bool = True, + use_s2s_endpoint: bool = False, max_queue_size: int = 2048, scheduled_delay_ms: int = 5000, exporter_timeout_ms: int = 30000, diff --git a/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md index b44208ce..069225a0 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md +++ b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md @@ -7,16 +7,14 @@ All notable changes to this package will be documented in this file. ### Breaking Changes - **Hosting OBS token cache requires an app-only resolver** — - `refresh_observability_token(agent_id, tenant_id, token_resolver)` (and the - compatibility alias `RefreshObservabilityToken`) replaces delegated - TurnContext/Authorization token exchange for OBS export. The resolver receives - the configured OBS scopes and must acquire a token for the exporting agent + `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. -- **Delegated OBS registration/refresh shapes are no-ops** — - `register_observability(...)` and untyped refresh calls with TurnContext / - Authorization arguments log one error and return no token. They never call - `Authorization.exchange_token` and never acquire a delegated OBS token. +- **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 diff --git a/libraries/microsoft-agents-a365-observability-hosting/README.md b/libraries/microsoft-agents-a365-observability-hosting/README.md index f9fcdf49..15988ebb 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/README.md +++ b/libraries/microsoft-agents-a365-observability-hosting/README.md @@ -13,11 +13,10 @@ pip install microsoft-agents-a365-observability-hosting 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. `RefreshObservabilityToken(...)` is also available -as a compatibility alias. +the exporting agent identity. ```python -from microsoft_agents_a365.observability.hosting import AgenticTokenCache +from microsoft_agents_a365.observability.hosting.token_cache_helpers import AgenticTokenCache cache = AgenticTokenCache() @@ -39,10 +38,19 @@ 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. -The previous delegated registration/refresh shapes using `TurnContext` and -`Authorization.exchange_token` are removed for OBS export. They are accepted only -as no-op compatibility shapes: the cache logs once, returns `None`, and never -calls `exchange_token`. +If you use an async exporter resolver, `_Agent365Exporter` runs it with +`asyncio.run` on the BatchSpanProcessor worker thread. Create async clients +inside that resolver; do not reuse `aiohttp` or `azure.identity.aio` clients +bound to your app's event loop. Also, `force_flush()` and `shutdown()` called +from inside a running event loop cannot await an async resolver, so that export +fails with a logged error; call them with `await asyncio.to_thread(...)`. A +synchronous, thread-safe cached resolver avoids both constraints and is the +pattern used by the Agent365-Samples Python samples. + +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 diff --git a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py index bca2ba08..6b41ec88 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py +++ b/libraries/microsoft-agents-a365-observability-hosting/microsoft_agents_a365/observability/hosting/__init__.py @@ -14,7 +14,6 @@ A365_PARENT_TRACEPARENT_KEY, OutputLoggingMiddleware, ) -from .token_cache_helpers import AgenticTokenCache, AgenticTokenStruct, ObservabilityTokenResolver __all__ = [ "BaggageMiddleware", @@ -22,7 +21,4 @@ "A365_PARENT_TRACEPARENT_KEY", "ObservabilityHostingManager", "ObservabilityHostingOptions", - "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 25dd631b..302559b4 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 @@ -14,7 +14,6 @@ from dataclasses import dataclass from inspect import isawaitable from threading import Lock -from typing import cast from microsoft_agents.hosting.core.app.oauth.authorization import Authorization from microsoft_agents.hosting.core.turn_context import TurnContext @@ -74,7 +73,7 @@ def __init__(self, observability_scopes: Sequence[str] | None = None) -> None: self._observability_scopes = ( None if observability_scopes is None else tuple(observability_scopes) ) - self._removed_overload_logged = False + self._removed_registration_logged = False @staticmethod def make_key(agent_id: str, tenant_id: str) -> str: @@ -89,30 +88,14 @@ def register_observability( observability_scopes: list[str], ) -> None: """Deprecated no-op for the removed delegated OBS registration flow.""" - self._log_removed_overload_once() - - async def RefreshObservabilityToken( - self, - agent_id: str, - tenant_id: str, - token_resolver: ObservabilityTokenResolver | object, - *removed_overload_args: object, - ) -> str | None: - """Compatibility alias for :meth:`refresh_observability_token`.""" - return await self.refresh_observability_token( - agent_id, - tenant_id, - token_resolver, - *removed_overload_args, - ) + self._log_removed_registration_once() async def refresh_observability_token( self, agent_id: str, tenant_id: str, - token_resolver: ObservabilityTokenResolver | object, - *removed_overload_args: object, - ) -> str | None: + token_resolver: ObservabilityTokenResolver, + ) -> str: """Refresh an app-only OBS token for the exporting agent identity. Args: @@ -120,19 +103,17 @@ async def refresh_observability_token( tenant_id: The exporting tenant identifier. token_resolver: App-only token resolver receiving ``(agent_id, tenant_id, scopes)``. - removed_overload_args: Present only for the removed delegated - TurnContext/Authorization overload; ignored after a one-time log. Returns: - The cached app-only OBS token, or ``None`` for the removed overload. + The cached app-only OBS token. Raises: + 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 removed_overload_args or not callable(token_resolver): - self._log_removed_overload_once() - return None + if not callable(token_resolver): + raise TypeError("token_resolver must be callable") 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") @@ -144,8 +125,7 @@ async def refresh_observability_token( if entry.token is not None and not self._is_expired(entry): return entry.token - resolver = cast(ObservabilityTokenResolver, token_resolver) - return await self._acquire_token(agent_id, tenant_id, entry, resolver) + return await self._acquire_token(agent_id, tenant_id, entry, token_resolver) def get_observability_token(self, agent_id: str, tenant_id: str) -> str | None: """Get a non-expired cached app-only OBS token.""" @@ -226,7 +206,7 @@ async def _acquire_token( ) try: result = resolver(agent_id, tenant_id, list(entry.scopes)) - token = await cast(Awaitable[str | None], result) if isawaitable(result) else result + 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" @@ -305,12 +285,12 @@ def _clear_token(self, entry: _Entry) -> None: entry.expires_on_ms = None entry.acquired_on_ms = None - def _log_removed_overload_once(self) -> None: - if self._removed_overload_logged: + def _log_removed_registration_once(self) -> None: + if self._removed_registration_logged: return - self._removed_overload_logged = True + self._removed_registration_logged = True logger.error( - "[AgenticTokenCache] Delegated OBS token registration/refresh was removed and " - "does nothing; S2S OBS needs an app-only token. Call " + "[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/tests/observability/core/test_agent365.py b/tests/observability/core/test_agent365.py index cd42414d..ca91e13b 100644 --- a/tests/observability/core/test_agent365.py +++ b/tests/observability/core/test_agent365.py @@ -83,22 +83,31 @@ 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_defaults_to_legacy_s2s_true(self): - """Deprecated use_s2s_endpoint option defaults to True but routing ignores the value.""" - self.assertTrue(Agent365ExporterOptions().use_s2s_endpoint) + 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_fails_when_exporter_enabled_without_token_resolver(self, mock_is_enabled): - """Enabled Agent 365 exporter requires an explicit app-only OBS resolver.""" + 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 - result = configure( - service_name="test-service", - service_namespace="test-namespace", - exporter_options=Agent365ExporterOptions(), - ) + 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.assertFalse(result, "configure() should fail without an app-only token resolver") + 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") 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 9277b9da..1dce0bc3 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 @@ -112,34 +112,6 @@ async def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: ) -@pytest.mark.asyncio -async def test_legacy_refresh_shape_logs_once_and_does_not_exchange( - token_cache, mock_authorization, mock_turn_context, caplog -): - """Removed TurnContext/Authorization refresh shape is a safe no-op.""" - with caplog.at_level("ERROR"): - result1 = await token_cache.RefreshObservabilityToken( - "agent", - "tenant", - mock_turn_context, - mock_authorization, - ["api://old-scope/.default"], - "agentic", - ) - result2 = await token_cache.RefreshObservabilityToken( - "agent", - "tenant", - mock_turn_context, - mock_authorization, - ) - - assert result1 is None - assert result2 is None - mock_authorization.exchange_token.assert_not_called() - assert token_cache.get_observability_token("agent", "tenant") is None - assert sum("Delegated OBS token" in record.message for record in caplog.records) == 1 - - def test_legacy_register_shape_logs_once_and_does_not_exchange( token_cache, mock_authorization, mock_turn_context, caplog ): @@ -158,6 +130,13 @@ def test_legacy_register_shape_logs_once_and_does_not_exchange( 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): From 5d40056050a50f212f6740050b91bffdae632e2a Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 14:57:04 +0100 Subject: [PATCH 05/13] fix(hosting): preserve async token cache read Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../CHANGELOG.md | 7 +-- .../README.md | 10 ++-- .../token_cache_helpers/agent_token_cache.py | 8 ++- .../test_agent_token_cache.py | 50 +++++++++++-------- 4 files changed, 45 insertions(+), 30 deletions(-) diff --git a/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md index 069225a0..ab08f4f8 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md +++ b/libraries/microsoft-agents-a365-observability-hosting/CHANGELOG.md @@ -12,15 +12,16 @@ All notable changes to this package will be documented in this file. 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 `refresh_observability_token(agent_id, tenant_id, app_only_token_resolver)` - from the exporter `token_resolver`, then return - `get_observability_token(agent_id, tenant_id)`. +- 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. diff --git a/libraries/microsoft-agents-a365-observability-hosting/README.md b/libraries/microsoft-agents-a365-observability-hosting/README.md index 15988ebb..b85baa88 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/README.md +++ b/libraries/microsoft-agents-a365-observability-hosting/README.md @@ -21,14 +21,13 @@ from microsoft_agents_a365.observability.hosting.token_cache_helpers import Agen cache = AgenticTokenCache() -async def refresh_app_only_obs(agent_id: str, tenant_id: str, scopes: list[str]) -> str: +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 resolver(agent_id: str, tenant_id: str) -> str | None: - await cache.refresh_observability_token(agent_id, tenant_id, refresh_app_only_obs) - return cache.get_observability_token(agent_id, 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 @@ -46,6 +45,9 @@ from inside a running event loop cannot await an async resolver, so that export fails with a logged error; call them with `await asyncio.to_thread(...)`. A synchronous, thread-safe cached resolver avoids both constraints and is the pattern used by the Agent365-Samples Python samples. +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(...)` 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 302559b4..2a41933e 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 @@ -127,8 +127,12 @@ async def refresh_observability_token( return await self._acquire_token(agent_id, tenant_id, entry, token_resolver) - def get_observability_token(self, agent_id: str, tenant_id: str) -> str | None: - """Get a non-expired cached app-only OBS token.""" + 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. + """ key = self.make_key(agent_id, tenant_id) with self._lock: entry = self._map.get(key) 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 1dce0bc3..ccccbd9e 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 @@ -62,9 +62,10 @@ def token_cache(): return AgenticTokenCache() -def test_get_observability_token_returns_none_without_entry(token_cache): +@pytest.mark.asyncio +async def test_get_observability_token_returns_none_without_entry(token_cache): """A cache miss returns None.""" - assert token_cache.get_observability_token("agent", "tenant") is None + assert await token_cache.get_observability_token("agent", "tenant") is None @pytest.mark.asyncio @@ -80,7 +81,7 @@ def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: refreshed = await token_cache.refresh_observability_token("agent", "tenant", resolver) assert refreshed == token - assert token_cache.get_observability_token("agent", "tenant") == token + assert await token_cache.get_observability_token("agent", "tenant") == token assert calls == [("agent", "tenant", ["api://9b975845-388f-4429-889e-eab1ef63949c/.default"])] @@ -112,7 +113,8 @@ async def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: ) -def test_legacy_register_shape_logs_once_and_does_not_exchange( +@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.""" @@ -126,7 +128,7 @@ def test_legacy_register_shape_logs_once_and_does_not_exchange( token_cache.register_observability("agent", "tenant", token_struct, ["scope"]) mock_authorization.exchange_token.assert_not_called() - assert token_cache.get_observability_token("agent", "tenant") is None + 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 @@ -170,7 +172,7 @@ async def test_refresh_surfaces_empty_resolver_result(token_cache, token): with pytest.raises(RuntimeError, match="returned no token"): await token_cache.refresh_observability_token("agent", "tenant", resolver) - assert token_cache.get_observability_token("agent", "tenant") is None + assert await token_cache.get_observability_token("agent", "tenant") is None @pytest.mark.asyncio @@ -184,7 +186,7 @@ async def test_refresh_propagates_permanent_resolver_failure(token_cache): with pytest.raises(RuntimeError, match="permission denied"): await token_cache.refresh_observability_token("agent", "tenant", resolver) - assert token_cache.get_observability_token("agent", "tenant") is None + assert await token_cache.get_observability_token("agent", "tenant") is None @pytest.mark.asyncio @@ -202,7 +204,7 @@ async def no_sleep(delay: float) -> None: assert refreshed == token assert resolver.call_count == 2 - assert token_cache.get_observability_token("agent", "tenant") == token + assert await token_cache.get_observability_token("agent", "tenant") == token @pytest.mark.asyncio @@ -235,12 +237,12 @@ async def test_refresh_reuses_cached_token_until_expiry_skew(token_cache): await token_cache.refresh_observability_token("agent", "tenant", resolver) assert resolver.call_count == 1 - assert token_cache.get_observability_token("agent", "tenant") == token + 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 token_cache.get_observability_token("agent", "tenant") is None + assert await token_cache.get_observability_token("agent", "tenant") is None @pytest.mark.asyncio @@ -250,10 +252,10 @@ async def test_opaque_token_uses_fresh_fallback_ttl(token_cache): token_cache.invalidate_token("agent", "tenant") await token_cache.refresh_observability_token("agent", "tenant", lambda *_: "opaque-token") - assert token_cache.get_observability_token("agent", "tenant") == "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 token_cache.get_observability_token("agent", "tenant") is None + assert await token_cache.get_observability_token("agent", "tenant") is None @pytest.mark.asyncio @@ -267,9 +269,15 @@ def resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str: await token_cache.refresh_observability_token("agent-two", "tenant-a", resolver) await token_cache.refresh_observability_token("agent-one", "tenant-b", resolver) - assert token_cache.get_observability_token("agent-one", "tenant-a") == "agent-one:tenant-a" - assert token_cache.get_observability_token("agent-two", "tenant-a") == "agent-two:tenant-a" - assert token_cache.get_observability_token("agent-one", "tenant-b") == "agent-one:tenant-b" + 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" + ) @pytest.mark.asyncio @@ -280,11 +288,11 @@ async def test_invalidate_one_then_all(token_cache): await token_cache.refresh_observability_token("two", "tenant", lambda *_: token) token_cache.invalidate_token("one", "tenant") - assert token_cache.get_observability_token("one", "tenant") is None - assert token_cache.get_observability_token("two", "tenant") == token + 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 token_cache.get_observability_token("two", "tenant") is None + assert await token_cache.get_observability_token("two", "tenant") is None @pytest.mark.asyncio @@ -296,6 +304,6 @@ async def test_cache_evicts_oldest_entry_when_capacity_is_reached(token_cache): await token_cache.refresh_observability_token("two", "tenant", lambda *_: token) await token_cache.refresh_observability_token("three", "tenant", lambda *_: token) - assert token_cache.get_observability_token("one", "tenant") is None - assert token_cache.get_observability_token("two", "tenant") == token - assert token_cache.get_observability_token("three", "tenant") == 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 From 97c820114993160a6f10fba09355c339f5891576 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:01:38 +0100 Subject: [PATCH 06/13] fix(hosting): isolate OBS cache keys and bound refresh locks - Key the cache by an (agent_id, tenant_id) tuple instead of "agent_id:tenant_id", so IDs that contain a colon cannot share another identity's cached token. make_key (new in this PR) becomes private. - Keep each refresh lock on its cache entry, so capacity eviction and invalidate_all() reclaim it. Eviction skips entries with a refresh in flight. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../token_cache_helpers/agent_token_cache.py | 47 +++++++-------- .../test_agent_token_cache.py | 60 ++++++++++++++++++- 2 files changed, 82 insertions(+), 25 deletions(-) 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 2a41933e..63b99ff0 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 @@ -11,7 +11,7 @@ import logging import time from collections.abc import Awaitable, Callable, Sequence -from dataclasses import dataclass +from dataclasses import dataclass, field from inspect import isawaitable from threading import Lock @@ -59,6 +59,7 @@ class _Entry: 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) _default_refresh_skew_ms = 60_000 _default_max_token_age_ms = 3_600_000 @@ -67,8 +68,7 @@ class _Entry: def __init__(self, observability_scopes: Sequence[str] | None = None) -> None: """Initialize the token cache.""" - self._map: dict[str, AgenticTokenCache._Entry] = {} - self._key_locks: dict[str, asyncio.Lock] = {} + self._map: dict[tuple[str, str], AgenticTokenCache._Entry] = {} self._lock = Lock() self._observability_scopes = ( None if observability_scopes is None else tuple(observability_scopes) @@ -76,9 +76,9 @@ def __init__(self, observability_scopes: Sequence[str] | None = None) -> None: self._removed_registration_logged = False @staticmethod - def make_key(agent_id: str, tenant_id: str) -> str: - """Create a cache key for an agent and tenant.""" - return f"{agent_id}:{tenant_id}" + 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, @@ -118,10 +118,9 @@ async def refresh_observability_token( 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") - key = self.make_key(agent_id, tenant_id) - lock = self._get_key_lock(key) - async with lock: - entry = self._get_or_create_entry(key) + key = self._make_key(agent_id, tenant_id) + entry = self._get_or_create_entry(key) + async with entry.lock: if entry.token is not None and not self._is_expired(entry): return entry.token @@ -133,7 +132,7 @@ async def get_observability_token(self, agent_id: str, tenant_id: str) -> str | This method is a pure cache read. It never acquires a token and never calls delegated token exchange. """ - key = self.make_key(agent_id, tenant_id) + key = self._make_key(agent_id, tenant_id) with self._lock: entry = self._map.get(key) @@ -147,7 +146,7 @@ async def get_observability_token(self, agent_id: str, tenant_id: str) -> str | def invalidate_token(self, agent_id: str, tenant_id: str) -> None: """Invalidate one cached token.""" - key = self.make_key(agent_id, tenant_id) + key = self._make_key(agent_id, tenant_id) with self._lock: entry = self._map.get(key) if entry is not None: @@ -158,7 +157,7 @@ def invalidate_all(self) -> None: with self._lock: self._map.clear() - def _get_or_create_entry(self, key: str) -> _Entry: + def _get_or_create_entry(self, key: tuple[str, str]) -> _Entry: with self._lock: entry = self._map.get(key) if entry is not None: @@ -171,22 +170,22 @@ def _get_or_create_entry(self, key: str) -> _Entry: raise ValueError("[AgenticTokenCache] No valid scopes") if len(self._map) >= self._max_cache_size: - oldest_key = next(iter(self._map), None) - if oldest_key is not None: - del self._map[oldest_key] + # Evict the oldest idle entry; an entry with a refresh in flight keeps its lock. + idle_key = next( + ( + existing_key + for existing_key, existing in self._map.items() + if not existing.lock.locked() + ), + None, + ) + if idle_key is not None: + del self._map[idle_key] entry = AgenticTokenCache._Entry(scopes=scopes) self._map[key] = entry return entry - def _get_key_lock(self, key: str) -> asyncio.Lock: - with self._lock: - lock = self._key_locks.get(key) - if lock is None: - lock = asyncio.Lock() - self._key_locks[key] = lock - return lock - def _get_effective_scopes(self) -> tuple[str, ...]: scopes = self._observability_scopes if scopes is None: 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 ccccbd9e..849de941 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 @@ -253,7 +253,7 @@ async def test_opaque_token_uses_fresh_fallback_ttl(token_cache): 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 = 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 @@ -307,3 +307,61 @@ async def test_cache_evicts_oldest_entry_when_capacity_is_reached(token_cache): 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_entry_with_refresh_in_flight(token_cache): + """Capacity eviction skips an identity whose refresh is still running.""" + 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" From 2a7b0065ff598879009492fe0133ea2c6a2b00e5 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:10:00 +0100 Subject: [PATCH 07/13] fix(hosting): trim OBS cache overflow and guard the log-once flag - When every cached entry has a refresh in flight, a new identity can take the cache past its bound. Evict idle entries until there is room, so the next insertion trims any overflow instead of keeping it. - Guard the one-time removed-registration log with the cache lock. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../token_cache_helpers/agent_token_cache.py | 17 ++++++++++------- .../test_agent_token_cache.py | 7 +++++-- 2 files changed, 15 insertions(+), 9 deletions(-) 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 63b99ff0..d15dee4c 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 @@ -169,8 +169,9 @@ def _get_or_create_entry(self, key: tuple[str, str]) -> _Entry: if not scopes: raise ValueError("[AgenticTokenCache] No valid scopes") - if len(self._map) >= self._max_cache_size: - # Evict the oldest idle entry; an entry with a refresh in flight keeps its lock. + # 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 @@ -179,8 +180,9 @@ def _get_or_create_entry(self, key: tuple[str, str]) -> _Entry: ), None, ) - if idle_key is not None: - del self._map[idle_key] + if idle_key is None: + break + del self._map[idle_key] entry = AgenticTokenCache._Entry(scopes=scopes) self._map[key] = entry @@ -289,9 +291,10 @@ def _clear_token(self, entry: _Entry) -> None: entry.acquired_on_ms = None def _log_removed_registration_once(self) -> None: - if self._removed_registration_logged: - return - self._removed_registration_logged = True + 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 " 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 849de941..285e85af 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 @@ -344,8 +344,8 @@ async def test_identity_churn_keeps_every_per_identity_registry_bounded(token_ca @pytest.mark.asyncio -async def test_eviction_keeps_entry_with_refresh_in_flight(token_cache): - """Capacity eviction skips an identity whose refresh is still running.""" +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() @@ -365,3 +365,6 @@ async def slow_resolver(agent_id: str, tenant_id: str, scopes: list[str]) -> str 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 From 8ae0e73d718a4cdb52131a3bf7cd429d691ae88d Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:16:36 +0100 Subject: [PATCH 08/13] fix(hosting): retry standard transient errors; type resolver scopes as list - Treat TimeoutError and ConnectionError (including ConnectionResetError) as transient before the message and status checks. Bare instances have empty messages, so they were not retried. - Type the ObservabilityTokenResolver scopes parameter as list[str], which is what the cache passes and what the README documents. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../token_cache_helpers/agent_token_cache.py | 5 ++++- .../test_agent_token_cache.py | 16 ++++++++++++++++ 2 files changed, 20 insertions(+), 1 deletion(-) 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 d15dee4c..615cbff4 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 @@ -21,7 +21,7 @@ logger = logging.getLogger(__name__) -ObservabilityTokenResolver = Callable[[str, str, Sequence[str]], str | Awaitable[str | None] | None] +ObservabilityTokenResolver = Callable[[str, str, list[str]], str | Awaitable[str | None] | None] @dataclass @@ -271,6 +271,9 @@ def _is_expired(self, entry: _Entry) -> bool: 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 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 285e85af..72f376dc 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 @@ -207,6 +207,22 @@ async def no_sleep(delay: float) -> None: 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.""" From 4485092ac987c4cf0a05c2c730c658adb0881432 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:22:56 +0100 Subject: [PATCH 09/13] fix(hosting): treat a bare scope string as one OBS scope AgenticTokenCache(observability_scopes="api://.../.default") split the string into one scope per character. Wrap a bare string as a single scope. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../hosting/token_cache_helpers/agent_token_cache.py | 2 ++ .../token_cache_helpers/test_agent_token_cache.py | 11 +++++++++++ 2 files changed, 13 insertions(+) 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 615cbff4..976a0ca8 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 @@ -68,6 +68,8 @@ class _Entry: def __init__(self, observability_scopes: Sequence[str] | None = None) -> None: """Initialize the token cache.""" + if isinstance(observability_scopes, str): + observability_scopes = (observability_scopes,) self._map: dict[tuple[str, str], AgenticTokenCache._Entry] = {} self._lock = Lock() self._observability_scopes = ( 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 72f376dc..0ba2507a 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 @@ -97,6 +97,17 @@ async def test_refresh_uses_custom_app_only_scopes(): 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.""" From 2b3710c073e48fead308e2bef68247f4f21cc091 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 22:37:59 +0100 Subject: [PATCH 10/13] fix(hosting): keep OBS token invalidation from being undone - invalidate_token() now removes the cache entry under the cache lock, as invalidate_all() and the .NET cache do. A refresh that is still in flight writes to the removed entry, so it can no longer repopulate the cache after the invalidation completes. - refresh_observability_token() returns the token it checked. A concurrent invalidation between the usability check and the return could otherwise make it return None. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../token_cache_helpers/agent_token_cache.py | 17 ++++--- .../test_agent_token_cache.py | 44 +++++++++++++++++++ 2 files changed, 55 insertions(+), 6 deletions(-) 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 976a0ca8..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 @@ -123,8 +123,9 @@ async def refresh_observability_token( key = self._make_key(agent_id, tenant_id) entry = self._get_or_create_entry(key) async with entry.lock: - if entry.token is not None and not self._is_expired(entry): - return entry.token + token = entry.token + if token is not None and not self._is_expired(entry): + return token return await self._acquire_token(agent_id, tenant_id, entry, token_resolver) @@ -147,12 +148,16 @@ async def get_observability_token(self, agent_id: str, tenant_id: str) -> str | return entry.token def invalidate_token(self, agent_id: str, tenant_id: str) -> None: - """Invalidate one cached token.""" + """Invalidate one cached token. + + The entry is removed, so a refresh already in flight for this identity + cannot repopulate the cache. + """ key = self._make_key(agent_id, tenant_id) with self._lock: - entry = self._map.get(key) - if entry is not None: - self._clear_token(entry) + entry = self._map.pop(key, None) + if entry is not None: + self._clear_token(entry) def invalidate_all(self) -> None: """Invalidate all cached tokens.""" 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 0ba2507a..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 @@ -322,6 +322,50 @@ async def test_invalidate_one_then_all(token_cache): 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.""" From 1e83d758fcb0d24ca9f84111ad1ea1feba40b642 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 22:49:48 +0100 Subject: [PATCH 11/13] build(hosting): declare microsoft-agents-a365-runtime as a direct dependency The hosting token cache now imports get_observability_authentication_scope from microsoft-agents-a365-runtime. It was only available transitively through observability-core; declare it directly so the published metadata matches the imports. The build pins it to the same version like other internal packages. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../microsoft-agents-a365-observability-hosting/pyproject.toml | 1 + uv.lock | 2 ++ 2 files changed, 3 insertions(+) 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/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" }, ] From 6b358f071281a85086a1ffdee44b5b6fe3d565f7 Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 22:57:05 +0100 Subject: [PATCH 12/13] docs(hosting): distinguish force_flush and shutdown in async-resolver caveat In opentelemetry-sdk 1.39, force_flush() exports on the calling thread, so an async resolver fails when it is called from a running event loop. shutdown() runs its final export on the batch worker thread and only blocks the caller. The README previously said both would fail. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../README.md | 21 ++++++++++++------- 1 file changed, 14 insertions(+), 7 deletions(-) diff --git a/libraries/microsoft-agents-a365-observability-hosting/README.md b/libraries/microsoft-agents-a365-observability-hosting/README.md index b85baa88..8984b3e5 100644 --- a/libraries/microsoft-agents-a365-observability-hosting/README.md +++ b/libraries/microsoft-agents-a365-observability-hosting/README.md @@ -38,13 +38,20 @@ 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 BatchSpanProcessor worker thread. Create async clients -inside that resolver; do not reuse `aiohttp` or `azure.identity.aio` clients -bound to your app's event loop. Also, `force_flush()` and `shutdown()` called -from inside a running event loop cannot await an async resolver, so that export -fails with a logged error; call them with `await asyncio.to_thread(...)`. A -synchronous, thread-safe cached resolver avoids both constraints and is the -pattern used by the Agent365-Samples Python samples. +`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. From ddd0df1ced64263fe8ea4b511c1f2a25805baf1b Mon Sep 17 00:00:00 2001 From: Krishnadheeraj <12496535+DheerajPannala@users.noreply.github.com> Date: Tue, 29 Sep 2026 23:03:57 +0100 Subject: [PATCH 13/13] fix(observability): cancel resolver futures when export can't await them When export runs inside an active event loop, the exporter can't await an async token resolver and fails that export. It closed bare coroutines but left a resolver-returned asyncio.Future or Task running, where it could still mutate a token cache or raise an unobserved exception. Cancel those too. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 5cbf5f6b-cc40-4b7e-a591-65848db73a12 --- .../core/exporters/agent365_exporter.py | 2 ++ .../core/test_agent365_exporter.py | 30 +++++++++++++++++++ 2 files changed, 32 insertions(+) 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 8ad3ead4..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 @@ -299,6 +299,8 @@ def _resolve_token(self, agent_id: str, tenant_id: str) -> str | None: 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 " diff --git a/tests/observability/core/test_agent365_exporter.py b/tests/observability/core/test_agent365_exporter.py index f9d72228..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 @@ -352,6 +354,34 @@ async def resolver(agent_id, tenant_id): 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):