This is an automated email from the ASF dual-hosted git repository.
o-nikolas 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 f7087bc2feb Fix empty status messages in deferrable Bedrock/OpenSearch
sensors (#72330)
f7087bc2feb is described below
commit f7087bc2febc6785013e211bd619fa77d5629d13
Author: PoAn Yang <[email protected]>
AuthorDate: Tue Sep 15 03:57:51 2026 +0900
Fix empty status messages in deferrable Bedrock/OpenSearch sensors (#72330)
Two Bedrock and an OpenSearch Trigger were not extracting the correct
status fields from responses.
---
.../providers/amazon/aws/triggers/bedrock.py | 4 +-
.../amazon/aws/triggers/opensearch_serverless.py | 6 ++-
.../tests/unit/amazon/aws/triggers/test_bedrock.py | 62 ++++++++++++++++++++++
.../aws/triggers/test_opensearch_serverless.py | 51 ++++++++++++++++++
4 files changed, 120 insertions(+), 3 deletions(-)
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/triggers/bedrock.py
b/providers/amazon/src/airflow/providers/amazon/aws/triggers/bedrock.py
index abf0bb4c9db..4671d62007e 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/triggers/bedrock.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/triggers/bedrock.py
@@ -91,7 +91,7 @@ class BedrockKnowledgeBaseActiveTrigger(AwsBaseWaiterTrigger):
waiter_args={"knowledgeBaseId": knowledge_base_id},
failure_message="Bedrock Knowledge Base creation failed.",
status_message="Status of Bedrock Knowledge Base job is",
- status_queries=["status"],
+ status_queries=["knowledgeBase.status",
"knowledgeBase.failureReasons"],
return_key="knowledge_base_id",
return_value=knowledge_base_id,
waiter_delay=waiter_delay,
@@ -177,7 +177,7 @@ class BedrockIngestionJobTrigger(AwsBaseWaiterTrigger):
},
failure_message="Bedrock ingestion job creation failed.",
status_message="Status of Bedrock ingestion job is",
- status_queries=["status"],
+ status_queries=["ingestionJob.status",
"ingestionJob.failureReasons"],
return_key="ingestion_job_id",
return_value=ingestion_job_id,
waiter_delay=waiter_delay,
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/triggers/opensearch_serverless.py
b/providers/amazon/src/airflow/providers/amazon/aws/triggers/opensearch_serverless.py
index 3b090b1ad91..7543ac0d64c 100644
---
a/providers/amazon/src/airflow/providers/amazon/aws/triggers/opensearch_serverless.py
+++
b/providers/amazon/src/airflow/providers/amazon/aws/triggers/opensearch_serverless.py
@@ -57,7 +57,11 @@ class
OpenSearchServerlessCollectionActiveTrigger(AwsBaseWaiterTrigger):
waiter_args={"ids": [collection_id]} if collection_id else
{"names": [collection_name]},
failure_message="OpenSearch Serverless Collection creation
failed.",
status_message="Status of OpenSearch Serverless Collection is",
- status_queries=["status"],
+ status_queries=[
+ "collectionDetails[0].status",
+ "collectionDetails[0].failureMessage",
+ "collectionErrorDetails[0].errorMessage",
+ ],
return_key="collection_id" if collection_id else "collection_name",
return_value=collection_id if collection_id else collection_name,
waiter_delay=waiter_delay,
diff --git a/providers/amazon/tests/unit/amazon/aws/triggers/test_bedrock.py
b/providers/amazon/tests/unit/amazon/aws/triggers/test_bedrock.py
index 8a865c74626..2ccb341fd08 100644
--- a/providers/amazon/tests/unit/amazon/aws/triggers/test_bedrock.py
+++ b/providers/amazon/tests/unit/amazon/aws/triggers/test_bedrock.py
@@ -20,6 +20,7 @@ from unittest import mock
from unittest.mock import AsyncMock
import pytest
+from botocore.exceptions import WaiterError
from airflow.providers.amazon.aws.hooks.bedrock import (
BedrockAgentCoreControlHook,
@@ -41,6 +42,7 @@ from airflow.triggers.base import TriggerEvent
from unit.amazon.aws.utils.test_waiter import assert_expected_waiter_type
BASE_TRIGGER_CLASSPATH = "airflow.providers.amazon.aws.triggers.bedrock."
+FAILURE_REASON = "Access denied when calling Bedrock. Check your request
permissions and retry."
class TestBaseBedrockTrigger:
@@ -138,6 +140,34 @@ class
TestBedrockKnowledgeBaseActiveTrigger(TestBaseBedrockTrigger):
assert_expected_waiter_type(mock_get_waiter, self.EXPECTED_WAITER_NAME)
mock_get_waiter().wait.assert_called_once()
+ @pytest.mark.asyncio
+ @mock.patch.object(BedrockAgentHook, "get_waiter")
+ @mock.patch.object(BedrockAgentHook, "get_async_conn")
+ async def test_run_failure_message_includes_status(self, mock_async_conn,
mock_get_waiter):
+ mock_async_conn.__aenter__.return_value = mock.MagicMock()
+ mock_get_waiter().wait = AsyncMock(
+ side_effect=WaiterError(
+ name=self.EXPECTED_WAITER_NAME,
+ reason='Waiter encountered a terminal failure state: For
expression "knowledgeBase.status"',
+ last_response={
+ "knowledgeBase": {
+ "knowledgeBaseId": self.KNOWLEDGE_BASE_NAME,
+ "status": "FAILED",
+ "failureReasons": [FAILURE_REASON],
+ }
+ },
+ )
+ )
+ trigger =
BedrockKnowledgeBaseActiveTrigger(knowledge_base_id=self.KNOWLEDGE_BASE_NAME)
+
+ response = await trigger.run().asend(None)
+
+ assert response.payload["status"] == "error"
+ assert (
+ response.payload["message"].splitlines()[0]
+ == f"Bedrock Knowledge Base creation failed.: FAILED -
['{FAILURE_REASON}']"
+ )
+
class TestBedrockIngestionJobTrigger(TestBaseBedrockTrigger):
EXPECTED_WAITER_NAME = "ingestion_job_complete"
@@ -178,6 +208,38 @@ class
TestBedrockIngestionJobTrigger(TestBaseBedrockTrigger):
assert response == TriggerEvent({"status": "success",
"ingestion_job_id": self.INGESTION_JOB_ID})
mock_get_waiter().wait.assert_called_once()
+ @pytest.mark.asyncio
+ @mock.patch.object(BedrockAgentHook, "get_waiter")
+ @mock.patch.object(BedrockAgentHook, "get_async_conn")
+ async def test_run_failure_message_includes_status(self, mock_async_conn,
mock_get_waiter):
+ mock_async_conn.__aenter__.return_value = mock.MagicMock()
+ mock_get_waiter().wait = AsyncMock(
+ side_effect=WaiterError(
+ name=self.EXPECTED_WAITER_NAME,
+ reason='Waiter encountered a terminal failure state: For
expression "ingestionJob.status"',
+ last_response={
+ "ingestionJob": {
+ "ingestionJobId": self.INGESTION_JOB_ID,
+ "status": "FAILED",
+ "failureReasons": [FAILURE_REASON],
+ }
+ },
+ )
+ )
+ trigger = BedrockIngestionJobTrigger(
+ knowledge_base_id=self.KNOWLEDGE_BASE_ID,
+ data_source_id=self.DATA_SOURCE_ID,
+ ingestion_job_id=self.INGESTION_JOB_ID,
+ )
+
+ response = await trigger.run().asend(None)
+
+ assert response.payload["status"] == "error"
+ assert (
+ response.payload["message"].splitlines()[0]
+ == f"Bedrock ingestion job creation failed.: FAILED -
['{FAILURE_REASON}']"
+ )
+
class TestBedrockAgentRuntimeReadyTrigger(TestBaseBedrockTrigger):
EXPECTED_WAITER_NAME = "agent_runtime_ready"
diff --git
a/providers/amazon/tests/unit/amazon/aws/triggers/test_opensearch_serverless.py
b/providers/amazon/tests/unit/amazon/aws/triggers/test_opensearch_serverless.py
index 37ce9d3843c..2dbe91091e4 100644
---
a/providers/amazon/tests/unit/amazon/aws/triggers/test_opensearch_serverless.py
+++
b/providers/amazon/tests/unit/amazon/aws/triggers/test_opensearch_serverless.py
@@ -20,17 +20,21 @@ from unittest import mock
from unittest.mock import AsyncMock
import pytest
+from botocore.exceptions import WaiterError
from airflow.providers.amazon.aws.hooks.opensearch_serverless import
OpenSearchServerlessHook
from airflow.providers.amazon.aws.triggers.opensearch_serverless import (
OpenSearchServerlessCollectionActiveTrigger,
)
+from airflow.providers.amazon.aws.utils.waiter_with_logging import
_LazyStatusFormatter
from airflow.triggers.base import TriggerEvent
from airflow.utils.helpers import prune_dict
from unit.amazon.aws.triggers.test_base import TestAwsBaseWaiterTrigger
BASE_TRIGGER_CLASSPATH =
"airflow.providers.amazon.aws.triggers.opensearch_serverless."
+FAILURE_MESSAGE = "The KMS key used to encrypt the collection is no longer
accessible."
+NOT_FOUND_MESSAGE = "Collection with name test_collection_name not found."
class TestBaseBedrockTrigger(TestAwsBaseWaiterTrigger):
@@ -87,3 +91,50 @@ class TestOpenSearchServerlessCollectionActiveTrigger:
assert response == TriggerEvent({"status": "success", "collection_id":
self.COLLECTION_ID})
assert mock_get_waiter().wait.call_count == 1
+
+ @pytest.mark.asyncio
+ @mock.patch.object(OpenSearchServerlessHook, "get_waiter")
+ @mock.patch.object(OpenSearchServerlessHook, "get_async_conn")
+ async def test_run_failure_message_includes_status(self, mock_async_conn,
mock_get_waiter):
+ mock_async_conn.__aenter__.return_value = mock.MagicMock()
+ mock_get_waiter().wait = AsyncMock(
+ side_effect=WaiterError(
+ name=self.EXPECTED_WAITER_NAME,
+ reason=(
+ "Waiter encountered a terminal failure state: "
+ 'For expression "collectionDetails[0].status"'
+ ),
+ last_response={
+ "collectionDetails": [
+ {
+ "id": self.COLLECTION_ID,
+ "status": "FAILED",
+ "failureCode": "KMS_KEY_INACCESSIBLE",
+ "failureMessage": FAILURE_MESSAGE,
+ }
+ ],
+ "collectionErrorDetails": [],
+ },
+ )
+ )
+ trigger =
OpenSearchServerlessCollectionActiveTrigger(collection_id=self.COLLECTION_ID)
+
+ response = await trigger.run().asend(None)
+
+ assert response.payload["status"] == "error"
+ assert (
+ response.payload["message"].splitlines()[0]
+ == f"OpenSearch Serverless Collection creation failed.: FAILED -
{FAILURE_MESSAGE}"
+ )
+
+ def test_status_queries_surface_collection_error_details(self):
+ """A collection that does not exist comes back as a 200 with only
collectionErrorDetails set."""
+ trigger =
OpenSearchServerlessCollectionActiveTrigger(collection_name=self.COLLECTION_NAME)
+ not_found = {
+ "collectionDetails": [],
+ "collectionErrorDetails": [
+ {"name": self.COLLECTION_NAME, "errorCode": "NOT_FOUND",
"errorMessage": NOT_FOUND_MESSAGE}
+ ],
+ }
+
+ assert str(_LazyStatusFormatter(trigger.status_queries, not_found)) ==
NOT_FOUND_MESSAGE