jason810496 commented on code in PR #69817:
URL: https://github.com/apache/airflow/pull/69817#discussion_r3608441078


##########
providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py:
##########
@@ -52,6 +53,86 @@ def s3mock():
         yield
 
 
+class TestS3RemoteLogIOFromConfig:
+    @conf_vars(
+        {
+            ("logging", "base_log_folder"): "~/airflow/logs",
+            ("logging", "remote_base_log_folder"): 
"s3://bucket/remote/log/location",
+            ("logging", "delete_local_logs"): "True",
+        }
+    )
+    def test_from_config(self):
+        subject = S3RemoteLogIO.from_config()
+
+        assert subject.remote_base == "s3://bucket/remote/log/location"
+        assert subject.base_log_folder == 
pathlib.Path(os.path.expanduser("~/airflow/logs"))
+        assert subject.delete_local_copy is True
+
+    @conf_vars(
+        {
+            ("logging", "base_log_folder"): "/tmp/airflow/logs",
+            ("logging", "remote_base_log_folder"): 
"s3://bucket/remote/log/location",
+            ("logging", "delete_local_logs"): "False",
+            ("logging", "remote_task_handler_kwargs"): '{"delete_local_copy": 
true, "max_bytes": 1024}',
+        }
+    )
+    def 
test_from_config_applies_io_kwargs_and_filters_file_handler_kwargs(self):
+        subject = S3RemoteLogIO.from_config()
+
+        assert subject.delete_local_copy is True
+        assert not hasattr(subject, "max_bytes")
+
+    @conf_vars({("logging", "remote_task_handler_kwargs"): '["not", "a", 
"dict"]'})
+    def test_from_config_rejects_non_dict_remote_task_handler_kwargs(self):
+        with pytest.raises(ValueError, match="remote_task_handler_kwargs"):
+            S3RemoteLogIO.from_config()
+
+    def test_provider_registers_s3_scheme(self):
+        from airflow.providers_manager import ProvidersManager
+
+        manager = ProvidersManager()
+        if not hasattr(manager, "remote_logging_handler_by_scheme"):
+            pytest.skip("Airflow core does not support remote logging provider 
dispatch")
+
+        info = manager.remote_logging_handler_by_scheme("s3")
+
+        assert info is not None
+        assert info.classpath == 
"airflow.providers.amazon.aws.log.s3_task_handler.S3RemoteLogIO"
+
+    @pytest.mark.parametrize(
+        "manager_classpath",
+        [
+            pytest.param("airflow.providers_manager.ProvidersManager", 
id="core"),
+            pytest.param(
+                
"airflow.sdk.providers_manager_runtime.ProvidersManagerTaskRuntime", 
id="task-runtime"
+            ),
+        ],
+    )
+    @conf_vars(
+        {
+            ("logging", "remote_logging"): "True",
+            ("logging", "remote_base_log_folder"): 
"s3://bucket/remote/log/location",
+            ("logging", "remote_log_conn_id"): "aws_default",
+        }
+    )
+    def 
test_resolve_remote_task_log_uses_provider_dispatch_not_local_settings(self, 
manager_classpath):

Review Comment:
   This test is exercising the following remote logging discovery that we 
should go through the `_build_remote_task_log_from_provider` path.
   
   
https://github.com/apache/airflow/blob/ffde96203e76c0c27ccf1ae6eeaad270653248bf/shared/logging/src/airflow_shared/logging/factory.py#L134-L182



##########
providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py:
##########
@@ -46,6 +47,31 @@ class S3RemoteLogIO(LoggingMixin):  # noqa: D101
 
     processors = ()
 
+    @classmethod
+    def from_config(cls) -> S3RemoteLogIO:
+        """Build the remote log IO from Airflow logging configuration."""
+        remote_task_handler_kwargs = conf.getjson("logging", 
"remote_task_handler_kwargs", fallback={})
+        if not isinstance(remote_task_handler_kwargs, dict):
+            raise ValueError(
+                "logging/remote_task_handler_kwargs must be a JSON object (a 
python dict), we got "
+                f"{type(remote_task_handler_kwargs)}"
+            )
+        # remote_task_handler_kwargs mixes FileTaskHandler kwargs with IO 
kwargs; only the
+        # latter belong to this class (same split as 
airflow_local_settings.py).
+        fth_params = 
frozenset(inspect.signature(FileTaskHandler.__init__).parameters) - {
+            "self",
+            "base_log_folder",
+        }
+        io_kwargs = {k: v for k, v in remote_task_handler_kwargs.items() if k 
not in fth_params}
+        return cls(
+            **{
+                "base_log_folder": 
os.path.expanduser(conf.get_mandatory_value("logging", "base_log_folder")),
+                "remote_base": conf.get_mandatory_value("logging", 
"remote_base_log_folder"),
+                "delete_local_copy": conf.getboolean("logging", 
"delete_local_logs"),
+            }
+            | io_kwargs,
+        )

Review Comment:
   The rest of the S3 specific kwargs are coming from:
   
   
https://github.com/apache/airflow/blob/ffde96203e76c0c27ccf1ae6eeaad270653248bf/airflow-core/src/airflow/config_templates/airflow_local_settings.py#L173-L187
   



##########
providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py:
##########
@@ -46,6 +47,31 @@ class S3RemoteLogIO(LoggingMixin):  # noqa: D101
 
     processors = ()
 
+    @classmethod
+    def from_config(cls) -> S3RemoteLogIO:
+        """Build the remote log IO from Airflow logging configuration."""
+        remote_task_handler_kwargs = conf.getjson("logging", 
"remote_task_handler_kwargs", fallback={})
+        if not isinstance(remote_task_handler_kwargs, dict):
+            raise ValueError(
+                "logging/remote_task_handler_kwargs must be a JSON object (a 
python dict), we got "
+                f"{type(remote_task_handler_kwargs)}"
+            )
+        # remote_task_handler_kwargs mixes FileTaskHandler kwargs with IO 
kwargs; only the
+        # latter belong to this class (same split as 
airflow_local_settings.py).
+        fth_params = 
frozenset(inspect.signature(FileTaskHandler.__init__).parameters) - {
+            "self",
+            "base_log_folder",
+        }

Review Comment:
   I intentionally duplicate the logic from core side:
   
   
https://github.com/apache/airflow/blob/ffde96203e76c0c27ccf1ae6eeaad270653248bf/airflow-core/src/airflow/config_templates/airflow_local_settings.py#L157-L170



-- 
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]

Reply via email to