This is an automated email from the ASF dual-hosted git repository.
potiuk 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 0361344c873 Add Cloud Run job links during execution (#70179)
0361344c873 is described below
commit 0361344c873a3dbb966858ba29fb6d9755fa4e6d
Author: Ulada Zakharava <[email protected]>
AuthorDate: Thu Aug 13 14:53:00 2026 +0200
Add Cloud Run job links during execution (#70179)
---
providers/google/provider.yaml | 1 +
.../providers/google/cloud/links/cloud_run.py | 34 +++++++++++++++++++-
.../providers/google/cloud/operators/cloud_run.py | 17 ++++++++--
.../airflow/providers/google/get_provider_info.py | 1 +
.../unit/google/cloud/links/test_cloud_run.py | 36 ++++++++++++++++++++--
.../unit/google/cloud/operators/test_cloud_run.py | 28 +++++++++++++++++
6 files changed, 111 insertions(+), 6 deletions(-)
diff --git a/providers/google/provider.yaml b/providers/google/provider.yaml
index c8aa881ad0d..80624bb5ffb 100644
--- a/providers/google/provider.yaml
+++ b/providers/google/provider.yaml
@@ -1385,6 +1385,7 @@ extra-links:
-
airflow.providers.google.cloud.links.compute.ComputeInstanceTemplateDetailsLink
-
airflow.providers.google.cloud.links.compute.ComputeInstanceGroupManagerDetailsLink
- airflow.providers.google.cloud.links.cloud_run.CloudRunJobLoggingLink
+ -
airflow.providers.google.cloud.links.cloud_run.CloudRunJobExecutionDetailsLink
- airflow.providers.google.cloud.links.cloud_tasks.CloudTasksQueueLink
- airflow.providers.google.cloud.links.cloud_tasks.CloudTasksLink
- airflow.providers.google.cloud.links.dataproc.DataprocLink
diff --git
a/providers/google/src/airflow/providers/google/cloud/links/cloud_run.py
b/providers/google/src/airflow/providers/google/cloud/links/cloud_run.py
index 11efa9d80ba..f3363af4152 100644
--- a/providers/google/src/airflow/providers/google/cloud/links/cloud_run.py
+++ b/providers/google/src/airflow/providers/google/cloud/links/cloud_run.py
@@ -16,7 +16,18 @@
# under the License.
from __future__ import annotations
-from airflow.providers.google.cloud.links.base import BaseGoogleLink
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.google.cloud.links.base import BASE_LINK, BaseGoogleLink
+
+if TYPE_CHECKING:
+ from airflow.providers.common.compat.sdk import Context
+
+
+CLOUD_RUN_BASE_LINK = "/run"
+CLOUD_RUN_JOB_EXECUTION_DETAILS_LINK = (
+ CLOUD_RUN_BASE_LINK +
"/jobs/details/{region}/{job_name}/executions?project={project_id}"
+)
class CloudRunJobLoggingLink(BaseGoogleLink):
@@ -25,3 +36,24 @@ class CloudRunJobLoggingLink(BaseGoogleLink):
name = "Cloud Run Job Logging"
key = "log_uri"
format_str = "{log_uri}"
+
+ @classmethod
+ def persist(cls, context: Context, **value: Any) -> None:
+ """Persist the complete URL expected by serialized operator links."""
+ context["ti"].xcom_push(key=cls.key, value=value["log_uri"])
+
+
+class CloudRunJobExecutionDetailsLink(BaseGoogleLink):
+ """Helper class for constructing a Cloud Run Job execution details link."""
+
+ name = "Cloud Run Job Execution Details"
+ key = "cloud_run_job_execution_details"
+ format_str = CLOUD_RUN_JOB_EXECUTION_DETAILS_LINK
+
+ @classmethod
+ def persist(cls, context: Context, **value: Any) -> None:
+ """Persist the complete URL so it is usable while the execution is
running."""
+ context["ti"].xcom_push(
+ key=cls.key,
+ value=BASE_LINK + cls.format_str.format(**value),
+ )
diff --git
a/providers/google/src/airflow/providers/google/cloud/operators/cloud_run.py
b/providers/google/src/airflow/providers/google/cloud/operators/cloud_run.py
index f1c307f55f4..babe37ee5e4 100644
--- a/providers/google/src/airflow/providers/google/cloud/operators/cloud_run.py
+++ b/providers/google/src/airflow/providers/google/cloud/operators/cloud_run.py
@@ -29,7 +29,10 @@ from google.cloud.run_v2 import Job, Service
from airflow.providers.common.compat.sdk import AirflowException, conf
from airflow.providers.google.cloud.hooks.cloud_run import CloudRunHook,
CloudRunServiceHook
-from airflow.providers.google.cloud.links.cloud_run import
CloudRunJobLoggingLink
+from airflow.providers.google.cloud.links.cloud_run import (
+ CloudRunJobExecutionDetailsLink,
+ CloudRunJobLoggingLink,
+)
from airflow.providers.google.cloud.operators.cloud_base import
GoogleCloudBaseOperator
from airflow.providers.google.cloud.triggers.cloud_run import
CloudRunJobFinishedTrigger, RunJobStatus
@@ -319,7 +322,10 @@ class CloudRunExecuteJobOperator(GoogleCloudBaseOperator):
Default: ``False``.
"""
- operator_extra_links = (CloudRunJobLoggingLink(),)
+ operator_extra_links = (
+ CloudRunJobLoggingLink(),
+ CloudRunJobExecutionDetailsLink(),
+ )
template_fields = (
"project_id",
"region",
@@ -381,6 +387,13 @@ class CloudRunExecuteJobOperator(GoogleCloudBaseOperator):
log_uri=self.operation.metadata.log_uri,
)
+ CloudRunJobExecutionDetailsLink.persist(
+ context=context,
+ region=self.region,
+ job_name=self.job_name,
+ project_id=self.project_id,
+ )
+
if not self.deferrable:
result: Execution = self._wait_for_operation(self.operation)
if self.verbose and result.name:
diff --git a/providers/google/src/airflow/providers/google/get_provider_info.py
b/providers/google/src/airflow/providers/google/get_provider_info.py
index 7334666783d..e2bfa861e24 100644
--- a/providers/google/src/airflow/providers/google/get_provider_info.py
+++ b/providers/google/src/airflow/providers/google/get_provider_info.py
@@ -1604,6 +1604,7 @@ def get_provider_info():
"airflow.providers.google.cloud.links.compute.ComputeInstanceTemplateDetailsLink",
"airflow.providers.google.cloud.links.compute.ComputeInstanceGroupManagerDetailsLink",
"airflow.providers.google.cloud.links.cloud_run.CloudRunJobLoggingLink",
+
"airflow.providers.google.cloud.links.cloud_run.CloudRunJobExecutionDetailsLink",
"airflow.providers.google.cloud.links.cloud_tasks.CloudTasksQueueLink",
"airflow.providers.google.cloud.links.cloud_tasks.CloudTasksLink",
"airflow.providers.google.cloud.links.dataproc.DataprocLink",
diff --git a/providers/google/tests/unit/google/cloud/links/test_cloud_run.py
b/providers/google/tests/unit/google/cloud/links/test_cloud_run.py
index c2bd362efa3..f7f620ca591 100644
--- a/providers/google/tests/unit/google/cloud/links/test_cloud_run.py
+++ b/providers/google/tests/unit/google/cloud/links/test_cloud_run.py
@@ -21,7 +21,11 @@ from unittest import mock
import pytest
-from airflow.providers.google.cloud.links.cloud_run import
CloudRunJobLoggingLink
+from airflow.providers.google.cloud.links.cloud_run import (
+ CLOUD_RUN_JOB_EXECUTION_DETAILS_LINK,
+ CloudRunJobExecutionDetailsLink,
+ CloudRunJobLoggingLink,
+)
from airflow.providers.google.cloud.operators.cloud_run import
CloudRunExecuteJobOperator
from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS
@@ -32,6 +36,9 @@ if AIRFLOW_V_3_0_PLUS:
TEST_LOG_URI = (
"https://console.cloud.google.com/run/jobs/logs?project=test-project®ion=test-region&job=test-job"
)
+TEST_EXECUTION_DETAILS_URI = (
+
"https://console.cloud.google.com/run/jobs/details/test-region/test-job/executions?project=test-project"
+)
class TestCloudRunJobLoggingLink:
@@ -52,7 +59,7 @@ class TestCloudRunJobLoggingLink:
mock_context["ti"].xcom_push.assert_called_once_with(
key=CloudRunJobLoggingLink.key,
- value={"log_uri": TEST_LOG_URI},
+ value=TEST_LOG_URI,
)
@pytest.mark.db_test
@@ -74,7 +81,30 @@ class TestCloudRunJobLoggingLink:
if mock_supervisor_comms:
mock_supervisor_comms.send.return_value = XComResult(
key="key",
- value={"log_uri": TEST_LOG_URI},
+ value=TEST_LOG_URI,
)
actual_url = link.get_link(operator=ti.task, ti_key=ti.key)
assert actual_url == TEST_LOG_URI
+
+
+class TestCloudRunJobExecutionDetailsLink:
+ def test_class_attributes(self):
+ assert CloudRunJobExecutionDetailsLink.key ==
"cloud_run_job_execution_details"
+ assert CloudRunJobExecutionDetailsLink.name == "Cloud Run Job
Execution Details"
+ assert CloudRunJobExecutionDetailsLink.format_str ==
CLOUD_RUN_JOB_EXECUTION_DETAILS_LINK
+
+ def test_persist(self):
+ mock_context = mock.MagicMock()
+ mock_context["ti"] = mock.MagicMock()
+
+ CloudRunJobExecutionDetailsLink.persist(
+ context=mock_context,
+ region="test-region",
+ job_name="test-job",
+ project_id="test-project",
+ )
+
+ mock_context["ti"].xcom_push.assert_called_once_with(
+ key=CloudRunJobExecutionDetailsLink.key,
+ value=TEST_EXECUTION_DETAILS_URI,
+ )
diff --git
a/providers/google/tests/unit/google/cloud/operators/test_cloud_run.py
b/providers/google/tests/unit/google/cloud/operators/test_cloud_run.py
index 00f339e9493..48e3a6358d6 100644
--- a/providers/google/tests/unit/google/cloud/operators/test_cloud_run.py
+++ b/providers/google/tests/unit/google/cloud/operators/test_cloud_run.py
@@ -43,6 +43,9 @@ from airflow.providers.google.cloud.triggers.cloud_run import
RunJobStatus
CLOUD_RUN_HOOK_PATH =
"airflow.providers.google.cloud.operators.cloud_run.CloudRunHook"
CLOUD_RUN_SERVICE_HOOK_PATH =
"airflow.providers.google.cloud.operators.cloud_run.CloudRunServiceHook"
GCP_LOGGING_PATH =
"airflow.providers.google.cloud.operators.cloud_run.gcp_logging"
+CLOUD_RUN_EXECUTION_DETAILS_LINK_PATH = (
+
"airflow.providers.google.cloud.operators.cloud_run.CloudRunJobExecutionDetailsLink"
+)
TASK_ID = "test"
PROJECT_ID = "testproject"
REGION = "us-central1"
@@ -219,6 +222,31 @@ class TestCloudRunExecuteJobOperator:
with pytest.raises(TaskDeferred):
operator.execute(mock.MagicMock())
+ @mock.patch(f"{CLOUD_RUN_EXECUTION_DETAILS_LINK_PATH}.persist")
+ @mock.patch(CLOUD_RUN_HOOK_PATH)
+ def test_execute_persists_job_execution_details_link(self, hook_mock,
persist_mock):
+ operation = mock.MagicMock()
+ operation.metadata.log_uri = None
+ hook_mock.return_value.execute_job.return_value = operation
+ context = mock.MagicMock()
+ operator = CloudRunExecuteJobOperator(
+ task_id=TASK_ID,
+ project_id=PROJECT_ID,
+ region=REGION,
+ job_name=JOB_NAME,
+ deferrable=True,
+ )
+
+ with pytest.raises(TaskDeferred):
+ operator.execute(context)
+
+ persist_mock.assert_called_once_with(
+ context=context,
+ region=REGION,
+ job_name=JOB_NAME,
+ project_id=PROJECT_ID,
+ )
+
@mock.patch(CLOUD_RUN_HOOK_PATH)
def test_execute_deferrable_execute_complete_method_timeout(self,
hook_mock):
operator = CloudRunExecuteJobOperator(