|
1 | 1 | import json |
| 2 | +from unittest.mock import MagicMock, patch |
| 3 | + |
| 4 | +import requests |
2 | 5 |
|
3 | 6 | from amiadapters.outputs.base import ExtractOutput |
4 | 7 | from amiadapters.models import DataclassJSONEncoder, GeneralMeter, GeneralMeterRead |
5 | | -from amiadapters.adapters.xylem_datalake import XylemDatalakeAdapter |
| 8 | +from amiadapters.adapters.xylem_datalake import ( |
| 9 | + XylemDatalakeAdapter, |
| 10 | + _create_superset_session_with_retry, |
| 11 | +) |
6 | 12 | from test.base_test_case import BaseTestCase |
7 | 13 |
|
8 | 14 |
|
@@ -235,3 +241,114 @@ def test_parse_flowtime_with_space(self): |
235 | 241 |
|
236 | 242 | result = self.adapter._parse_flowtime("2026-04-14 08:00:00") |
237 | 243 | self.assertEqual(result, datetime(2026, 4, 14, 8, 0, 0)) |
| 244 | + |
| 245 | + |
| 246 | +def _query_response( |
| 247 | + status_code, json_data=None, text="", content_type="application/json" |
| 248 | +): |
| 249 | + resp = MagicMock(spec=requests.Response) |
| 250 | + resp.status_code = status_code |
| 251 | + resp.headers = {"Content-Type": content_type} |
| 252 | + resp.text = text |
| 253 | + resp.json.return_value = json_data if json_data is not None else {} |
| 254 | + if 400 <= status_code < 600: |
| 255 | + resp.raise_for_status.side_effect = requests.exceptions.HTTPError( |
| 256 | + f"{status_code} Error", response=resp |
| 257 | + ) |
| 258 | + else: |
| 259 | + resp.raise_for_status.return_value = None |
| 260 | + return resp |
| 261 | + |
| 262 | + |
| 263 | +class TestXylemDatalakeQueryRetry(BaseTestCase): |
| 264 | + """Tests for retry behavior in _query and the auth retry wrapper.""" |
| 265 | + |
| 266 | + def setUp(self): |
| 267 | + self.adapter = XylemDatalakeAdapter( |
| 268 | + org_id="test_org", |
| 269 | + org_timezone="America/Los_Angeles", |
| 270 | + pipeline_configuration=None, |
| 271 | + configured_task_output_controller=self.TEST_TASK_OUTPUT_CONTROLLER_CONFIGURATION, |
| 272 | + configured_metrics=self.TEST_METRICS_CONFIGURATION, |
| 273 | + agency_code="hlsbo", |
| 274 | + database_id=1, |
| 275 | + client_id="test-client-id", |
| 276 | + username="test", |
| 277 | + password="test", |
| 278 | + configured_sinks=[], |
| 279 | + ) |
| 280 | + self.adapter._session = MagicMock() |
| 281 | + self.adapter._requests_since_csrf = 0 |
| 282 | + |
| 283 | + @patch("amiadapters.adapters.xylem_datalake.time.sleep") |
| 284 | + def test_query_retries_on_5xx_then_succeeds(self, _mock_sleep): |
| 285 | + success = _query_response( |
| 286 | + 200, json_data={"status": "success", "data": [{"x": 1}]} |
| 287 | + ) |
| 288 | + self.adapter._session.post.side_effect = [ |
| 289 | + _query_response(500, text="boom"), |
| 290 | + _query_response(502, text="bad gateway"), |
| 291 | + success, |
| 292 | + ] |
| 293 | + rows = self.adapter._query("SELECT 1") |
| 294 | + self.assertEqual(rows, [{"x": 1}]) |
| 295 | + self.assertEqual(self.adapter._session.post.call_count, 3) |
| 296 | + |
| 297 | + @patch("amiadapters.adapters.xylem_datalake.time.sleep") |
| 298 | + def test_query_raises_after_max_5xx_attempts(self, _mock_sleep): |
| 299 | + self.adapter._session.post.return_value = _query_response(500, text="boom") |
| 300 | + with self.assertRaises(requests.exceptions.HTTPError): |
| 301 | + self.adapter._query("SELECT 1") |
| 302 | + |
| 303 | + @patch("amiadapters.adapters.xylem_datalake.time.sleep") |
| 304 | + def test_query_retries_on_connection_error_then_succeeds(self, _mock_sleep): |
| 305 | + success = _query_response( |
| 306 | + 200, json_data={"status": "success", "data": [{"x": 1}]} |
| 307 | + ) |
| 308 | + self.adapter._session.post.side_effect = [ |
| 309 | + requests.exceptions.ConnectionError("dns failure"), |
| 310 | + requests.exceptions.Timeout("read timeout"), |
| 311 | + success, |
| 312 | + ] |
| 313 | + rows = self.adapter._query("SELECT 1") |
| 314 | + self.assertEqual(rows, [{"x": 1}]) |
| 315 | + self.assertEqual(self.adapter._session.post.call_count, 3) |
| 316 | + |
| 317 | + |
| 318 | +class TestCreateSupersetSessionWithRetry(BaseTestCase): |
| 319 | + """Tests for the auth retry wrapper.""" |
| 320 | + |
| 321 | + @patch("amiadapters.adapters.xylem_datalake.time.sleep") |
| 322 | + @patch("amiadapters.adapters.xylem_datalake._create_superset_session") |
| 323 | + def test_returns_session_on_first_success(self, mock_create, _mock_sleep): |
| 324 | + sentinel = MagicMock(name="session") |
| 325 | + mock_create.return_value = sentinel |
| 326 | + result = _create_superset_session_with_retry( |
| 327 | + "https://dl", "https://sup", "client", "user", "pw" |
| 328 | + ) |
| 329 | + self.assertIs(result, sentinel) |
| 330 | + self.assertEqual(mock_create.call_count, 1) |
| 331 | + |
| 332 | + @patch("amiadapters.adapters.xylem_datalake.time.sleep") |
| 333 | + @patch("amiadapters.adapters.xylem_datalake._create_superset_session") |
| 334 | + def test_retries_then_succeeds(self, mock_create, _mock_sleep): |
| 335 | + sentinel = MagicMock(name="session") |
| 336 | + mock_create.side_effect = [ |
| 337 | + RuntimeError("Superset login returned unexpected shape (status 401)"), |
| 338 | + RuntimeError("Superset login returned unexpected shape (status 401)"), |
| 339 | + sentinel, |
| 340 | + ] |
| 341 | + result = _create_superset_session_with_retry( |
| 342 | + "https://dl", "https://sup", "client", "user", "pw" |
| 343 | + ) |
| 344 | + self.assertIs(result, sentinel) |
| 345 | + self.assertEqual(mock_create.call_count, 3) |
| 346 | + |
| 347 | + @patch("amiadapters.adapters.xylem_datalake.time.sleep") |
| 348 | + @patch("amiadapters.adapters.xylem_datalake._create_superset_session") |
| 349 | + def test_raises_after_max_attempts(self, mock_create, _mock_sleep): |
| 350 | + mock_create.side_effect = RuntimeError("persistent failure") |
| 351 | + with self.assertRaises(RuntimeError): |
| 352 | + _create_superset_session_with_retry( |
| 353 | + "https://dl", "https://sup", "client", "user", "pw" |
| 354 | + ) |
0 commit comments