amoghrajesh commented on code in PR #73135:
URL: https://github.com/apache/airflow/pull/73135#discussion_r4062171974
##########
airflow-core/src/airflow/jobs/triggerer_job_runner.py:
##########
@@ -1399,6 +1425,18 @@ async def create_triggers(self):
trigger_instance.task_instance = runtime_ti
else:
trigger_instance.task_instance = ti
+
+ # Pass the task_state_store through to the Trigger so it
can be used in the Trigger
Review Comment:
```suggestion
```
Not needed
##########
airflow-core/tests/unit/jobs/test_triggerer_job.py:
##########
@@ -1033,6 +1043,254 @@ async def
test_create_triggers_asset_state_store_accessor_reads_and_writes(
await runner.cleanup_finished_triggers()
[email protected]
+def make_deferred_trigger():
+ """Factory fixture: call with a list to get a BaseTrigger subclass that
appends each new instance."""
+
+ def factory(injected_instances):
+ class DeferredTrigger(BaseTrigger):
+ def __init__(self, **kwargs):
+ super().__init__(**kwargs)
+ injected_instances.append(self)
+
+ def serialize(self):
+ return (f"{type(self).__module__}.{type(self).__qualname__}",
{})
+
+ async def run(self):
+ yield TriggerEvent("done")
+
+ return DeferredTrigger
+
+ return factory
+
+
+def _ti_dto(map_index=-1):
+ return TaskInstanceDTO(
+ id=uuid.uuid4(),
+ dag_version_id=uuid.uuid4(),
+ task_id="my_task",
+ dag_id="test_dag",
+ run_id="test_run",
+ try_number=1,
+ map_index=map_index,
+ pool_slots=1,
+ queue="default",
+ priority_weight=1,
+ )
+
+
+def test_task_instance_dto_rejects_a_null_map_index():
Review Comment:
This asserts that pydantic rejects `None` for a field this PR does not
touch, so it cannot fail without the PR. It is also in the wrong file: the
model lives in
airflow-`core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py`,
so a future change to `int | None` breaks a triggerer test with no hint about
the triggerer. Drop it, or move it next to the model with a comment saying the
triggerer depends on the field staying non-optional.
##########
airflow-core/tests/unit/triggers/test_base_trigger.py:
##########
@@ -278,6 +278,40 @@ def
test_base_event_trigger_asset_state_store_independent_across_instances():
assert b.asset_state_store is None
+def test_base_trigger_task_state_store_initialized_to_none():
+ """task_state_store is None before the triggerer injects one."""
Review Comment:
```suggestion
```
title is good enough
##########
airflow-core/tests/unit/triggers/test_base_trigger.py:
##########
@@ -278,6 +278,40 @@ def
test_base_event_trigger_asset_state_store_independent_across_instances():
assert b.asset_state_store is None
+def test_base_trigger_task_state_store_initialized_to_none():
+ """task_state_store is None before the triggerer injects one."""
+ trigger = DummyTrigger(name="Dummy Trigger")
+
+ assert trigger.task_state_store is None
+
+
+def test_base_trigger_task_state_store_can_be_set():
+ """task_state_store can be set once the Trigger is initialized."""
Review Comment:
```suggestion
```
title is good enough
##########
airflow-core/tests/unit/jobs/test_triggerer_job.py:
##########
@@ -1033,6 +1043,254 @@ async def
test_create_triggers_asset_state_store_accessor_reads_and_writes(
await runner.cleanup_finished_triggers()
[email protected]
+def make_deferred_trigger():
+ """Factory fixture: call with a list to get a BaseTrigger subclass that
appends each new instance."""
+
+ def factory(injected_instances):
+ class DeferredTrigger(BaseTrigger):
+ def __init__(self, **kwargs):
+ super().__init__(**kwargs)
+ injected_instances.append(self)
+
+ def serialize(self):
+ return (f"{type(self).__module__}.{type(self).__qualname__}",
{})
+
+ async def run(self):
+ yield TriggerEvent("done")
+
+ return DeferredTrigger
+
+ return factory
+
+
+def _ti_dto(map_index=-1):
+ return TaskInstanceDTO(
+ id=uuid.uuid4(),
+ dag_version_id=uuid.uuid4(),
+ task_id="my_task",
+ dag_id="test_dag",
+ run_id="test_run",
+ try_number=1,
+ map_index=map_index,
+ pool_slots=1,
+ queue="default",
+ priority_weight=1,
+ )
+
+
+def test_task_instance_dto_rejects_a_null_map_index():
+ """If the map_index is None, it should raise an error (it should be -1)."""
+ with pytest.raises(ValidationError):
+ _ti_dto(map_index=None)
+
+
[email protected]
[email protected]("map_index", [-1, 3], ids=["unmapped", "mapped"])
+@patch("airflow.jobs.triggerer_job_runner.TriggerRunner.get_trigger_by_classpath")
+async def
test_create_triggers_injects_task_state_store_scoped_to_the_deferring_ti(
+ mock_get_classpath, session, make_deferred_trigger, map_index
+):
+ """task_state_store is populated, and scoped to the task instance that
deferred."""
Review Comment:
The name already says "injects task state store scoped to the deferring ti".
##########
airflow-core/tests/unit/triggers/test_base_trigger.py:
##########
@@ -278,6 +278,40 @@ def
test_base_event_trigger_asset_state_store_independent_across_instances():
assert b.asset_state_store is None
+def test_base_trigger_task_state_store_initialized_to_none():
+ """task_state_store is None before the triggerer injects one."""
+ trigger = DummyTrigger(name="Dummy Trigger")
+
+ assert trigger.task_state_store is None
+
+
+def test_base_trigger_task_state_store_can_be_set():
+ """task_state_store can be set once the Trigger is initialized."""
+ trigger = DummyTrigger(name="Dummy Trigger")
+
+ mock_store = create_autospec(TaskStateStoreAccessor, instance=True)
+ trigger.task_state_store = mock_store
+
+ assert trigger.task_state_store is mock_store
+
+
+def test_base_trigger_task_state_store_independent_across_instances():
+ """a.task_state_store does not impact b.task_state_store."""
Review Comment:
```suggestion
```
title is good enough
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]