From 2ad20857805540bebb1250218f4bd621b9cdfc07 Mon Sep 17 00:00:00 2001 From: Bernd Verst Date: Sun, 26 Jul 2026 23:11:54 -0700 Subject: [PATCH] Performance: cache Azure managed SDK version lookup Both the sync and async DTS interceptor constructors called importlib.metadata.version('durabletask-azuremanaged') directly, which walks distribution metadata on disk. That work was repeated on every interceptor construction, i.e. on every client and worker construction. Resolve the version once in a module-level lru_cache'd helper and reuse it. The "unknown" fallback, the user-agent string format, and the metadata key names and their order are all unchanged. Verified with a probe that the number of importlib.metadata.version() calls across three interceptor constructions drops from 3 to 1. Fixes #195 Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e7499a54-f262-49a6-bf91-274af28543ae --- .../internal/durabletask_grpc_interceptor.py | 33 +++++----- .../test_durabletask_grpc_interceptor.py | 60 ++++++++++++++++++- 2 files changed, 78 insertions(+), 15 deletions(-) diff --git a/durabletask-azuremanaged/durabletask/azuremanaged/internal/durabletask_grpc_interceptor.py b/durabletask-azuremanaged/durabletask/azuremanaged/internal/durabletask_grpc_interceptor.py index e7cbb10a..3c4aecc7 100644 --- a/durabletask-azuremanaged/durabletask/azuremanaged/internal/durabletask_grpc_interceptor.py +++ b/durabletask-azuremanaged/durabletask/azuremanaged/internal/durabletask_grpc_interceptor.py @@ -1,6 +1,7 @@ # Copyright (c) Microsoft Corporation. # Licensed under the MIT License. +from functools import lru_cache from importlib.metadata import version import grpc @@ -17,6 +18,22 @@ ) +@lru_cache(maxsize=1) +def _get_sdk_version() -> str: + """Return the installed version of the azuremanaged package. + + Resolving the version walks distribution metadata on disk, so the result is + cached and shared by every interceptor instance instead of being recomputed + on each client or worker construction. Falls back to ``"unknown"`` when the + version cannot be determined. + """ + try: + return version('durabletask-azuremanaged') + except Exception: + # Fallback if version cannot be determined + return "unknown" + + class DTSDefaultClientInterceptorImpl (DefaultClientInterceptorImpl): """The class implements a UnaryUnaryClientInterceptor, UnaryStreamClientInterceptor, StreamUnaryClientInterceptor and StreamStreamClientInterceptor from grpc to add an @@ -27,13 +44,7 @@ def __init__( token_credential: TokenCredential | None, taskhub_name: str, worker_id: str | None = None): - try: - # Get the version of the azuremanaged package - sdk_version = version('durabletask-azuremanaged') - except Exception: - # Fallback if version cannot be determined - sdk_version = "unknown" - user_agent = f"durabletask-python/{sdk_version}" + user_agent = f"durabletask-python/{_get_sdk_version()}" self._metadata = [ ("taskhub", taskhub_name), ("x-user-agent", user_agent)] # 'user-agent' is a reserved header; use 'x-user-agent' @@ -81,13 +92,7 @@ class DTSAsyncDefaultClientInterceptorImpl(DefaultAsyncClientInterceptorImpl): (task hub name, user agent, and authentication token) to all async calls.""" def __init__(self, token_credential: AsyncTokenCredential | None, taskhub_name: str): - try: - # Get the version of the azuremanaged package - sdk_version = version('durabletask-azuremanaged') - except Exception: - # Fallback if version cannot be determined - sdk_version = "unknown" - user_agent = f"durabletask-python/{sdk_version}" + user_agent = f"durabletask-python/{_get_sdk_version()}" self._metadata = [ ("taskhub", taskhub_name), ("x-user-agent", user_agent)] diff --git a/tests/durabletask-azuremanaged/test_durabletask_grpc_interceptor.py b/tests/durabletask-azuremanaged/test_durabletask_grpc_interceptor.py index 878253a9..856e8772 100644 --- a/tests/durabletask-azuremanaged/test_durabletask_grpc_interceptor.py +++ b/tests/durabletask-azuremanaged/test_durabletask_grpc_interceptor.py @@ -4,15 +4,21 @@ import unittest from concurrent import futures from datetime import datetime, timedelta, timezone -from importlib.metadata import version +from importlib.metadata import PackageNotFoundError, version import threading import time +from unittest import mock import grpc from azure.core.credentials import AccessToken from durabletask.azuremanaged.client import DurableTaskSchedulerClient +from durabletask.azuremanaged.internal import durabletask_grpc_interceptor from durabletask.azuremanaged.internal.access_token_manager import AccessTokenManager +from durabletask.azuremanaged.internal.durabletask_grpc_interceptor import ( + DTSAsyncDefaultClientInterceptorImpl, + DTSDefaultClientInterceptorImpl, +) from durabletask.azuremanaged.worker import DurableTaskSchedulerWorker from durabletask.internal.grpc_interceptor import DefaultClientInterceptorImpl from durabletask.internal import orchestrator_service_pb2 as pb @@ -141,6 +147,58 @@ def test_worker_includes_workerid_header(self): self.assertTrue(metadata["workerid"]) +class TestSdkVersionCaching(unittest.TestCase): + """Tests that the azuremanaged SDK version is resolved once and reused.""" + + def setUp(self): + durabletask_grpc_interceptor._get_sdk_version.cache_clear() + + def tearDown(self): + # Drop any patched value so later tests observe the real package version. + durabletask_grpc_interceptor._get_sdk_version.cache_clear() + + def test_version_resolved_once_across_interceptor_constructions(self): + """The distribution metadata lookup happens once, not per interceptor.""" + with mock.patch.object( + durabletask_grpc_interceptor, "version", return_value="1.2.3") as mock_version: + client_interceptor = DTSDefaultClientInterceptorImpl(None, "test-taskhub") + worker_interceptor = DTSDefaultClientInterceptorImpl( + None, "test-taskhub", worker_id="test-worker-id") + async_interceptor = DTSAsyncDefaultClientInterceptorImpl(None, "test-taskhub") + + mock_version.assert_called_once_with('durabletask-azuremanaged') + + # The cached value is reused verbatim, and the metadata keys and their + # order are unchanged. + self.assertEqual( + [("taskhub", "test-taskhub"), ("x-user-agent", "durabletask-python/1.2.3")], + client_interceptor._metadata) + self.assertEqual( + [("taskhub", "test-taskhub"), + ("x-user-agent", "durabletask-python/1.2.3"), + ("workerid", "test-worker-id")], + worker_interceptor._metadata) + self.assertEqual( + [("taskhub", "test-taskhub"), ("x-user-agent", "durabletask-python/1.2.3")], + async_interceptor._metadata) + + def test_unknown_fallback_when_version_cannot_be_determined(self): + """A missing distribution still yields the 'unknown' user agent fallback.""" + with mock.patch.object( + durabletask_grpc_interceptor, + "version", + side_effect=PackageNotFoundError('durabletask-azuremanaged')) as mock_version: + client_interceptor = DTSDefaultClientInterceptorImpl(None, "test-taskhub") + async_interceptor = DTSAsyncDefaultClientInterceptorImpl(None, "test-taskhub") + + # The failed lookup is cached too, so it is not retried per interceptor. + mock_version.assert_called_once_with('durabletask-azuremanaged') + self.assertEqual( + "durabletask-python/unknown", dict(client_interceptor._metadata)["x-user-agent"]) + self.assertEqual( + "durabletask-python/unknown", dict(async_interceptor._metadata)["x-user-agent"]) + + class _TestTokenCredential: def __init__(self): self._lock = threading.Lock()