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]

Reply via email to