From 0bc5a2d9b3c285f774331abab7af1b7312e36fe0 Mon Sep 17 00:00:00 2001 From: Karthikeyan Singaravelan Date: Tue, 1 Nov 2022 09:51:29 +0530 Subject: [PATCH 1/4] Handle exception during serializing incompatible executor_config objects. --- .../schemas/task_instance_schema.py | 11 +++++ .../endpoints/test_task_instance_endpoint.py | 40 +++++++++++++++++++ 2 files changed, 51 insertions(+) diff --git a/airflow/api_connexion/schemas/task_instance_schema.py b/airflow/api_connexion/schemas/task_instance_schema.py index 970ef9a0fd4d6..bd2ed427efe64 100644 --- a/airflow/api_connexion/schemas/task_instance_schema.py +++ b/airflow/api_connexion/schemas/task_instance_schema.py @@ -33,6 +33,17 @@ from airflow.utils.state import State +class _ExecutorConfigField(fields.String): + def _super_serialize(self, value, attr, obj): + return super()._serialize(value, attr, obj) + + def _serialize(self, value, attr, obj, **kwargs): + try: + return self._super_serialize(value, attr, obj) + except Exception: + return "{}" + + class TaskInstanceSchema(SQLAlchemySchema): """Task instance schema.""" diff --git a/tests/api_connexion/endpoints/test_task_instance_endpoint.py b/tests/api_connexion/endpoints/test_task_instance_endpoint.py index 03819ffaa818f..a9b6c4b0180cc 100644 --- a/tests/api_connexion/endpoints/test_task_instance_endpoint.py +++ b/tests/api_connexion/endpoints/test_task_instance_endpoint.py @@ -220,6 +220,46 @@ def test_should_respond_200(self, username, session): "triggerer_job": None, } + @provide_session + @mock.patch( + "airflow.api_connexion.schemas.task_instance_schema.ExecutorConfigField._super_serialize", + side_effect=AttributeError(), + ) + def test_task_instance_executor_config_exception(self, mock_executor_config, session): + self.create_task_instances(session) + response = self.client.get( + "/api/v1/dags/example_python_operator/dagRuns/TEST_DAG_RUN_ID/taskInstances/print_the_context", + environ_overrides={"REMOTE_USER": "test"}, + ) + assert response.status_code == 200 + assert response.json == { + "dag_id": "example_python_operator", + "duration": 10000.0, + "end_date": "2020-01-03T00:00:00+00:00", + "execution_date": "2020-01-01T00:00:00+00:00", + "executor_config": "{}", + "hostname": "", + "map_index": -1, + "max_tries": 0, + "operator": "_PythonDecoratedOperator", + "pid": 100, + "pool": "default_pool", + "pool_slots": 1, + "priority_weight": 9, + "queue": "default_queue", + "queued_when": None, + "sla_miss": None, + "start_date": "2020-01-02T00:00:00+00:00", + "state": "running", + "task_id": "print_the_context", + "try_number": 0, + "unixname": getuser(), + "dag_run_id": "TEST_DAG_RUN_ID", + "rendered_fields": {}, + "trigger": None, + "triggerer_job": None, + } + def test_should_respond_200_with_task_state_in_deferred(self, session): now = pendulum.now("UTC") ti = self.create_task_instances( From 4c0b5efc0d44b81ef775e15a26806bba73614712 Mon Sep 17 00:00:00 2001 From: Karthikeyan Singaravelan Date: Thu, 3 Nov 2022 12:23:51 +0530 Subject: [PATCH 2/4] Update tests to patch executor_config instead of serialize method. --- .../api_connexion/schemas/task_instance_schema.py | 5 +---- .../endpoints/test_task_instance_endpoint.py | 12 +++++------- 2 files changed, 6 insertions(+), 11 deletions(-) diff --git a/airflow/api_connexion/schemas/task_instance_schema.py b/airflow/api_connexion/schemas/task_instance_schema.py index bd2ed427efe64..846d63caaaf5d 100644 --- a/airflow/api_connexion/schemas/task_instance_schema.py +++ b/airflow/api_connexion/schemas/task_instance_schema.py @@ -34,12 +34,9 @@ class _ExecutorConfigField(fields.String): - def _super_serialize(self, value, attr, obj): - return super()._serialize(value, attr, obj) - def _serialize(self, value, attr, obj, **kwargs): try: - return self._super_serialize(value, attr, obj) + return super().serialize(value, attr, obj, **kwargs) except Exception: return "{}" diff --git a/tests/api_connexion/endpoints/test_task_instance_endpoint.py b/tests/api_connexion/endpoints/test_task_instance_endpoint.py index a9b6c4b0180cc..3bd847caf2137 100644 --- a/tests/api_connexion/endpoints/test_task_instance_endpoint.py +++ b/tests/api_connexion/endpoints/test_task_instance_endpoint.py @@ -221,12 +221,10 @@ def test_should_respond_200(self, username, session): } @provide_session - @mock.patch( - "airflow.api_connexion.schemas.task_instance_schema.ExecutorConfigField._super_serialize", - side_effect=AttributeError(), - ) - def test_task_instance_executor_config_exception(self, mock_executor_config, session): - self.create_task_instances(session) + def test_task_instance_executor_config_exception(self, session): + tis = self.create_task_instances(session) + print_the_context = next(ti for ti in tis if ti.task_id == "print_the_context") + print_the_context.executor_config = Exception() response = self.client.get( "/api/v1/dags/example_python_operator/dagRuns/TEST_DAG_RUN_ID/taskInstances/print_the_context", environ_overrides={"REMOTE_USER": "test"}, @@ -245,7 +243,7 @@ def test_task_instance_executor_config_exception(self, mock_executor_config, ses "pid": 100, "pool": "default_pool", "pool_slots": 1, - "priority_weight": 9, + "priority_weight": 11, "queue": "default_queue", "queued_when": None, "sla_miss": None, From d53396e775bdd9c0d48968b82b5c84dd0e30938c Mon Sep 17 00:00:00 2001 From: Karthikeyan Singaravelan Date: Thu, 3 Nov 2022 12:57:58 +0530 Subject: [PATCH 3/4] Fix tests and typo in method call. --- .../schemas/task_instance_schema.py | 2 +- .../endpoints/test_task_instance_endpoint.py | 17 ++++++++++++++--- 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/airflow/api_connexion/schemas/task_instance_schema.py b/airflow/api_connexion/schemas/task_instance_schema.py index 846d63caaaf5d..02be499d5b2b7 100644 --- a/airflow/api_connexion/schemas/task_instance_schema.py +++ b/airflow/api_connexion/schemas/task_instance_schema.py @@ -36,7 +36,7 @@ class _ExecutorConfigField(fields.String): def _serialize(self, value, attr, obj, **kwargs): try: - return super().serialize(value, attr, obj, **kwargs) + return super()._serialize(value, attr, obj, **kwargs) except Exception: return "{}" diff --git a/tests/api_connexion/endpoints/test_task_instance_endpoint.py b/tests/api_connexion/endpoints/test_task_instance_endpoint.py index 3bd847caf2137..65315be532ce3 100644 --- a/tests/api_connexion/endpoints/test_task_instance_endpoint.py +++ b/tests/api_connexion/endpoints/test_task_instance_endpoint.py @@ -40,6 +40,11 @@ DEFAULT_DATETIME_STR_2 = "2020-01-02T00:00:00+00:00" +class _BadExecutorConfig: + def __str__(self): + raise Exception() + + @pytest.fixture(scope="module") def configured_app(minimal_app_for_api): app = minimal_app_for_api @@ -220,11 +225,17 @@ def test_should_respond_200(self, username, session): "triggerer_job": None, } + @parameterized.expand( + [ + ("Exception serialization", _BadExecutorConfig(), "{}"), + ("Dict serialization", {"foo": "bar"}, "{'foo': 'bar'}"), + ], + ) @provide_session - def test_task_instance_executor_config_exception(self, session): + def test_task_instance_executor_config_exception(self, _, executor_config, expected_config, session): tis = self.create_task_instances(session) print_the_context = next(ti for ti in tis if ti.task_id == "print_the_context") - print_the_context.executor_config = Exception() + print_the_context.executor_config = executor_config response = self.client.get( "/api/v1/dags/example_python_operator/dagRuns/TEST_DAG_RUN_ID/taskInstances/print_the_context", environ_overrides={"REMOTE_USER": "test"}, @@ -235,7 +246,7 @@ def test_task_instance_executor_config_exception(self, session): "duration": 10000.0, "end_date": "2020-01-03T00:00:00+00:00", "execution_date": "2020-01-01T00:00:00+00:00", - "executor_config": "{}", + "executor_config": expected_config, "hostname": "", "map_index": -1, "max_tries": 0, From 39582158e2031f3a929e9655d504f1378f2d8b3d Mon Sep 17 00:00:00 2001 From: Karthikeyan Singaravelan Date: Sat, 7 Jan 2023 09:52:20 +0530 Subject: [PATCH 4/4] Fix tests. --- airflow/api_connexion/schemas/task_instance_schema.py | 2 +- .../endpoints/test_task_instance_endpoint.py | 11 ++++++----- 2 files changed, 7 insertions(+), 6 deletions(-) diff --git a/airflow/api_connexion/schemas/task_instance_schema.py b/airflow/api_connexion/schemas/task_instance_schema.py index 02be499d5b2b7..a1444a897e4a0 100644 --- a/airflow/api_connexion/schemas/task_instance_schema.py +++ b/airflow/api_connexion/schemas/task_instance_schema.py @@ -69,7 +69,7 @@ class Meta: operator = auto_field() queued_dttm = auto_field(data_key="queued_when") pid = auto_field() - executor_config = auto_field() + executor_config = _ExecutorConfigField() note = auto_field() sla_miss = fields.Nested(SlaMissSchema, dump_default=None) rendered_fields = JsonObjectField(dump_default={}) diff --git a/tests/api_connexion/endpoints/test_task_instance_endpoint.py b/tests/api_connexion/endpoints/test_task_instance_endpoint.py index 65315be532ce3..4c20833d11912 100644 --- a/tests/api_connexion/endpoints/test_task_instance_endpoint.py +++ b/tests/api_connexion/endpoints/test_task_instance_endpoint.py @@ -225,14 +225,14 @@ def test_should_respond_200(self, username, session): "triggerer_job": None, } - @parameterized.expand( + @pytest.mark.parametrize( + "executor_config, expected_config", [ - ("Exception serialization", _BadExecutorConfig(), "{}"), - ("Dict serialization", {"foo": "bar"}, "{'foo': 'bar'}"), + pytest.param(_BadExecutorConfig(), "{}", id="Exception serialization"), + pytest.param({"foo": "bar"}, "{'foo': 'bar'}", id="Dict serialization"), ], ) - @provide_session - def test_task_instance_executor_config_exception(self, _, executor_config, expected_config, session): + def test_task_instance_executor_config_exception(self, executor_config, expected_config, session): tis = self.create_task_instances(session) print_the_context = next(ti for ti in tis if ti.task_id == "print_the_context") print_the_context.executor_config = executor_config @@ -251,6 +251,7 @@ def test_task_instance_executor_config_exception(self, _, executor_config, expec "map_index": -1, "max_tries": 0, "operator": "_PythonDecoratedOperator", + "note": "placeholder-note", "pid": 100, "pool": "default_pool", "pool_slots": 1,