11from __future__ import annotations
22
3+ import os
34from unittest .mock import MagicMock , patch
45
56from pyspark .sql import DataFrame , SparkSession
@@ -200,11 +201,32 @@ def test_read_dlo_deltas(self, reset_client, mock_spark):
200201 reader .read_dlo_deltas .return_value = mock_df
201202
202203 client = Client (reader = reader , writer = writer )
203- result = client .read_dlo_deltas ("test_dlo" )
204+ with patch .dict ("os.environ" , {}, clear = False ):
205+ os .environ .pop ("BYOC_STREAMING_SOURCE_NAME" , None )
206+ result = client .read_dlo_deltas ()
204207
205- reader .read_dlo_deltas .assert_called_once_with ("test_dlo" )
208+ reader .read_dlo_deltas .assert_called_once_with ()
206209 assert result is mock_df
207- assert "test_dlo" in client ._data_layer_history [DataCloudObjectType .DLO ]
210+ assert (
211+ "<streaming delta source>"
212+ in client ._data_layer_history [DataCloudObjectType .DLO ]
213+ )
214+
215+ def test_read_dlo_deltas_records_runtime_source_name (
216+ self , reset_client , mock_spark
217+ ):
218+ """The runtime source env var populates the access-history entry."""
219+ reader = MagicMock (spec = BaseDataCloudReader )
220+ writer = MagicMock (spec = BaseDataCloudWriter )
221+ reader .read_dlo_deltas .return_value = MagicMock (spec = DataFrame )
222+
223+ client = Client (reader = reader , writer = writer )
224+ with patch .dict (
225+ "os.environ" , {"BYOC_STREAMING_SOURCE_NAME" : "Account_std__dll" }
226+ ):
227+ client .read_dlo_deltas ()
228+
229+ assert "Account_std__dll" in client ._data_layer_history [DataCloudObjectType .DLO ]
208230
209231 def test_read_dmo_deltas (self , reset_client , mock_spark ):
210232 reader = MagicMock (spec = BaseDataCloudReader )
@@ -213,11 +235,16 @@ def test_read_dmo_deltas(self, reset_client, mock_spark):
213235 reader .read_dmo_deltas .return_value = mock_df
214236
215237 client = Client (reader = reader , writer = writer )
216- result = client .read_dmo_deltas ("test_dmo" )
238+ with patch .dict (
239+ "os.environ" , {"BYOC_STREAMING_SOURCE_NAME" : "Account_model__dlm" }
240+ ):
241+ result = client .read_dmo_deltas ()
217242
218- reader .read_dmo_deltas .assert_called_once_with ("test_dmo" )
243+ reader .read_dmo_deltas .assert_called_once_with ()
219244 assert result is mock_df
220- assert "test_dmo" in client ._data_layer_history [DataCloudObjectType .DMO ]
245+ assert (
246+ "Account_model__dlm" in client ._data_layer_history [DataCloudObjectType .DMO ]
247+ )
221248
222249 def test_write_dlo_deltas (self , reset_client , mock_spark ):
223250 reader = MagicMock (spec = BaseDataCloudReader )
@@ -262,10 +289,11 @@ def test_streaming_read_write_flow(self, reset_client, mock_spark):
262289
263290 client = Client (reader = reader , writer = writer )
264291
265- df = client .read_dlo_deltas ("source_dll" )
292+ with patch .dict ("os.environ" , {"BYOC_STREAMING_SOURCE_NAME" : "source_dll" }):
293+ df = client .read_dlo_deltas ()
266294 client .write_dlo_deltas ("target_dll" , df )
267295
268- reader .read_dlo_deltas .assert_called_once_with ("source_dll" )
296+ reader .read_dlo_deltas .assert_called_once_with ()
269297 writer .write_dlo_deltas .assert_called_once_with ("target_dll" , stream_df )
270298 assert "source_dll" in client ._data_layer_history [DataCloudObjectType .DLO ]
271299
0 commit comments