Skip to content
Merged
Show file tree
Hide file tree
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
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ def get_storage_type(self) -> StorageType:
def create_object_store(self, path: str, connection_config: ConnectionConfig | None = None):
"""Create an S3 object store using DataFusion's AmazonS3."""
if connection_config is None:
raise ValueError("connection_config must be provided for %s", self.get_storage_type)
raise ValueError(f"connection_config must be provided for {self.get_storage_type}")

try:
credentials = connection_config.credentials
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -156,3 +156,10 @@ def test_get_object_storage_provider(self):

with pytest.raises(ValueError, match="Unsupported storage type"):
get_object_storage_provider("invalid")

def test_s3_provider_requires_connection_config(self):
"""The message names the storage type rather than rendering as a tuple of format args."""
provider = S3ObjectStorageProvider()

with pytest.raises(ValueError, match="connection_config must be provided for s3"):
provider.create_object_store("s3://demo-data/path", connection_config=None)
Original file line number Diff line number Diff line change
Expand Up @@ -253,7 +253,7 @@ def execute(self, context: Context):
self.conf = json.loads(self.conf)
json.dumps(self.conf)
except (TypeError, JSONDecodeError):
raise ValueError("conf parameter should be JSON Serializable %s", self.conf)
raise ValueError(f"conf parameter should be JSON Serializable: {self.conf}")

if self.openlineage_inject_parent_info:
self.log.debug("Checking if OpenLineage information can be safely injected into dagrun conf.")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -144,5 +144,5 @@ def execute(self, context: Context) -> None:

def execute_complete(self, context: Context, event: bool | None = None) -> None:
if not event:
raise AirflowException("%s task failed as %s not found.", self.task_id, self.filepath)
raise AirflowException(f"{self.task_id} task failed as {self.filepath} not found.")
self.log.info("%s completed successfully as %s found.", self.task_id, self.filepath)
Original file line number Diff line number Diff line change
Expand Up @@ -277,7 +277,7 @@ def test_trigger_dagrun_operator_templated_invalid_conf(self, dag_maker):
dag_maker.sync_dagbag_to_db()
parse_and_sync_to_db(self.f_name)
dr = dag_maker.create_dagrun()
with pytest.raises(ValueError, match="conf parameter should be JSON Serializable"):
with pytest.raises(ValueError, match="conf parameter should be JSON Serializable: "):
dag_maker.run_ti(task.task_id, dr)

def test_trigger_dagrun_with_no_failed_state(self, dag_maker):
Expand Down Expand Up @@ -455,7 +455,7 @@ def test_trigger_dagrun_with_str_conf_error(self):
conf="{'foo': 'bar', 'key': 123}",
)

with pytest.raises(ValueError, match="conf parameter should be JSON Serializable"):
with pytest.raises(ValueError, match="conf parameter should be JSON Serializable: "):
task.execute(context={})

@pytest.mark.parametrize("original_conf", (None, {}, {"foo": "bar"}))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,12 @@
import pytest

from airflow.models.dag import DAG
from airflow.providers.common.compat.sdk import AirflowSensorTimeout, TaskDeferred, timezone
from airflow.providers.common.compat.sdk import (
AirflowException,
AirflowSensorTimeout,
TaskDeferred,
timezone,
)
from airflow.providers.standard.sensors.filesystem import FileSensor
from airflow.providers.standard.triggers.file import FileTrigger

Expand Down Expand Up @@ -281,3 +286,10 @@ def test_non_deferrable_sensor_does_not_defer_after_the_sync_path(self):
task.execute({})

assert mock_poke.call_count == 1, "the sensor poked again after the sync path completed"

def test_execute_complete_failure_names_the_task_and_path(self):
"""The message interpolates the task and path rather than rendering as a tuple of format args."""
sensor = FileSensor(task_id="waiting_for_drop", filepath="incoming_data.csv")

with pytest.raises(AirflowException, match="waiting_for_drop task failed as .*incoming_data.csv"):
sensor.execute_complete(context={}, event=False)
Loading