dabla commented on code in PR #62922:
URL: https://github.com/apache/airflow/pull/62922#discussion_r4145617467


##########
task-sdk/src/airflow/sdk/execution_time/supervisor.py:
##########
@@ -2404,6 +2404,16 @@ def send(self, msg: BaseModel):
 
         return self._get_response()
 
+    async def asend(self, msg: BaseModel):

Review Comment:
   Reproduced: 8 sync items with `task_concurrency=4` pushing and reading XComs 
failed under the in-process supervisor with `ImportError: cannot import name 
'SUPERVISOR_COMMS'` every time, and concurrent senders could take each other's 
answers. 467c7ee292 does what you suggested. A lock is held from 
`_handle_request` to `popleft`. The module-level delete is replaced by a 
per-thread flag: `task_runner.supervisor_comms()` returns `None` only on the 
thread serving a request. The readers that decided "server side" from the 
attribute now go through it: `models.Variable` and `models.Connection` via a 
new `airflow.utils.helpers.in_task_execution_context()`, which still avoids 
importing the task runner, plus `mask_secret` forwarding and the 
secrets-backend choice. Tests cover crossed answers, per-thread visibility and 
the DB case at full concurrency.
   
   Drafted-by: Claude Opus 5.5; reviewed by @dabla before posting
   



-- 
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