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]