Skip to content
Draft
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
2 changes: 2 additions & 0 deletions .release_notes/.unreleased.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,3 +5,5 @@
## New Features

## Bug Fixes

- Telemetry is now best-effort and never crashes a storage operation. Provider/exporter construction failures disable metrics/traces (and cache the disabled state); a failing metric record latches metrics off for the storage client instead of raising; and metrics init is skipped entirely when `MSC_TELEMETRY_DISABLED` or `OTEL_SDK_DISABLED` is set. The `_otlp_mtls_vault` exporter also fails fast with a clear error when the installed OTLP exporter is too old for its mTLS client-certificate options.
16 changes: 16 additions & 0 deletions multi-storage-client-docs/src/user_guide/telemetry.rst
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,22 @@ If the default telemetry provider creation doesn't behave as desired, you can ma
# Use an MSC shortcut to create a storage client for a profile and open an object/file.
multistorageclient.open("msc://data/file.txt")

*******************
Disabling Telemetry
*******************

Telemetry is best-effort observability and never crashes a storage operation. Any telemetry failure (e.g. an exporter that fails to construct, missing optional dependencies, or an unreachable telemetry process) is logged and metrics are disabled for the affected storage client; the storage operation always proceeds.

Telemetry can also be disabled explicitly, before any telemetry process, exporter, or IPC setup, by setting either of these environment variables to a truthy value (``1``, ``true``, or ``yes``):

.. code-block:: shell

# MSC-specific switch.
export MSC_TELEMETRY_DISABLED=true

# OpenTelemetry SDK standard switch (also honored).
export OTEL_SDK_DISABLED=true

**************
Authentication
**************
Expand Down
33 changes: 32 additions & 1 deletion multi-storage-client/src/multistorageclient/providers/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -243,6 +243,10 @@ def __init__(
self._metrics_dropped_count = 0
self._metrics_dropped_count_lock = threading.Lock()

# Latched off after an unrecoverable metric-record failure so telemetry never
# crashes a storage operation and failures aren't logged per-operation.
self._metrics_disabled = False

def __str__(self) -> str:
return self._provider_name

Expand All @@ -254,6 +258,15 @@ def __del__(self) -> None:
except Exception as e:
logger.warning(f"Failed to shutdown async telemetry: {e}", exc_info=True)

@staticmethod
def _telemetry_disabled_via_env() -> bool:
"""Whether telemetry is explicitly disabled via environment variable."""
for name in ("MSC_TELEMETRY_DISABLED", "OTEL_SDK_DISABLED"):
value = os.environ.get(name)
if value is not None and value.strip().lower() in ("1", "true", "yes"):
return True
return False

def _init_metrics(self) -> None:
"""
Initialize metrics.
Expand All @@ -266,6 +279,10 @@ def _init_metrics(self) -> None:
"""
with self._metric_init_lock:
if not self._metric_init_event.is_set():
if self._telemetry_disabled_via_env():
logger.debug("Telemetry disabled via environment variable; skipping metrics initialization.")
self._metric_init_event.set()
return
if self._config_dict is not None and self._telemetry_provider is not None:
opentelemetry_config: Optional[dict[str, Any]] = self._config_dict.get("opentelemetry")
if opentelemetry_config is not None:
Expand Down Expand Up @@ -571,7 +588,14 @@ def _dispatch_metrics(
Unlike :meth:`_emit_metrics` which wraps a callable, this method accepts
pre-computed metric values. Used by :meth:`_emit_metrics_sync`, :meth:`_emit_metrics_async`,
and directly when the operation is performed outside the standard wrapper (e.g. async Rust downloads).

Telemetry is best-effort: a failing synchronous record (e.g. a dead telemetry manager
or a disabled instrument) must never crash the storage operation. Any failure latches
metrics off so subsequent operations silently no-op instead of logging per-operation.
"""
if self._metrics_disabled:
return

if self._async_metrics_enabled and self._metrics_queue is not None:
metric_data = {
"operation": operation,
Expand All @@ -585,7 +609,14 @@ def _dispatch_metrics(
with self._metrics_dropped_count_lock:
self._metrics_dropped_count += 1
else:
self._record_metrics(operation, latency, data_size, error_type)
try:
self._record_metrics(operation, latency, data_size, error_type)
except (EOFError, BrokenPipeError, ConnectionError):
self._metrics_disabled = True
logger.warning("Telemetry manager connection closed; disabling metrics.")
except Exception:
self._metrics_disabled = True
logger.warning("Failed to record metrics; disabling metrics.", exc_info=True)

def _append_delimiter(self, s: str, delimiter: str = "/") -> str:
if not s.endswith(delimiter):
Expand Down
26 changes: 16 additions & 10 deletions multi-storage-client/src/multistorageclient/telemetry/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,8 @@ class CounterName(enum.Enum):
}

# Map of config as a sorted JSON string (since dictionaries can't be hashed) to meter provider.
_meter_provider_cache: dict[str, api_metrics.MeterProvider]
# A cached ``None`` means the config is disabled; it prevents re-running (failed) exporter construction.
_meter_provider_cache: dict[str, api_metrics.MeterProvider | None]
_meter_provider_cache_lock: threading.Lock
# Map of config as a sorted JSON string (since dictionaries can't be hashed) to meter.
_meter_cache: dict[str, api_metrics.Meter]
Expand All @@ -126,7 +127,8 @@ class CounterName(enum.Enum):
_counter_cache: dict[str, dict[CounterName, api_metrics.Counter]]
_counter_cache_lock: threading.Lock
# Map of config as a sorted JSON string (since dictionaries can't be hashed) to tracer provider.
_tracer_provider_cache: dict[str, api_trace.TracerProvider]
# A cached ``None`` means the config is disabled; it prevents re-running (failed) exporter construction.
_tracer_provider_cache: dict[str, api_trace.TracerProvider | None]
_tracer_provider_cache_lock: threading.Lock
# Map of config as a sorted JSON string (since dictionaries can't be hashed) to tracer.
_tracer_cache: dict[str, api_trace.Tracer]
Expand Down Expand Up @@ -194,15 +196,17 @@ def meter_provider(self, config: dict[str, Any]) -> Optional[api_metrics.MeterPr
return self._meter_provider_cache.setdefault(
config_json, sdk_metrics.MeterProvider(metric_readers=[reader])
)
except (AttributeError, ImportError):
except Exception:
logger.error(
"Failed to import OpenTelemetry Python SDK or exporter! Disabling metrics.", exc_info=True
"Failed to initialize the OpenTelemetry meter provider or exporter! Disabling metrics.",
exc_info=True,
)
return None
# Cache the disabled state so (possibly expensive) exporter construction isn't retried.
return self._meter_provider_cache.setdefault(config_json, None)
else:
# Don't return a no-op meter provider to avoid unnecessary overhead.
logger.error("No exporter configured! Disabling metrics.")
return None
return self._meter_provider_cache.setdefault(config_json, None)

def meter(self, config: dict[str, Any]) -> Optional[api_metrics.Meter]:
"""
Expand Down Expand Up @@ -301,14 +305,16 @@ def tracer_provider(self, config: dict[str, Any]) -> Optional[api_trace.TracerPr
config_json,
sdk_trace.TracerProvider(active_span_processor=processor, sampler=sampler),
)
except (AttributeError, ImportError):
except Exception:
logger.error(
"Failed to import OpenTelemetry Python SDK or exporter! Disabling traces.", exc_info=True
"Failed to initialize the OpenTelemetry tracer provider or exporter! Disabling traces.",
exc_info=True,
)
return None
# Cache the disabled state so (possibly expensive) exporter construction isn't retried.
return self._tracer_provider_cache.setdefault(config_json, None)
else:
logger.error("No exporter configured! Disabling traces.")
return None
return self._tracer_provider_cache.setdefault(config_json, None)

def tracer(self, config: dict[str, Any]) -> Optional[api_trace.Tracer]:
"""
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import inspect
import logging
from typing import Any

Expand Down Expand Up @@ -67,7 +68,21 @@ def __init__(
- key_key: Key name for client key (default: "key")
- ca_key: Key name for CA certificate (default: "ca")
:param exporter: OTLP metric exporter config dictionary (passed through to OTLPMetricExporter).
:raises RuntimeError: If the installed OTLP exporter does not support mTLS client certificate options.
"""
# mTLS client certificate options were added in opentelemetry-exporter-otlp-proto-http 1.32.1.
# Fail fast with a clear error (instead of an opaque TypeError from the parent) on older versions.
parameters = inspect.signature(OTLPMetricExporter.__init__).parameters
supports_client_certificate = "client_certificate_file" in parameters or any(
parameter.kind is inspect.Parameter.VAR_KEYWORD for parameter in parameters.values()
)
if not supports_client_certificate:
raise RuntimeError(
"The installed opentelemetry-exporter-otlp-proto-http does not support the mTLS client "
"certificate options (client_certificate_file / client_key_file). Upgrade to "
"opentelemetry-exporter-otlp-proto-http >= 1.32.1 to use the _otlp_mtls_vault exporter."
)

provider = VaultCertificateProvider(**auth)
cert_paths = provider.get_certificates()

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@
# limitations under the License.

import asyncio
import logging
import multiprocessing
import tempfile
import time
from collections.abc import Iterator
Expand All @@ -24,7 +26,7 @@
import pytest

from multistorageclient.providers.base import BaseStorageProvider
from multistorageclient.telemetry import Telemetry
from multistorageclient.telemetry import Telemetry, TelemetryManager
from multistorageclient.types import BatchTransferError, ObjectMetadata, Range, RetryableError, SymlinkHandling


Expand Down Expand Up @@ -1206,3 +1208,150 @@ def test_download_file_follows_symlink():
import os

os.unlink(tmp_path)


# ---------------------------------------------------------------------------
# Telemetry crash-safety.
#
# Telemetry is best-effort observability and must never crash a storage
# operation. These mirror ``test_async_metrics_handles_errors_in_worker`` for
# the synchronous record path, plus the manager-mode disabled-instrument and
# kill-switch cases.
# ---------------------------------------------------------------------------


def _sync_metrics_config() -> dict[str, Any]:
return {
"opentelemetry": {
"metrics": {
"exporter": {"type": "console"},
"reader": {"options": {}},
}
}
}


@pytest.mark.parametrize("failing_instrument", ["counter", "gauge"])
@pytest.mark.parametrize("error_type", [EOFError, BrokenPipeError, ConnectionResetError, RuntimeError, Exception])
def test_sync_metrics_recording_error_does_not_propagate(error_type, failing_instrument, caplog):
"""A failing synchronous metric record must not propagate out of a storage operation,
must be logged once, and must latch metrics off."""
caplog.set_level(logging.WARNING, logger="multistorageclient.providers.base")

# Distinct instruments so either the counter's .add() or the gauge's .set() can be the failing call.
mock_gauge = Mock()
mock_counter = Mock()
if failing_instrument == "gauge":
mock_gauge.set.side_effect = error_type("telemetry gone")
else:
mock_counter.add.side_effect = error_type("telemetry gone")
mock_telemetry = Mock(spec=Telemetry)
mock_telemetry.gauge = Mock(return_value=mock_gauge)
mock_telemetry.counter = Mock(return_value=mock_counter)

provider = MockBaseStorageProvider(
base_path="bucket",
provider_name="mock",
config_dict=_sync_metrics_config(),
telemetry_provider=lambda: mock_telemetry,
)
provider._init_metrics()
assert provider._async_metrics_enabled is False

# The storage operation must return normally despite the record failure.
assert provider.get_object("file.txt") == b""
assert provider._metrics_disabled is True
# The failure is logged (the "logged warning" half of the best-effort contract).
assert [record for record in caplog.records if "disabling metrics" in record.getMessage().lower()]

# Metrics latch off after the failure; subsequent operations silently no-op and do not re-log.
calls_after_first = mock_counter.add.call_count + mock_gauge.set.call_count
assert provider.get_object("file.txt") == b""
assert mock_counter.add.call_count + mock_gauge.set.call_count == calls_after_first
assert len([record for record in caplog.records if "disabling metrics" in record.getMessage().lower()]) == 1


def test_async_rust_transfer_metrics_error_does_not_fail_transfer():
"""A failing sync record in the async Rust download finally must not fail the transfer."""
mock_gauge = Mock()
mock_gauge.set.side_effect = BrokenPipeError("telemetry gone")
mock_counter = Mock()
mock_counter.add.side_effect = BrokenPipeError("telemetry gone")
mock_telemetry = Mock(spec=Telemetry)
mock_telemetry.gauge = Mock(return_value=mock_gauge)
mock_telemetry.counter = Mock(return_value=mock_counter)

provider = MockBaseStorageProvider(
base_path="bucket",
provider_name="mock",
config_dict=_sync_metrics_config(),
telemetry_provider=lambda: mock_telemetry,
)

mock_rust_client = MagicMock()
mock_rust_client.download = AsyncMock(return_value=100)
provider._rust_client = mock_rust_client

metadata = [ObjectMetadata(key="remote-a", content_length=1, last_modified=datetime.now())]
with patch("multistorageclient.providers.base.safe_makedirs"):
# Must not raise BatchTransferError from the telemetry record failure.
provider.download_files(remote_paths=["remote-a"], local_paths=["/tmp/local-a"], metadata=metadata)

assert mock_rust_client.download.await_count == 1
assert provider._metrics_disabled is True


def test_manager_mode_disabled_metrics_do_not_crash_record():
"""A disabled meter returned through a real manager proxy is an auto-proxy wrapping ``None``.
Recording through it must not crash the storage operation: the record-time latch disables
metrics after the first failure."""
# No exporter configured -> meter_provider returns None -> disabled instruments.
config = {"opentelemetry": {"metrics": {"reader": {"options": {}}}}}

manager = TelemetryManager(address=("127.0.0.1", 0), ctx=multiprocessing.get_context("spawn"))
manager.start()
try:
telemetry_proxy = manager.Telemetry() # pyright: ignore [reportAttributeAccessIssue]

provider = MockBaseStorageProvider(
base_path="bucket",
provider_name="mock",
config_dict=config,
telemetry_provider=lambda: telemetry_proxy,
)
provider._init_metrics()

# The storage operation must not crash on record; the first failing record latches metrics off.
assert provider.get_object("file.txt") == b""
assert provider._metrics_disabled is True
finally:
manager.shutdown()


@pytest.mark.parametrize("env_var", ["OTEL_SDK_DISABLED", "MSC_TELEMETRY_DISABLED"])
def test_init_metrics_skipped_when_disabled_via_env(monkeypatch, env_var):
"""An explicit kill-switch short-circuits metrics init before any telemetry work."""
monkeypatch.setenv(env_var, "true")

provider_calls: list[bool] = []

def _telemetry_provider() -> Telemetry:
provider_calls.append(True)
return Mock(spec=Telemetry)

provider = MockBaseStorageProvider(
base_path="bucket",
provider_name="mock",
config_dict=_sync_metrics_config(),
telemetry_provider=_telemetry_provider,
)
provider._init_metrics()

# No exporter/manager/proxy work happened.
assert provider_calls == []
assert provider._metric_gauges == {}
assert provider._metric_counters == {}
assert provider._async_metrics_enabled is False

# The storage operation still succeeds with metrics disabled.
assert provider.get_object("file.txt") == b""
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,37 @@ def test_exporter_does_not_modify_original_config(self, mock_vault_provider):
assert "client_key_file" not in exporter_config
assert "certificate_file" not in exporter_config

def test_exporter_raises_clear_error_when_client_certificate_unsupported(self, mock_vault_provider):
"""An OTLP exporter lacking client_certificate_file support must raise a clear
RuntimeError, not an opaque TypeError, into the caller."""

# Models opentelemetry-exporter-otlp-proto-http < 1.32.1: no mTLS client cert options.
def _legacy_init(self, endpoint=None, certificate_file=None, headers=None, timeout=None, compression=None):
pass

with patch(
"multistorageclient.telemetry.metrics.exporters.otlp_mtls_vault.VaultCertificateProvider",
return_value=mock_vault_provider,
):
with patch(
"opentelemetry.exporter.otlp.proto.http.metric_exporter.OTLPMetricExporter.__init__",
_legacy_init,
):
from multistorageclient.telemetry.metrics.exporters.otlp_mtls_vault import (
_OTLPmTLSVaultMetricExporter,
)

auth_config = {
"vault_endpoint": "https://vault.example.com",
"vault_namespace": "test-namespace",
"approle_id": "test-role-id",
"approle_secret": "test-secret-id",
}
exporter_config = {"endpoint": "https://otlp.example.com/v1/metrics"}

with pytest.raises(RuntimeError, match="client_certificate_file"):
_OTLPmTLSVaultMetricExporter(auth=auth_config, exporter=exporter_config)

def test_vault_provider_receives_auth_config(self, mock_vault_provider):
"""Test that VaultCertificateProvider is initialized with auth config."""
with patch(
Expand Down
Loading