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
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,7 @@ def _transmit(self, envelopes: List[TelemetryItem]) -> ExportResult:
self.storage.put(envelopes_to_store, 0)
self._consecutive_redirects = 0
return ExportResult.FAILED_RETRYABLE

return ExportResult.FAILED_NOT_RETRYABLE
except HttpResponseError as response_error:
if _is_retryable_code(response_error.status_code):
return ExportResult.FAILED_RETRYABLE
Expand Down Expand Up @@ -183,7 +183,6 @@ def _transmit(self, envelopes: List[TelemetryItem]) -> ExportResult:
"Envelopes could not be exported and are not retryable: %s.", ex
)
return ExportResult.FAILED_NOT_RETRYABLE
return ExportResult.FAILED_NOT_RETRYABLE
# No spans to export
self._consecutive_redirects = 0
return ExportResult.SUCCESS
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,13 @@
BaseExporter,
ExportResult,
)
from azure.monitor.opentelemetry.exporter._storage import LocalFileStorage
from azure.monitor.opentelemetry.exporter._generated import AzureMonitorClient
from azure.monitor.opentelemetry.exporter._generated.models import TelemetryItem, TrackResponse
from azure.monitor.opentelemetry.exporter._generated.models import (
TelemetryErrorDetails,
TelemetryItem,
TrackResponse,
)


def throw(exc_type, *args, **kwargs):
Expand Down Expand Up @@ -68,6 +73,7 @@ def test_constructor(self):
storage_maintenance_period=30,
storage_max_size=1000,
storage_min_retry_interval=100,
storage_path="test/path",
storage_retention_period=2000,
)
self.assertEqual(
Expand All @@ -81,49 +87,60 @@ def test_constructor(self):
self.assertEqual(base._timeout, 10)
self.assertEqual(base._api_version, "2021-02-10_Preview")
self.assertEqual(base._storage_min_retry_interval, 100)


@unittest.skip("transient storage")
def test_transmit_from_storage_failed_retryable(self):
envelopes_to_store = [x.as_dict() for x in self._envelopes_to_export]
self._base.storage.put(envelopes_to_store)
# Timeout in HTTP request is a retryable case
with mock.patch("requests.Session.request", throw(requests.Timeout)):
self._base._transmit_from_storage()
# File would be locked for 1 second
self.assertIsNone(self._base.storage.get())
# File still present
self.assertGreaterEqual(len(os.listdir(self._base.storage._path)), 1)

@unittest.skip("transient storage")
def test_transmit_from_storage_failed_not_retryable(self):
envelopes_to_store = [x.as_dict() for x in self._envelopes_to_export]
self._base.storage.put(envelopes_to_store)
with mock.patch("requests.Session.request") as post:
# Do not retry with internal server error responses
post.return_value = MockResponse(400, "{}")
self._base._transmit_from_storage()
self._base.storage.get()
# File no longer present
self.assertEqual(len(os.listdir(self._base.storage._path)), 0)
self.assertEqual(base._storage_path, "test/path")

def test_transmit_from_storage_success(self):
exporter = BaseExporter()
exporter.storage = mock.Mock()
blob_mock = mock.Mock()
blob_mock.lease.return_value = True
envelope_mock = {"name":"test","time":"time"}
blob_mock.get.return_value = [envelope_mock]
exporter.storage.gets.return_value = [blob_mock]
with mock.patch.object(AzureMonitorClient, 'track') as post:
post.return_value = TrackResponse(
items_received=1,
items_accepted=1,
errors=[],
)
Comment thread
lzchen marked this conversation as resolved.
exporter._transmit_from_storage()
exporter.storage.gets.assert_called_once()
blob_mock.lease.assert_called_once()
blob_mock.delete.assert_called_once()

def test_transmit_from_storage_store_again(self):
exporter = BaseExporter()
exporter.storage = mock.Mock()
blob_mock = mock.Mock()
blob_mock.lease.return_value = True
envelope_mock = {"name":"test","time":"time"}
blob_mock.get.return_value = [envelope_mock]
exporter.storage.gets.return_value = [blob_mock]
with mock.patch("azure.monitor.opentelemetry.exporter.export._base._is_retryable_code"):
with mock.patch.object(AzureMonitorClient, 'track', throw(HttpResponseError)):
Comment thread
lzchen marked this conversation as resolved.
exporter._transmit_from_storage()
exporter.storage.gets.assert_called_once()
blob_mock.lease.assert_called()
blob_mock.delete.assert_not_called()

def test_transmit_from_storage_nothing(self):
with mock.patch("requests.Session.request") as post:
post.return_value = None
self._base._transmit_from_storage()

@unittest.skip("transient storage")
@mock.patch("requests.Session.request", return_value=mock.Mock())
def test_transmit_from_storage_lease_failure(self, requests_mock):
requests_mock.return_value = MockResponse(200, "unknown")
envelopes_to_store = [x.as_dict() for x in self._envelopes_to_export]
self._base.storage.put(envelopes_to_store)
with mock.patch(
"azure.monitor.opentelemetry.exporter._storage.LocalFileBlob.lease"
) as lease: # noqa: E501
lease.return_value = False
self._base._transmit_from_storage()
self.assertTrue(self._base.storage.get())
def test_transmit_from_storage_lease_failure(self):
exporter = BaseExporter()
exporter.storage = mock.Mock()
blob_mock = mock.Mock()
blob_mock.lease.return_value = False
exporter.storage.gets.return_value = [blob_mock]
transmit_mock = mock.Mock()
exporter._transmit = transmit_mock
exporter._transmit_from_storage()
exporter.storage.gets.assert_called_once()
transmit_mock.assert_not_called()
blob_mock.lease.assert_called_once()
blob_mock.delete.assert_not_called()

def test_transmit_http_error_retryable(self):
with mock.patch("azure.monitor.opentelemetry.exporter.export._base._is_retryable_code") as m:
Expand Down Expand Up @@ -176,77 +193,54 @@ def test_transmission_200(self):
result = self._base._transmit(self._envelopes_to_export)
self.assertEqual(result, ExportResult.SUCCESS)

def test_transmission_206(self):
with mock.patch("requests.Session.request") as post:
post.return_value = MockResponse(206, "unknown")
result = self._base._transmit(self._envelopes_to_export)
self.assertEqual(result, ExportResult.FAILED_RETRYABLE)

def test_transmission_206_500(self):
def test_transmission_206_retry(self):
exporter = BaseExporter()
exporter.storage = mock.Mock()
test_envelope = TelemetryItem(name="testEnvelope", time=datetime.now())
custom_envelopes_to_export = [TelemetryItem(name="Test", time=datetime.now(
)), TelemetryItem(name="Test", time=datetime.now()), test_envelope]
with mock.patch("requests.Session.request") as post:
post.return_value = MockResponse(
206,
json.dumps(
{
"itemsReceived": 5,
"itemsAccepted": 3,
"errors": [
{"index": 0, "statusCode": 400, "message": ""},
{
"index": 2,
"statusCode": 500,
"message": "Internal Server Error",
},
],
}
),
with mock.patch.object(AzureMonitorClient, 'track') as post:
post.return_value = TrackResponse(
items_received=3,
items_accepted=1,
errors=[
TelemetryErrorDetails(
index=0,
status_code=400,
message="should drop",
),
TelemetryErrorDetails(
index=2,
status_code=500,
message="should retry"
)
],
)
result = self._base._transmit(custom_envelopes_to_export)
result = exporter._transmit(custom_envelopes_to_export)
self.assertEqual(result, ExportResult.FAILED_RETRYABLE)
self.assertEqual(
self._base.storage.get().get()[0]["name"], "testEnvelope"
)
exporter.storage.put.assert_called_once()

def test_transmission_206_no_retry(self):
envelopes_to_export = map(lambda x: x.as_dict(), tuple(
[TelemetryItem(name="testEnvelope", time=datetime.now())]))
self._base.storage.put(envelopes_to_export)
with mock.patch("requests.Session.request") as post:
post.return_value = MockResponse(
206,
json.dumps(
{
"itemsReceived": 3,
"itemsAccepted": 2,
"errors": [
{"index": 0, "statusCode": 400, "message": ""}
],
}
),
)
result = self._base._transmit(self._envelopes_to_export)
self.assertEqual(result, ExportResult.FAILED_NOT_RETRYABLE)

def test_transmission_206_bogus(self):
envelopes_to_export = map(lambda x: x.as_dict(), tuple(
[TelemetryItem(name="testEnvelope", time=datetime.now())]))
self._base.storage.put(envelopes_to_export)
with mock.patch("requests.Session.request") as post:
post.return_value = MockResponse(
206,
json.dumps(
{
"itemsReceived": 5,
"itemsAccepted": 3,
"errors": [{"foo": 0, "bar": 1}],
}
),
exporter = BaseExporter()
exporter.storage = mock.Mock()
test_envelope = TelemetryItem(name="testEnvelope", time=datetime.now())
custom_envelopes_to_export = [TelemetryItem(name="Test", time=datetime.now(
)), TelemetryItem(name="Test", time=datetime.now()), test_envelope]
with mock.patch.object(AzureMonitorClient, 'track') as post:
post.return_value = TrackResponse(
items_received=3,
items_accepted=2,
errors=[
TelemetryErrorDetails(
index=0,
status_code=400,
message="should drop",
),
],
)
result = self._base._transmit(self._envelopes_to_export)
result = self._base._transmit(custom_envelopes_to_export)
self.assertEqual(result, ExportResult.FAILED_NOT_RETRYABLE)
exporter.storage.put.assert_not_called()

def test_transmission_400(self):
with mock.patch("requests.Session.request") as post:
Expand Down