Skip to content

Commit ddd0df1

Browse files
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
1 parent 6b358f0 commit ddd0df1

2 files changed

Lines changed: 32 additions & 0 deletions

File tree

‎libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/exporters/agent365_exporter.py‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -299,6 +299,8 @@ def _resolve_token(self, agent_id: str, tenant_id: str) -> str | None:
299299
return asyncio.run(_await_token(cast(Awaitable[str | None], token)))
300300
if inspect.iscoroutine(token):
301301
token.close()
302+
elif isinstance(token, asyncio.Future):
303+
token.cancel()
302304
raise RuntimeError(
303305
"Agent365Exporter cannot await an async token_resolver while running "
304306
"inside an active event loop; use a synchronous cached resolver or refresh "

‎tests/observability/core/test_agent365_exporter.py‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
# Copyright (c) Microsoft Corporation.
22
# Licensed under the MIT License.
33

4+
import asyncio
5+
import inspect
46
import json
57
import os
68
import unittest
@@ -352,6 +354,34 @@ async def resolver(agent_id, tenant_id):
352354
self.assertEqual(result, SpanExportResult.SUCCESS)
353355
mock_post.assert_called_once()
354356

357+
def test_awaitable_resolver_in_active_event_loop_fails_and_is_released(self):
358+
"""Inside a running loop, export fails and the resolver's coroutine, future or task is released."""
359+
360+
async def pending_token():
361+
await asyncio.sleep(3600)
362+
return "never-used"
363+
364+
async def run_in_loop():
365+
coroutine = pending_token()
366+
future = asyncio.get_running_loop().create_future()
367+
task = asyncio.ensure_future(pending_token())
368+
for awaitable in (coroutine, future, task):
369+
exporter = _Agent365Exporter(
370+
token_resolver=lambda _agent_id, _tenant_id, value=awaitable: value,
371+
cluster_category="test",
372+
)
373+
spans = [self._create_mock_span("active_loop_span")]
374+
with patch.object(exporter, "_post_with_retries", return_value=True) as mock_post:
375+
self.assertEqual(exporter.export(spans), SpanExportResult.FAILURE)
376+
mock_post.assert_not_called()
377+
378+
await asyncio.wait([task], timeout=1)
379+
self.assertEqual(inspect.getcoroutinestate(coroutine), inspect.CORO_CLOSED)
380+
self.assertTrue(future.cancelled())
381+
self.assertTrue(task.cancelled())
382+
383+
asyncio.run(run_in_loop())
384+
355385
def test_auth_not_found_errors_do_not_fallback_to_delegated_route(self):
356386
"""401/403/404 failures do not trigger a delegated route fallback."""
357387
for status_code in (401, 403, 404):

0 commit comments

Comments
 (0)