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:
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?
Code of Conduct
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.idis 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()returnsNoneand the caller dies on the next attribute access. We have hit this twice on the same deployment with different call sites:Neither call site can legitimately get
Noneback. For the first one, the client wraps the API reply asTaskRescheduleStartDate(start_date=resp.json()), so even a JSONnullfrom the server arrives as a populated model, notNone. The only waysend()returnsNonefor 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 bareos.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 inmode="reschedule", third poke):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 touchesSUPERVISOR_COMMSfrom a forked child, or anything else that leaves an unread response on the socket, produces silent wrong data rather than an error. ANonethat crashes is the lucky outcome. The same off-by-one can hand a caller aConnectionResultorXComResultthat belongs to a different request.What you think should happen instead?
send()/asend()should compareframe.idwith 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_lockfor the round trip), so a mismatch is always a bug and never a reordering to wait out.The triggerer already does the equivalent:
TriggerCommsDecoderkeys pending futures byframe.idand logs "Got response for unknown request frame"._get_response()itself should stay as it is, because it is also used to read the unsolicitedStartupDetails/DagFileParseRequestframes, which arrive withrequest_id=0and 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 onframe.idcatches 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:
In a real deployment: Airflow 3.2.x,
apache-airflow-providers-openlineage2.19 / 2.20 withexecute_in_threadleft 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?
Code of Conduct