Skip to content
Merged
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
22 changes: 12 additions & 10 deletions task-sdk/tests/task_sdk/execution_time/test_task_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -5678,7 +5678,7 @@ def _watcher_side_effect(msg=None, *args, **kwargs):
return AssetResult(name=actual.uri, uri=actual.uri, group="asset")
return OKResponse(ok=True)

def test_asset_state_get_and_set(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_get_and_set(self, create_runtime_ti, mock_supervisor_comms):
watched = Asset(name="my_asset", uri="s3://bucket/data")

class WatcherOperator(BaseOperator):
Expand All @@ -5697,7 +5697,9 @@ def execute(self, context):
)
mock_supervisor_comms.send.assert_any_call(GetAssetStateStoreByName(name="my_asset", key="watermark"))

def test_asset_state_get_returns_default_when_key_missing(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_get_returns_default_when_key_missing(
self, create_runtime_ti, mock_supervisor_comms
):
watched = Asset(name="my_asset", uri="s3://bucket/data")
captured = {}

Expand All @@ -5718,7 +5720,7 @@ def execute(self, context):

assert captured["result"] == "2026-01-01T00:00:00+00:00"

def test_asset_state_delete(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_delete(self, create_runtime_ti, mock_supervisor_comms):
watched = Asset(name="my_asset", uri="s3://bucket/data")

class WatcherOperator(BaseOperator):
Expand All @@ -5735,7 +5737,7 @@ def execute(self, context):
DeleteAssetStateStoreByName(name="my_asset", key="watermark")
)

def test_asset_state_clear(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_clear(self, create_runtime_ti, mock_supervisor_comms):
watched = Asset(name="my_asset", uri="s3://bucket/data")

class WatcherOperator(BaseOperator):
Expand All @@ -5750,7 +5752,7 @@ def execute(self, context):

mock_supervisor_comms.send.assert_any_call(ClearAssetStateStoreByName(name="my_asset"))

def test_asset_state_uri_ref_inlet(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_uri_ref_inlet(self, create_runtime_ti, mock_supervisor_comms):
watched = AssetUriRef(uri="s3://bucket/data")

class WatcherOperator(BaseOperator):
Expand All @@ -5771,7 +5773,7 @@ def execute(self, context):
GetAssetStateStoreByUri(uri="s3://bucket/data", key="watermark")
)

def test_asset_state_alias_as_inlet(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_alias_as_inlet(self, create_runtime_ti, mock_supervisor_comms):
alias = AssetAlias(name="my_alias")
resolved = Asset(name="resolved_asset", uri="s3://bucket/resolved")

Expand All @@ -5796,7 +5798,7 @@ def side_effect(msg):
SetAssetStateStoreByName(name="resolved_asset", key="watermark", value="2026-05-01")
)

def test_asset_state_alias_inlet_no_resolved_assets(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_alias_inlet_no_resolved_assets(self, create_runtime_ti, mock_supervisor_comms):
alias = AssetAlias(name="empty_alias")

class WatcherOperator(BaseOperator):
Expand All @@ -5815,7 +5817,7 @@ def side_effect(msg):

run(runtime_ti, context=runtime_ti.get_template_context(), log=mock.MagicMock())

def test_asset_state_keyed_access_single_inlet(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_keyed_access_single_inlet(self, create_runtime_ti, mock_supervisor_comms):
watched = Asset(name="my_asset", uri="s3://bucket/data")

class WatcherOperator(BaseOperator):
Expand All @@ -5833,7 +5835,7 @@ def execute(self, context):
SetAssetStateStoreByName(name="my_asset", key="watermark", value="2026-05-01")
)

def test_asset_state_multi_inlet(self, create_runtime_ti, mock_supervisor_comms):
def test_asset_state_store_multi_inlet(self, create_runtime_ti, mock_supervisor_comms):
asset_a = Asset(name="asset_a", uri="s3://bucket/a")
asset_b = Asset(name="asset_b", uri="s3://bucket/b")

Expand All @@ -5855,7 +5857,7 @@ def execute(self, context):
SetAssetStateStoreByName(name="asset_b", key="watermark_b", value="2026-05-02")
)

def test_asset_state_set_sends_reference_via_custom_backend(
def test_asset_state_store_set_sends_reference_via_custom_backend(
self, create_runtime_ti, mock_supervisor_comms
):
"""When a worker backend is configured, asset state set() sends a reference, not the actual value."""
Expand Down
Loading