hkc-8010 opened a new issue, #74009:
URL: https://github.com/apache/airflow/issues/74009

   ### Apache Airflow version
   
   3.2.2 (also reproduced against `main`)
   
   ### What happened?
   
   `CommsDecoder.send()` / `asend()` in the Task SDK never check that the 
response frame they read is the response to the request they just sent. They 
read the next frame off the socket and decode it. `_ResponseFrame.id` is set by 
the supervisor (it echoes the request id), but nothing on the task side 
compares it.
   
   When the stream gets out of step, every later call gets the reply meant for 
the previous request. If that reply is an empty acknowledgement (`body=None, 
error=None`, which the supervisor sends for requests that return nothing), 
`send()` returns `None` and the caller dies on the next attribute access. We 
have hit this twice on the same deployment with different call sites:
   
   ```
   AttributeError: 'NoneType' object has no attribute 'start_date'
     File ".../airflow/sdk/bases/sensor.py", line 186, in execute
     File ".../airflow/sdk/execution_time/task_runner.py", line 522, in 
get_first_reschedule_date
   ```
   
   ```
   AttributeError: 'NoneType' object has no attribute 'iter_asset_event_results'
     File ".../airflow/sdk/execution_time/context.py", line 659
   ```
   
   Neither call site can legitimately get `None` back. For the first one, the 
client wraps the API reply as 
`TaskRescheduleStartDate(start_date=resp.json())`, so even a JSON `null` from 
the server arrives as a populated model, not `None`. The only way `send()` 
returns `None` for these requests is by reading a frame that belongs to another 
request.
   
   What desynchronised the stream in our case was the OpenLineage provider's 
fork mode (`execute_in_thread=False`, still the default in 2.20.x). 
`_fork_execute()` does a bare `os.fork()`, the child inherits the supervisor 
socket, and the listener kills the child when it exceeds its timeout. Task log 
from the failing reschedule poke (sensor in `mode="reschedule"`, third poke):
   
   ```
   01:16:02.233 This process (pid=257176) is multi-threaded, use of fork() may 
lead to deadlocks in the child.
   01:16:02.249 OpenLineage will process 12 SQL hook lineage extra(s).
   01:16:02.25 - 01:16:06.03  (child resolving connections, which goes through 
SUPERVISOR_COMMS)
   01:16:12.275 OpenLineage process with pid `257278` expired and will be 
terminated by listener.
   01:16:15.296 ::endgroup::   (end of pre-execute)
   01:16:22.215 Task failed with exception ... AttributeError: 'NoneType' 
object has no attribute 'start_date'
   ```
   
   The child was killed with a request in flight, the supervisor answered it 
anyway, and that orphan answer was the first thing the parent read on its next 
`send()`. The asset sensor case has the same shape, and there the task raised 
100 ms before the API server had even finished serving the request it was 
supposedly answering.
   
   OpenLineage has since added `execute_in_thread` (#68708) and that is the 
right fix for that provider. This issue is about the Task SDK side: any code 
that touches `SUPERVISOR_COMMS` from a forked child, or anything else that 
leaves an unread response on the socket, produces silent wrong data rather than 
an error. A `None` that crashes is the lucky outcome. The same off-by-one can 
hand a caller a `ConnectionResult` or `XComResult` that belongs to a different 
request.
   
   ### What you think should happen instead?
   
   `send()` / `asend()` should compare `frame.id` with the id of the request 
they sent and raise a clear error when they differ, instead of returning 
someone else's response. Responses on this channel are strictly in request 
order (`send()` holds `_thread_lock` for the round trip), so a mismatch is 
always a bug and never a reordering to wait out.
   
   The triggerer already does the equivalent: `TriggerCommsDecoder` keys 
pending futures by `frame.id` and logs "Got response for unknown request frame".
   
   `_get_response()` itself should stay as it is, because it is also used to 
read the unsolicited `StartupDetails` / `DagFileParseRequest` frames, which 
arrive with `request_id=0` and no matching request.
   
   One limitation worth stating: a forked child starts with a copy of the 
parent's `id_counter`, so the child's first request id can equal the parent's 
next one. A check on `frame.id` catches the common case (the child had already 
made one or more requests before it was killed, as above) but cannot catch that 
exact collision. It still turns most of these from wrong data into an error 
that names the cause.
   
   ### How to reproduce
   
   Put a response frame with an unexpected id on the socket before a request:
   
   ```python
   import socket
   import msgspec
   from airflow.sdk.execution_time.comms import CommsDecoder, GetVariable, 
_ResponseFrame
   
   r, w = socket.socketpair()
   stray = msgspec.msgpack.encode(_ResponseFrame(7, None, None))
   w.sendall(len(stray).to_bytes(4, "big") + stray)
   
   print(CommsDecoder(socket=r).send(GetVariable(key="a")))  # None, request id 
was 0
   ```
   
   In a real deployment: Airflow 3.2.x, `apache-airflow-providers-openlineage` 
2.19 / 2.20 with `execute_in_thread` left at its default, a task that emits SQL 
hook lineage slowly enough to hit the OpenLineage execution timeout, and a 
reschedule-mode sensor or asset sensor that makes a supervisor call after 
pre-execute.
   
   ### Operating System
   
   Debian (Astro Runtime image)
   
   ### Versions of Apache Airflow Providers
   
   apache-airflow-providers-openlineage 2.20.x (fork mode), 
apache-airflow-providers-sftp
   
   ### Deployment
   
   Astronomer
   
   ### Anything else?
   
   I have a PR ready that adds the check to `send()` / `asend()` with tests.
   
   ### Are you willing to submit PR?
   
   - [x] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's Code of Conduct
   


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