amoghrajesh commented on code in PR #72329:
URL: https://github.com/apache/airflow/pull/72329#discussion_r4142899285
##########
task-sdk/src/airflow/sdk/execution_time/context.py:
##########
@@ -454,6 +539,61 @@ def _set_variable(key: str, value: Any, description: str |
None = None, serializ
SecretCache.invalidate_variable(key)
+async def _async_set_variable(
+ key: str,
+ value: Any,
+ description: str | None = None,
+ serialize_json: bool = False,
+) -> None:
+ # TODO: This should probably be moved to a separate module like
`airflow.sdk.execution_time.comms`
+ # or `airflow.sdk.execution_time.variable`
+ # A reason to not move it to `airflow.sdk.execution_time.comms` is that
it
+ # will make that module depend on Task SDK, which is not ideal because
we intend to
+ # keep Task SDK as a separate package than execution time mods.
+ import json
Review Comment:
Already imported at module level
##########
task-sdk/src/airflow/sdk/definitions/variable.py:
##########
@@ -58,12 +58,37 @@ def get(cls, key: str, default: Any = NOTSET,
deserialize_json: bool = False):
return default
raise
+ @classmethod
+ async def aget(cls, key: str, default: Any = NOTSET, deserialize_json:
bool = False):
+ from airflow.sdk.exceptions import AirflowRuntimeError, ErrorType
+ from airflow.sdk.execution_time.context import _async_get_variable
+
+ try:
+ return await _async_get_variable(key,
deserialize_json=deserialize_json)
+ except AirflowRuntimeError as e:
+ if e.error.error == ErrorType.VARIABLE_NOT_FOUND and default is
not NOTSET:
+ await amask_secret(default, name=key)
Review Comment:
The fallback isnt tested in any test.
##########
task-sdk/src/airflow/sdk/execution_time/context.py:
##########
@@ -382,6 +400,52 @@ def _get_variable(key: str, deserialize_json: bool) -> Any:
)
+async def _async_get_variable(key: str, deserialize_json: bool) -> Any:
+ from airflow.sdk.execution_time.cache import SecretCache
+ from airflow.sdk.execution_time.supervisor import
ensure_secrets_backend_loaded
+
+ # Check cache first
+ try:
+ var_val = SecretCache.get_variable(key)
Review Comment:
Write a test similar to sync counterpart: `test_var_json_masks_from_cache`
##########
task-sdk/src/airflow/sdk/execution_time/context.py:
##########
@@ -343,6 +344,23 @@ def _mask_and_deserialize_variable(raw: str, key: str,
deserialize_json: bool) -
return val
+async def _async_mask_and_deserialize_variable(raw: str, key: str,
deserialize_json: bool) -> Any:
Review Comment:
And few more
##########
task-sdk/src/airflow/sdk/execution_time/context.py:
##########
@@ -343,6 +344,23 @@ def _mask_and_deserialize_variable(raw: str, key: str,
deserialize_json: bool) -
return val
+async def _async_mask_and_deserialize_variable(raw: str, key: str,
deserialize_json: bool) -> Any:
Review Comment:
This one has no tests. Refer to tests like `test_var_json_masks_from_cache`
and `test_var_value_masks_secret`
##########
task-sdk/src/airflow/sdk/execution_time/context.py:
##########
@@ -454,6 +539,61 @@ def _set_variable(key: str, value: Any, description: str |
None = None, serializ
SecretCache.invalidate_variable(key)
+async def _async_set_variable(
+ key: str,
+ value: Any,
+ description: str | None = None,
+ serialize_json: bool = False,
+) -> None:
+ # TODO: This should probably be moved to a separate module like
`airflow.sdk.execution_time.comms`
+ # or `airflow.sdk.execution_time.variable`
+ # A reason to not move it to `airflow.sdk.execution_time.comms` is that
it
+ # will make that module depend on Task SDK, which is not ideal because
we intend to
+ # keep Task SDK as a separate package than execution time mods.
+ import json
+
+ from airflow.sdk.execution_time.cache import SecretCache
+ from airflow.sdk.execution_time.secrets.execution_api import (
+ ExecutionAPISecretsBackend,
+ )
+ from airflow.sdk.execution_time.supervisor import
ensure_secrets_backend_loaded
+ from airflow.sdk.execution_time.task_runner import SUPERVISOR_COMMS
+
+ # check for write conflicts on the worker
Review Comment:
This isnt tested. Both `test_async_set_variable` cases patch
`ensure_secrets_backend_loaded` to return `[]`, so the loop body is skipped.
That block is about 20 lines of new code, including a warning message users
will see.
--
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]