This is an automated email from the ASF dual-hosted git repository.

shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new c1e98dc0c96 Add dedicated InfluxDB 3 sensor (#73522)
c1e98dc0c96 is described below

commit c1e98dc0c96eb9177f6086a4aa551d0200b65377
Author: Subhramit Basu <[email protected]>
AuthorDate: Fri Sep 25 12:41:11 2026 +0530

    Add dedicated InfluxDB 3 sensor (#73522)
---
 providers/influxdb/docs/index.rst                  |   1 +
 providers/influxdb/docs/sensors/index.rst          |  40 +++++
 providers/influxdb/provider.yaml                   |   5 +
 .../providers/influxdb/get_provider_info.py        |   6 +
 .../influxdb/{utils.py => sensors/__init__.py}     |  14 --
 .../providers/influxdb/sensors/influxdb3.py        | 108 +++++++++++++
 .../providers/influxdb/triggers/influxdb3.py       |  52 ++++++-
 .../src/airflow/providers/influxdb/utils.py        |  16 ++
 .../tests/system/influxdb/example_influxdb3.py     |  14 +-
 .../unit/influxdb/sensors/__init__.py}             |  14 --
 .../tests/unit/influxdb/sensors/test_influxdb3.py  | 170 +++++++++++++++++++++
 .../influxdb/tests/unit/influxdb/test_utils.py     |  29 +++-
 .../tests/unit/influxdb/triggers/test_influxdb3.py | 123 ++++++++++++++-
 13 files changed, 560 insertions(+), 32 deletions(-)

diff --git a/providers/influxdb/docs/index.rst 
b/providers/influxdb/docs/index.rst
index dfc9b71ad07..634405ee527 100644
--- a/providers/influxdb/docs/index.rst
+++ b/providers/influxdb/docs/index.rst
@@ -37,6 +37,7 @@
     Connection types (InfluxDB 2.x) <connections/influxdb>
     Connection types (InfluxDB 3.x) <connections/influxdb3>
     Operators <operators/index>
+    Sensors <sensors/index>
 
 .. toctree::
     :hidden:
diff --git a/providers/influxdb/docs/sensors/index.rst 
b/providers/influxdb/docs/sensors/index.rst
new file mode 100644
index 00000000000..dc81fa62663
--- /dev/null
+++ b/providers/influxdb/docs/sensors/index.rst
@@ -0,0 +1,40 @@
+.. Licensed to the Apache Software Foundation (ASF) under one
+   or more contributor license agreements.  See the NOTICE file
+   distributed with this work for additional information
+   regarding copyright ownership.  The ASF licenses this file
+   to you under the Apache License, Version 2.0 (the
+   "License"); you may not use this file except in compliance
+   with the License.  You may obtain a copy of the License at
+
+..   http://www.apache.org/licenses/LICENSE-2.0
+
+.. Unless required by applicable law or agreed to in writing,
+   software distributed under the License is distributed on an
+   "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+   KIND, either express or implied.  See the License for the
+   specific language governing permissions and limitations
+   under the License.
+
+.. _howto/sensor:InfluxDB3Sensor:
+
+InfluxDB3Sensor
+===============
+
+Use :class:`~airflow.providers.influxdb.sensors.influxdb3.InfluxDB3Sensor` to 
wait until an
+InfluxDB 3.x SQL query returns a truthy first cell. Prefer an efficient 
existence query that returns
+one value and limits the result to one row.
+
+.. exampleinclude:: /../../influxdb/tests/system/influxdb/example_influxdb3.py
+    :language: python
+    :start-after: [START howto_sensor_influxdb3]
+    :end-before: [END howto_sensor_influxdb3]
+
+An empty result, a missing value, numeric or string zero, and an empty string 
are treated as false.
+Set ``fail_on_empty=True`` to fail immediately when the query returns no rows.
+
+Deferrable mode
+^^^^^^^^^^^^^^^
+
+Set ``deferrable=True`` to release the worker slot between queries. The
+:class:`~airflow.providers.influxdb.triggers.influxdb3.InfluxDB3SensorTrigger` 
repeats the query
+at the configured ``poke_interval`` until the condition is met or Airflow 
reaches the sensor timeout.
diff --git a/providers/influxdb/provider.yaml b/providers/influxdb/provider.yaml
index b637a1d2f9a..ddc3d971df5 100644
--- a/providers/influxdb/provider.yaml
+++ b/providers/influxdb/provider.yaml
@@ -92,6 +92,11 @@ operators:
     python-modules:
       - airflow.providers.influxdb.operators.influxdb3
 
+sensors:
+  - integration-name: InfluxDB 3
+    python-modules:
+      - airflow.providers.influxdb.sensors.influxdb3
+
 triggers:
   - integration-name: InfluxDB 3
     python-modules:
diff --git 
a/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py 
b/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py
index 0b2011da40f..2fbde52a3f9 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/get_provider_info.py
@@ -56,6 +56,12 @@ def get_provider_info():
                 "python-modules": 
["airflow.providers.influxdb.operators.influxdb3"],
             },
         ],
+        "sensors": [
+            {
+                "integration-name": "InfluxDB 3",
+                "python-modules": 
["airflow.providers.influxdb.sensors.influxdb3"],
+            }
+        ],
         "triggers": [
             {
                 "integration-name": "InfluxDB 3",
diff --git a/providers/influxdb/src/airflow/providers/influxdb/utils.py 
b/providers/influxdb/src/airflow/providers/influxdb/sensors/__init__.py
similarity index 63%
copy from providers/influxdb/src/airflow/providers/influxdb/utils.py
copy to providers/influxdb/src/airflow/providers/influxdb/sensors/__init__.py
index 79abb8d6ee2..13a83393a91 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/utils.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/sensors/__init__.py
@@ -14,17 +14,3 @@
 # KIND, either express or implied.  See the License for the
 # specific language governing permissions and limitations
 # under the License.
-"""Internal helpers shared across the InfluxDB provider."""
-
-from __future__ import annotations
-
-import json
-from typing import TYPE_CHECKING, Any
-
-if TYPE_CHECKING:
-    import pandas as pd
-
-
-def _convert_dataframe_to_records(dataframe: pd.DataFrame) -> list[dict[str, 
Any]]:
-    """Convert a query result DataFrame into a JSON-serializable list of 
dictionaries."""
-    return json.loads(dataframe.to_json(orient="records", date_format="iso"))
diff --git 
a/providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py 
b/providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py
new file mode 100644
index 00000000000..88067a0bb48
--- /dev/null
+++ b/providers/influxdb/src/airflow/providers/influxdb/sensors/influxdb3.py
@@ -0,0 +1,108 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""Sensor waiting for a SQL query to return a truthy first cell in InfluxDB 
3.x."""
+
+from __future__ import annotations
+
+from collections.abc import Sequence
+from datetime import timedelta
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.common.compat.sdk import AirflowFailException, 
BaseSensorOperator, conf
+from airflow.providers.influxdb.hooks.influxdb3 import InfluxDB3Hook
+from airflow.providers.influxdb.triggers.influxdb3 import 
InfluxDB3SensorTrigger
+from airflow.providers.influxdb.utils import _first_cell_is_truthy
+
+if TYPE_CHECKING:
+    from airflow.sdk.definitions.context import Context
+
+
+class InfluxDB3Sensor(BaseSensorOperator):
+    """
+    Wait until an InfluxDB 3.x SQL query returns a truthy first cell.
+
+    .. seealso::
+        For more information on how to use this sensor, take a look at the 
guide:
+        :ref:`howto/sensor:InfluxDB3Sensor`
+
+    :param sql: The SQL query to poll.
+    :param influxdb3_conn_id: Reference to :ref:`InfluxDB 3 connection id 
<howto/connection:influxdb3>`.
+        Defaults to ``influxdb3_default``.
+    :param fail_on_empty: Fail instead of waiting when the query returns no 
rows. Defaults to ``False``.
+    :param deferrable: Run polling in the triggerer. Defaults to the
+        ``operators.default_deferrable`` configuration (``False`` if unset).
+    """
+
+    template_fields: Sequence[str] = ("sql", "influxdb3_conn_id")
+    template_ext: Sequence[str] = (".sql",)
+
+    def __init__(
+        self,
+        *,
+        sql: str,
+        influxdb3_conn_id: str = "influxdb3_default",
+        fail_on_empty: bool = False,
+        deferrable: bool = conf.getboolean("operators", "default_deferrable", 
fallback=False),
+        **kwargs,
+    ) -> None:
+        super().__init__(**kwargs)
+        self.sql = sql
+        self.influxdb3_conn_id = influxdb3_conn_id
+        self.fail_on_empty = fail_on_empty
+        self.deferrable = deferrable
+
+    def poke(self, context: Context) -> bool:
+        """Return whether the query result meets the sensor condition."""
+        self.log.info("Poking with SQL query: %s", self.sql)
+        dataframe = 
InfluxDB3Hook(conn_id=self.influxdb3_conn_id).query(self.sql)
+        if dataframe.empty and self.fail_on_empty:
+            raise AirflowFailException("No rows returned, raising as per 
fail_on_empty flag")
+        return _first_cell_is_truthy(dataframe)
+
+    def execute(self, context: Context) -> None:
+        if not self.deferrable:
+            super().execute(context)
+            return
+
+        if self.poke(context):
+            return
+
+        self.defer(
+            timeout=timedelta(seconds=self.timeout),
+            trigger=InfluxDB3SensorTrigger(
+                sql=self.sql,
+                influxdb3_conn_id=self.influxdb3_conn_id,
+                poll_interval=self.poke_interval,
+                fail_on_empty=self.fail_on_empty,
+            ),
+            method_name="execute_complete",
+        )
+
+    def execute_complete(self, context: Context, event: dict[str, Any] | None 
= None) -> None:
+        """Complete after the trigger reports that the condition was met."""
+        if event is None:
+            raise RuntimeError("InfluxDB 3 sensor did not return an event")
+
+        status = event.get("status")
+        if status == "fail":
+            raise AirflowFailException(event.get("message", "InfluxDB 3 sensor 
failed"))
+        if status == "error":
+            raise RuntimeError(event.get("message", "InfluxDB 3 sensor 
failed"))
+        if status != "success":
+            raise RuntimeError(f"InfluxDB 3 sensor returned unexpected status: 
{status!r}")
+
+        self.log.info("InfluxDB 3 sensor condition met")
diff --git 
a/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py 
b/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py
index f9bf999fca2..6bdeabf43e9 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/triggers/influxdb3.py
@@ -22,7 +22,7 @@ import asyncio
 from typing import TYPE_CHECKING, Any
 
 from airflow.providers.influxdb.hooks.influxdb3 import InfluxDB3Hook
-from airflow.providers.influxdb.utils import _convert_dataframe_to_records
+from airflow.providers.influxdb.utils import _convert_dataframe_to_records, 
_first_cell_is_truthy
 from airflow.triggers.base import BaseTrigger, TriggerEvent
 
 if TYPE_CHECKING:
@@ -76,3 +76,53 @@ class InfluxDB3QueryTrigger(BaseTrigger):
             return
 
         yield TriggerEvent({"status": "success", "records": records})
+
+
+class InfluxDB3SensorTrigger(BaseTrigger):
+    """Poll an InfluxDB 3.x SQL query until its first cell meets the sensor 
condition."""
+
+    def __init__(
+        self,
+        sql: str,
+        influxdb3_conn_id: str = "influxdb3_default",
+        poll_interval: float = 60,
+        fail_on_empty: bool = False,
+    ) -> None:
+        super().__init__()
+        self.sql = sql
+        self.influxdb3_conn_id = influxdb3_conn_id
+        self.poll_interval = poll_interval
+        self.fail_on_empty = fail_on_empty
+
+    def serialize(self) -> tuple[str, dict[str, Any]]:
+        return (
+            
"airflow.providers.influxdb.triggers.influxdb3.InfluxDB3SensorTrigger",
+            {
+                "sql": self.sql,
+                "influxdb3_conn_id": self.influxdb3_conn_id,
+                "poll_interval": self.poll_interval,
+                "fail_on_empty": self.fail_on_empty,
+            },
+        )
+
+    async def run(self) -> AsyncIterator[TriggerEvent]:
+        hook = InfluxDB3Hook(conn_id=self.influxdb3_conn_id)
+        while True:
+            try:
+                dataframe = await hook.query_async(self.sql)
+            except Exception as error:
+                self.log.exception("InfluxDB 3 sensor query failed")
+                yield TriggerEvent({"status": "error", "message": str(error)})
+                return
+
+            if dataframe.empty and self.fail_on_empty:
+                yield TriggerEvent(
+                    {"status": "fail", "message": "No rows returned, raising 
as per fail_on_empty flag"}
+                )
+                return
+
+            if _first_cell_is_truthy(dataframe):
+                yield TriggerEvent({"status": "success"})
+                return
+
+            await asyncio.sleep(self.poll_interval)
diff --git a/providers/influxdb/src/airflow/providers/influxdb/utils.py 
b/providers/influxdb/src/airflow/providers/influxdb/utils.py
index 79abb8d6ee2..6a64526f56d 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/utils.py
+++ b/providers/influxdb/src/airflow/providers/influxdb/utils.py
@@ -28,3 +28,19 @@ if TYPE_CHECKING:
 def _convert_dataframe_to_records(dataframe: pd.DataFrame) -> list[dict[str, 
Any]]:
     """Convert a query result DataFrame into a JSON-serializable list of 
dictionaries."""
     return json.loads(dataframe.to_json(orient="records", date_format="iso"))
+
+
+def _first_cell_is_truthy(dataframe: pd.DataFrame) -> bool:
+    """Return whether the first cell meets the sensor condition."""
+    import pandas as pd
+
+    if dataframe.empty or dataframe.shape[1] == 0:
+        return False
+
+    value = dataframe.iat[0, 0]
+    if not pd.api.types.is_scalar(value):
+        raise TypeError("The first query result cell must be a scalar value")
+    if pd.isna(value):
+        return False
+
+    return value not in (0, "0", "", None)
diff --git a/providers/influxdb/tests/system/influxdb/example_influxdb3.py 
b/providers/influxdb/tests/system/influxdb/example_influxdb3.py
index ef8029e58b9..c443f7df8ce 100644
--- a/providers/influxdb/tests/system/influxdb/example_influxdb3.py
+++ b/providers/influxdb/tests/system/influxdb/example_influxdb3.py
@@ -35,6 +35,7 @@ except ImportError:
 from airflow.models.dag import DAG
 from airflow.providers.influxdb.hooks.influxdb3 import InfluxDB3Hook
 from airflow.providers.influxdb.operators.influxdb3 import InfluxDB3Operator
+from airflow.providers.influxdb.sensors.influxdb3 import InfluxDB3Sensor
 
 
 @task(task_id="write_data")
@@ -66,6 +67,17 @@ deferrable_query_task = InfluxDB3Operator(
 )
 # [END howto_operator_influxdb3_deferrable]
 
+# [START howto_sensor_influxdb3]
+wait_for_data = InfluxDB3Sensor(
+    task_id="wait_for_data",
+    sql="""SELECT 1 FROM "temperature" WHERE time > now() - INTERVAL '1 hour' 
LIMIT 1""",
+    influxdb3_conn_id="influxdb3_default",
+    poke_interval=60,
+    timeout=3600,
+    deferrable=True,
+)
+# [END howto_sensor_influxdb3]
+
 ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID")
 DAG_ID = "influxdb3_example_dag"
 
@@ -77,7 +89,7 @@ with DAG(
     tags=["example", "influxdb3"],
 ) as dag:
     write_task = write_to_influxdb3()
-    write_task >> [query_task, deferrable_query_task]
+    write_task >> wait_for_data >> [query_task, deferrable_query_task]
 
     from tests_common.test_utils.watcher import watcher
 
diff --git a/providers/influxdb/src/airflow/providers/influxdb/utils.py 
b/providers/influxdb/tests/unit/influxdb/sensors/__init__.py
similarity index 63%
copy from providers/influxdb/src/airflow/providers/influxdb/utils.py
copy to providers/influxdb/tests/unit/influxdb/sensors/__init__.py
index 79abb8d6ee2..13a83393a91 100644
--- a/providers/influxdb/src/airflow/providers/influxdb/utils.py
+++ b/providers/influxdb/tests/unit/influxdb/sensors/__init__.py
@@ -14,17 +14,3 @@
 # KIND, either express or implied.  See the License for the
 # specific language governing permissions and limitations
 # under the License.
-"""Internal helpers shared across the InfluxDB provider."""
-
-from __future__ import annotations
-
-import json
-from typing import TYPE_CHECKING, Any
-
-if TYPE_CHECKING:
-    import pandas as pd
-
-
-def _convert_dataframe_to_records(dataframe: pd.DataFrame) -> list[dict[str, 
Any]]:
-    """Convert a query result DataFrame into a JSON-serializable list of 
dictionaries."""
-    return json.loads(dataframe.to_json(orient="records", date_format="iso"))
diff --git a/providers/influxdb/tests/unit/influxdb/sensors/test_influxdb3.py 
b/providers/influxdb/tests/unit/influxdb/sensors/test_influxdb3.py
new file mode 100644
index 00000000000..f87d89d5fa3
--- /dev/null
+++ b/providers/influxdb/tests/unit/influxdb/sensors/test_influxdb3.py
@@ -0,0 +1,170 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+from datetime import timedelta
+from unittest import mock
+
+import pandas as pd
+import pytest
+
+from airflow.providers.common.compat.sdk import AirflowFailException, 
AirflowSensorTimeout, TaskDeferred
+from airflow.providers.influxdb.sensors.influxdb3 import InfluxDB3Sensor
+from airflow.providers.influxdb.triggers.influxdb3 import 
InfluxDB3SensorTrigger
+
+SQL = """SELECT 1 FROM "events" WHERE time > now() - INTERVAL '1 hour' LIMIT 
1"""
+CONN_ID = "test_influxdb3_conn"
+HOOK_PATH = "airflow.providers.influxdb.sensors.influxdb3.InfluxDB3Hook"
+
+
+class TestInfluxDB3Sensor:
+    def test_init(self):
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL)
+
+        assert sensor.sql == SQL
+        assert sensor.influxdb3_conn_id == "influxdb3_default"
+        assert sensor.fail_on_empty is False
+        assert sensor.deferrable is False
+        assert sensor.template_fields == ("sql", "influxdb3_conn_id")
+        assert sensor.template_ext == (".sql",)
+
+    @pytest.mark.parametrize(
+        ("dataframe", "expected"),
+        [
+            pytest.param(pd.DataFrame({"literal": [1]}), True, 
id="numeric-one"),
+            pytest.param(pd.DataFrame({"literal": ["ready"]}), True, 
id="non-empty-string"),
+            pytest.param(pd.DataFrame({"literal": []}), False, id="no-rows"),
+            pytest.param(pd.DataFrame({"count": [0]}), False, 
id="numeric-zero"),
+            pytest.param(pd.DataFrame({"count": ["0"]}), False, 
id="string-zero"),
+            pytest.param(pd.DataFrame({"value": [None]}), False, id="none"),
+        ],
+    )
+    @mock.patch(HOOK_PATH, autospec=True)
+    def test_poke(self, mock_hook_class, dataframe, expected):
+        mock_hook_class.return_value.query.return_value = dataframe
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, 
influxdb3_conn_id=CONN_ID)
+
+        assert sensor.poke(context={}) is expected
+        mock_hook_class.assert_called_once_with(conn_id=CONN_ID)
+        mock_hook_class.return_value.query.assert_called_once_with(SQL)
+
+    @mock.patch(HOOK_PATH, autospec=True)
+    def test_poke_fail_on_empty(self, mock_hook_class):
+        mock_hook_class.return_value.query.return_value = 
pd.DataFrame({"literal": []})
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, fail_on_empty=True)
+
+        with pytest.raises(AirflowFailException, match="fail_on_empty"):
+            sensor.poke(context={})
+
+    @mock.patch.object(InfluxDB3Sensor, "defer", autospec=True)
+    @mock.patch(HOOK_PATH, autospec=True)
+    def test_execute_deferrable_fail_on_empty_before_deferral(self, 
mock_hook_class, mock_defer):
+        mock_hook_class.return_value.query.return_value = 
pd.DataFrame({"literal": []})
+        sensor = InfluxDB3Sensor(
+            task_id="wait",
+            sql=SQL,
+            fail_on_empty=True,
+            deferrable=True,
+        )
+
+        with pytest.raises(AirflowFailException, match="fail_on_empty"):
+            sensor.execute(context={})
+
+        mock_hook_class.return_value.query.assert_called_once_with(SQL)
+        mock_defer.assert_not_called()
+
+    @mock.patch(HOOK_PATH, autospec=True)
+    def test_execute_times_out_when_condition_is_never_met(self, 
mock_hook_class):
+        mock_hook_class.return_value.query.return_value = 
pd.DataFrame({"literal": []})
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, poke_interval=0, 
timeout=0)
+
+        with pytest.raises(AirflowSensorTimeout):
+            sensor.execute(context={})
+
+    @mock.patch(HOOK_PATH, autospec=True)
+    def test_execute_deferrable_completes_after_initial_match(self, 
mock_hook_class):
+        mock_hook_class.return_value.query.return_value = 
pd.DataFrame({"literal": [1]})
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+        assert sensor.execute(context={}) is None
+        mock_hook_class.return_value.query.assert_called_once_with(SQL)
+
+    @mock.patch(HOOK_PATH, autospec=True)
+    def test_execute_deferrable_defers_with_sensor_settings(self, 
mock_hook_class):
+        mock_hook_class.return_value.query.return_value = 
pd.DataFrame({"literal": []})
+        sensor = InfluxDB3Sensor(
+            task_id="wait",
+            sql=SQL,
+            influxdb3_conn_id=CONN_ID,
+            fail_on_empty=False,
+            deferrable=True,
+            poke_interval=30,
+            timeout=600,
+        )
+
+        with pytest.raises(TaskDeferred) as exc:
+            sensor.execute(context={})
+
+        trigger = exc.value.trigger
+        assert isinstance(trigger, InfluxDB3SensorTrigger)
+        assert trigger.sql == SQL
+        assert trigger.influxdb3_conn_id == CONN_ID
+        assert trigger.poll_interval == 30
+        assert trigger.fail_on_empty is False
+        assert exc.value.method_name == "execute_complete"
+        assert exc.value.timeout == timedelta(minutes=10)
+
+    def test_execute_complete_success(self):
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+        assert sensor.execute_complete(context={}, event={"status": 
"success"}) is None
+
+    @pytest.mark.parametrize(
+        ("event", "match"),
+        [
+            pytest.param(
+                {"status": "fail", "message": "No rows returned, raising as 
per fail_on_empty flag"},
+                "fail_on_empty",
+                id="fail-with-message",
+            ),
+            pytest.param(
+                {"status": "fail"},
+                "InfluxDB 3 sensor failed",
+                id="fail-without-message",
+            ),
+        ],
+    )
+    def test_execute_complete_fail_on_empty(self, event, match):
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+        with pytest.raises(AirflowFailException, match=match):
+            sensor.execute_complete(context={}, event=event)
+
+    @pytest.mark.parametrize(
+        ("event", "match"),
+        [
+            pytest.param(None, "did not return an event", id="missing-event"),
+            pytest.param({"status": "error", "message": "boom"}, "boom", 
id="error-with-message"),
+            pytest.param({"status": "error"}, "InfluxDB 3 sensor failed", 
id="error-without-message"),
+            pytest.param({"status": "cancelled"}, "unexpected status", 
id="unexpected-status"),
+        ],
+    )
+    def test_execute_complete_failures(self, event, match):
+        sensor = InfluxDB3Sensor(task_id="wait", sql=SQL, deferrable=True)
+
+        with pytest.raises(RuntimeError, match=match):
+            sensor.execute_complete(context={}, event=event)
diff --git a/providers/influxdb/tests/unit/influxdb/test_utils.py 
b/providers/influxdb/tests/unit/influxdb/test_utils.py
index 270474453b1..1e59ea7b639 100644
--- a/providers/influxdb/tests/unit/influxdb/test_utils.py
+++ b/providers/influxdb/tests/unit/influxdb/test_utils.py
@@ -17,8 +17,9 @@
 from __future__ import annotations
 
 import pandas as pd
+import pytest
 
-from airflow.providers.influxdb.utils import _convert_dataframe_to_records
+from airflow.providers.influxdb.utils import _convert_dataframe_to_records, 
_first_cell_is_truthy
 
 
 def test_convert_dataframe_to_records_serializes_rows_and_timestamps():
@@ -33,3 +34,29 @@ def 
test_convert_dataframe_to_records_serializes_rows_and_timestamps():
         {"col1": 1, "timestamp": "2024-01-01T00:00:00.000Z"},
         {"col1": 2, "timestamp": "2024-01-02T03:04:05.000Z"},
     ]
+
+
[email protected](
+    ("dataframe", "expected"),
+    [
+        pytest.param(pd.DataFrame({"literal": [1]}), True, id="numeric-one"),
+        pytest.param(pd.DataFrame({"count": [42]}), True, id="positive-count"),
+        pytest.param(pd.DataFrame({"flag": ["ready"]}), True, 
id="non-empty-string"),
+        pytest.param(pd.DataFrame({"literal": []}), False, id="no-rows"),
+        pytest.param(pd.DataFrame(), False, id="no-columns"),
+        pytest.param(pd.DataFrame({"count": [0]}), False, id="numeric-zero"),
+        pytest.param(pd.DataFrame({"count": ["0"]}), False, id="string-zero"),
+        pytest.param(pd.DataFrame({"value": [float("nan")]}), False, id="nan"),
+        pytest.param(pd.DataFrame({"value": [None]}), False, id="none"),
+        pytest.param(pd.DataFrame({"first": [0, 1], "second": [1, 1]}), False, 
id="first-cell-only"),
+    ],
+)
+def test_first_cell_is_truthy(dataframe, expected):
+    assert _first_cell_is_truthy(dataframe) is expected
+
+
+def test_first_cell_is_truthy_rejects_non_scalar_value():
+    dataframe = pd.DataFrame({"value": [[1, 2]]})
+
+    with pytest.raises(TypeError, match="must be a scalar"):
+        _first_cell_is_truthy(dataframe)
diff --git a/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py 
b/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
index 18669c504cf..89de9cf467c 100644
--- a/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
+++ b/providers/influxdb/tests/unit/influxdb/triggers/test_influxdb3.py
@@ -22,7 +22,7 @@ from unittest import mock
 import pandas as pd
 import pytest
 
-from airflow.providers.influxdb.triggers.influxdb3 import InfluxDB3QueryTrigger
+from airflow.providers.influxdb.triggers.influxdb3 import 
InfluxDB3QueryTrigger, InfluxDB3SensorTrigger
 from airflow.triggers.base import TriggerEvent
 
 SQL = 'SELECT "duration" FROM "pyexample"'
@@ -79,3 +79,124 @@ class TestInfluxDB3QueryTrigger:
 
         with pytest.raises(asyncio.CancelledError):
             await anext(trigger.run())
+
+
+class TestInfluxDB3SensorTrigger:
+    def test_serialization(self):
+        """Trigger serializes its constructor arguments."""
+        trigger = InfluxDB3SensorTrigger(
+            sql=SQL,
+            influxdb3_conn_id=CONN_ID,
+            poll_interval=30,
+            fail_on_empty=True,
+        )
+        classpath, kwargs = trigger.serialize()
+
+        assert classpath == 
"airflow.providers.influxdb.triggers.influxdb3.InfluxDB3SensorTrigger"
+        assert kwargs == {
+            "sql": SQL,
+            "influxdb3_conn_id": CONN_ID,
+            "poll_interval": 30,
+            "fail_on_empty": True,
+        }
+
+    @pytest.mark.asyncio
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep", 
new_callable=mock.AsyncMock)
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook", 
autospec=True)
+    async def test_run_succeeds_without_sleep_when_condition_is_met(self, 
mock_hook_class, mock_sleep):
+        mock_hook = mock_hook_class.return_value
+        mock_hook.query_async = 
mock.AsyncMock(return_value=pd.DataFrame({"literal": [1]}))
+
+        events = [event async for event in 
InfluxDB3SensorTrigger(sql=SQL).run()]
+
+        mock_hook_class.assert_called_once_with(conn_id="influxdb3_default")
+        mock_hook.query_async.assert_awaited_once_with(SQL)
+        mock_sleep.assert_not_awaited()
+        assert events == [TriggerEvent({"status": "success"})]
+
+    @pytest.mark.asyncio
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep", 
new_callable=mock.AsyncMock)
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook", 
autospec=True)
+    async def test_run_polls_until_condition_is_met(self, mock_hook_class, 
mock_sleep):
+        mock_hook = mock_hook_class.return_value
+        mock_hook.query_async = mock.AsyncMock(
+            side_effect=[
+                pd.DataFrame({"literal": []}),
+                pd.DataFrame({"count": [0]}),
+                pd.DataFrame({"literal": [1]}),
+            ]
+        )
+
+        events = [event async for event in InfluxDB3SensorTrigger(sql=SQL, 
poll_interval=30).run()]
+
+        assert mock_hook.query_async.await_count == 3
+        assert mock_sleep.await_args_list == [mock.call(30), mock.call(30)]
+        assert events == [TriggerEvent({"status": "success"})]
+
+    @pytest.mark.asyncio
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep", 
new_callable=mock.AsyncMock)
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook", 
autospec=True)
+    async def test_run_fail_on_empty(self, mock_hook_class, mock_sleep):
+        mock_hook = mock_hook_class.return_value
+        mock_hook.query_async = 
mock.AsyncMock(return_value=pd.DataFrame({"literal": []}))
+
+        events = [event async for event in InfluxDB3SensorTrigger(sql=SQL, 
fail_on_empty=True).run()]
+
+        mock_hook.query_async.assert_awaited_once_with(SQL)
+        mock_sleep.assert_not_awaited()
+        assert events == [
+            TriggerEvent({"status": "fail", "message": "No rows returned, 
raising as per fail_on_empty flag"})
+        ]
+
+    @pytest.mark.asyncio
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook", 
autospec=True)
+    async def test_run_failure(self, mock_hook_class):
+        mock_hook = mock_hook_class.return_value
+        mock_hook.query_async = mock.AsyncMock(side_effect=ValueError("boom"))
+
+        events = [event async for event in 
InfluxDB3SensorTrigger(sql=SQL).run()]
+
+        mock_hook_class.assert_called_once_with(conn_id="influxdb3_default")
+        mock_hook.query_async.assert_awaited_once_with(SQL)
+        assert events == [TriggerEvent({"status": "error", "message": "boom"})]
+
+    @pytest.mark.asyncio
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep", 
new_callable=mock.AsyncMock)
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook", 
autospec=True)
+    async def test_run_failure_after_unsuccessful_poll(self, mock_hook_class, 
mock_sleep):
+        mock_hook = mock_hook_class.return_value
+        mock_hook.query_async = mock.AsyncMock(
+            side_effect=[pd.DataFrame({"literal": []}), ValueError("boom")]
+        )
+
+        events = [event async for event in InfluxDB3SensorTrigger(sql=SQL, 
poll_interval=30).run()]
+
+        assert mock_hook.query_async.await_count == 2
+        mock_sleep.assert_awaited_once_with(30)
+        assert events == [TriggerEvent({"status": "error", "message": "boom"})]
+
+    @pytest.mark.asyncio
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook", 
autospec=True)
+    async def test_run_propagates_cancellation_during_query(self, 
mock_hook_class):
+        mock_hook = mock_hook_class.return_value
+        mock_hook.query_async = 
mock.AsyncMock(side_effect=asyncio.CancelledError())
+
+        with pytest.raises(asyncio.CancelledError):
+            await anext(InfluxDB3SensorTrigger(sql=SQL).run())
+
+    @pytest.mark.asyncio
+    @mock.patch(
+        "airflow.providers.influxdb.triggers.influxdb3.asyncio.sleep",
+        new_callable=mock.AsyncMock,
+        side_effect=asyncio.CancelledError(),
+    )
+    @mock.patch("airflow.providers.influxdb.triggers.influxdb3.InfluxDB3Hook", 
autospec=True)
+    async def test_run_propagates_cancellation_during_sleep(self, 
mock_hook_class, mock_sleep):
+        mock_hook = mock_hook_class.return_value
+        mock_hook.query_async = 
mock.AsyncMock(return_value=pd.DataFrame({"literal": []}))
+
+        with pytest.raises(asyncio.CancelledError):
+            await anext(InfluxDB3SensorTrigger(sql=SQL, 
poll_interval=30).run())
+
+        mock_hook.query_async.assert_awaited_once_with(SQL)
+        mock_sleep.assert_awaited_once_with(30)

Reply via email to