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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 13 additions & 3 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,15 +8,25 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/).

### Added

- Non-secret configuration variables are written to the log when the service
starts.

### Changed

### Fixed

### Removed

## [1.4.13] - 2026-08-27

### Added

- Non-secret configuration variables are written to the log when the service
starts.

### Changed

- Azure blob service clients are now obtained from the service factory, which
applies an explicit exponential retry policy rather than relying on the
Azure SDK defaults.

## [1.4.12] - 2026-08-03

### Changed
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "bulk-data-service"
version = "1.4.12"
version = "1.4.13"
requires-python = ">= 3.12.6"
readme = "README.md"
dependencies = [
Expand Down
4 changes: 1 addition & 3 deletions src/bulk_data_service/dataset_indexing.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,6 @@
from datetime import datetime
from typing import Any

from azure.storage.blob import BlobServiceClient

from bulk_data_service.data_converters import (
convert_reporting_org_to_reporting_org_dto,
get_full_dataset_check_result_dto,
Expand Down Expand Up @@ -43,7 +41,7 @@ def create_and_upload_indices(

def upload_index_json_to_azure(context: BDSContext, index_name: str, index_json: str):

az_blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
az_blob_service = context.service_factory.get_azure_blob_service_client()

azure_upload_to_blob(
az_blob_service, context["AZURE_STORAGE_BLOB_CONTAINER_NAME"], index_name, index_json, "application/json"
Expand Down
4 changes: 2 additions & 2 deletions src/bulk_data_service/dataset_remover.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ def remove_deleted_datasets_from_bds(

db_conn = get_db_connection(context)

az_blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
az_blob_service = context.service_factory.get_azure_blob_service_client()

ids_to_delete = [k for k in datasets_in_bds.keys() if k not in registered_datasets]

Expand Down Expand Up @@ -49,7 +49,7 @@ def remove_expired_downloads(context: BDSContext, datasets_in_bds: dict[uuid.UUI

db_conn = get_db_connection(context)

az_blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
az_blob_service = context.service_factory.get_azure_blob_service_client()

expired_datasets = 0

Expand Down
2 changes: 1 addition & 1 deletion src/bulk_data_service/dataset_updater.py
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ def add_or_update_dataset_batch(

db_conn = get_db_connection(context)

az_blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
az_blob_service = context.service_factory.get_azure_blob_service_client()

session = get_requests_session(context)

Expand Down
3 changes: 1 addition & 2 deletions src/bulk_data_service/zipper.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@
import uuid

from azure.core.exceptions import ResourceNotFoundError
from azure.storage.blob import BlobServiceClient

from bulk_data_service.zippers import CodeforIATILegacyZipper, IATIBulkDataServiceZipper
from config.bds_context import BDSContext
Expand Down Expand Up @@ -175,7 +174,7 @@ def remove_datasets_without_dls_from_working_dir(

def download_new_or_updated_to_working_dir(context: BDSContext, updated_datasets: dict[uuid.UUID, dict]):

az_blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
az_blob_service = context.service_factory.get_azure_blob_service_client()

xml_container_name = get_azure_container_name(context, "xml")

Expand Down
2 changes: 1 addition & 1 deletion src/bulk_data_service/zippers.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ def get_zip_local_pathname_no_extension(self) -> str:
class IATIBulkDataServiceZipper(IATIDataZipper):

def prepare(self):
az_blob_service = BlobServiceClient.from_connection_string(self.context["AZURE_STORAGE_CONNECTION_STRING"])
az_blob_service = self.context.service_factory.get_azure_blob_service_client()

self.download_dataset_index_to_working_dir(az_blob_service, "minimal")

Expand Down
28 changes: 28 additions & 0 deletions src/config/service_factory.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
from azure.servicebus import ServiceBusClient
from azure.storage.blob import BlobServiceClient, ExponentialRetry
from libsuitecrm import SuiteCRM # type: ignore

from .service_factory_interface import IServiceFactory
Expand All @@ -16,3 +17,30 @@ def get_suitecrm_client(self) -> SuiteCRM:
self._config["DATA_REGISTRY_SUITECRM_CLIENT_SECRET"],
secure=self._config["DATA_REGISTRY_SUITECRM_SECURE"] != "false",
)

def get_azure_blob_service_client(self) -> BlobServiceClient:
"""Returns a blob service client which retries failed requests. These are the
values the Azure SDK applies by default, set out here so that the behaviour is
visible and can be changed deliberately.

A request is retried on a connection or read error, on an HTTP 408, and on a 5xx
response other than 501 and 505.

The wait before each retry is `initial_backoff + increment_base ** retry_count`,
varied by up to `random_jitter_range` seconds either way. The SDK raises the
retry count before working out the wait, so the count starts at one and
`initial_backoff` is the base of every wait rather than the first one: the three
retries are attempted after roughly 18, 24 and 42 seconds, delaying a request
which never succeeds by about 84 seconds in total."""

retry_policy = ExponentialRetry(
initial_backoff=15,
increment_base=3,
retry_total=3,
retry_to_secondary=False,
random_jitter_range=3,
)

return BlobServiceClient.from_connection_string(
self._config["AZURE_STORAGE_CONNECTION_STRING"], retry_policy=retry_policy
)
5 changes: 5 additions & 0 deletions src/config/service_factory_interface.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
from typing import Any

from azure.servicebus import ServiceBusClient
from azure.storage.blob import BlobServiceClient
from libsuitecrm import SuiteCRM # type: ignore


Expand All @@ -17,3 +18,7 @@ def get_service_bus_client(self, conn_str: str, enable_logging: bool = False) ->
@abc.abstractmethod
def get_suitecrm_client(self) -> SuiteCRM:
raise NotImplementedError

@abc.abstractmethod
def get_azure_blob_service_client(self) -> BlobServiceClient:
raise NotImplementedError
6 changes: 3 additions & 3 deletions src/utilities/azure.py
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ def azure_upload_to_blob(


def create_azure_blob_containers(context: BDSContext):
blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
blob_service = context.service_factory.get_azure_blob_service_client()

containers = blob_service.list_containers()
container_names = [c.name for c in containers]
Expand All @@ -120,7 +120,7 @@ def create_azure_blob_containers(context: BDSContext):


def delete_azure_blob_containers(context: BDSContext):
blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
blob_service = context.service_factory.get_azure_blob_service_client()

containers = blob_service.list_containers()
container_names = [c.name for c in containers]
Expand Down Expand Up @@ -214,7 +214,7 @@ def send_message_to_iati_mq(context: BDSContext, topic_name, msg_payload):


def upload_zip_to_azure(context: BDSContext, zip_local_pathname: str, zip_azure_filename: str):
az_blob_service = BlobServiceClient.from_connection_string(context["AZURE_STORAGE_CONNECTION_STRING"])
az_blob_service = context.service_factory.get_azure_blob_service_client()

blob_client = az_blob_service.get_blob_client(context["AZURE_STORAGE_BLOB_CONTAINER_NAME"], zip_azure_filename)

Expand Down
9 changes: 8 additions & 1 deletion tests/helpers/helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@

from config.bds_context import BDSContext
from config.config import get_app_version
from config.service_factory import ServiceFactory
from config.service_factory_interface import IServiceFactory
from utilities.azure import (
create_azure_blob_containers,
Expand Down Expand Up @@ -84,7 +85,13 @@ def get_and_clear_up_context():
for metric in get_metrics_definitions():
config["prom_metrics"][metric[0]] = mock.Mock() # type: ignore

context = BDSContext(config, logger, mock.create_autospec(IServiceFactory))
# the Service Bus and SuiteCRM clients are mocked, so that tests can stub them, but
# blob storage is exercised for real against the Azurite emulator, so the blob client
# is delegated to the real service factory
service_factory = mock.create_autospec(IServiceFactory)
service_factory.get_azure_blob_service_client.side_effect = ServiceFactory(config).get_azure_blob_service_client

context = BDSContext(config, logger, service_factory)

context["TEST_TMP_ZIP_UNPACK"] = "tests/tmp_zip_unpack"

Expand Down
14 changes: 5 additions & 9 deletions tests/integration/test_dataset_update.py
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,7 @@ def test_update_dataset_registration_details(get_and_clear_up_context, field, or
assert datasets_in_bds[dataset_id][field] == expected


def test_update_dataset_mq_message_send_doesnt_crash_on_error(monkeypatch, get_and_clear_up_context): # noqa: F811
def test_update_dataset_mq_message_send_doesnt_crash_on_error(get_and_clear_up_context): # noqa: F811

context = get_and_clear_up_context

Expand All @@ -220,14 +220,10 @@ def get_topic_sender(self, *args, **kwargs):
def close(self):
return None

class FakeServiceFactory:
def get_service_bus_client(self, *args, **kwargs):
return FakeServiceBusClient()

def get_suitecrm_client(self):
raise NotImplementedError

monkeypatch.setattr(context, "_service_factory", FakeServiceFactory())
# only the Service Bus client is faked: the rest of the factory, including the blob
# storage client, is left as the fixture set it up
context.service_factory.get_service_bus_client.return_value = FakeServiceBusClient()
context.service_factory.get_suitecrm_client.side_effect = NotImplementedError

dataset_id = uuid.UUID("c8a40aa5-9f31-4bcf-a36f-51c1fc2cc159")

Expand Down
87 changes: 87 additions & 0 deletions tests/unit/test_service_factory.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
from typing import Any
from unittest import mock

from azure.storage.blob import BlobServiceClient, ExponentialRetry

from config.service_factory import ServiceFactory

# a syntactically valid connection string: these tests make no connection
AZURE_STORAGE_CONNECTION_STRING = (
"DefaultEndpointsProtocol=http;"
"AccountName=devstoreaccount1;"
"AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;"
"BlobEndpoint=http://127.0.0.1:10000/devstoreaccount1;"
)


def get_blob_service_client_arguments(monkeypatch) -> dict[str, Any]:
"""Returns the keyword arguments which the service factory passes when it builds a
blob service client. Taken from the call rather than read back off the client, whose
retry policy is only reachable through private attributes."""

arguments: dict[str, Any] = {}

def capture_arguments(conn_str: str, **kwargs: Any) -> Any:
arguments["connection_string"] = conn_str
arguments.update(kwargs)
return mock.Mock()

monkeypatch.setattr(BlobServiceClient, "from_connection_string", capture_arguments)

service_factory = ServiceFactory({"AZURE_STORAGE_CONNECTION_STRING": AZURE_STORAGE_CONNECTION_STRING})

service_factory.get_azure_blob_service_client()

return arguments


def get_retry_policy(monkeypatch) -> ExponentialRetry:
retry_policy = get_blob_service_client_arguments(monkeypatch)["retry_policy"]

assert isinstance(retry_policy, ExponentialRetry)

return retry_policy


def test_blob_service_client_is_given_the_configured_connection_string(monkeypatch):

arguments = get_blob_service_client_arguments(monkeypatch)

assert arguments["connection_string"] == AZURE_STORAGE_CONNECTION_STRING


def test_blob_service_client_is_given_an_explicit_retry_policy(monkeypatch):

assert "retry_policy" in get_blob_service_client_arguments(monkeypatch)


def test_blob_service_client_retries_with_exponential_backoff(monkeypatch):

assert isinstance(get_blob_service_client_arguments(monkeypatch)["retry_policy"], ExponentialRetry)


def test_blob_service_client_retry_values_are_set_explicitly(monkeypatch):

retry_policy = get_retry_policy(monkeypatch)

assert retry_policy.total_retries == 3
assert retry_policy.connect_retries == 3
assert retry_policy.read_retries == 3
assert retry_policy.status_retries == 3
assert retry_policy.retry_to_secondary is False
assert retry_policy.initial_backoff == 15
assert retry_policy.increment_base == 3
assert retry_policy.random_jitter_range == 3


def test_blob_service_client_waits_as_documented_before_each_retry(monkeypatch):
"""Checks the backoff values which get_azure_blob_service_client's docstring quotes.
The SDK raises the retry count before asking for the backoff time, so the first
retry is calculated with a count of one rather than zero."""

retry_policy = get_retry_policy(monkeypatch)

for retry_count, expected_wait in [(1, 18), (2, 24), (3, 42)]:
wait = retry_policy.get_backoff_time({"count": retry_count})

assert expected_wait - 3 <= wait <= expected_wait + 3
Loading