Skip to content

Task SDK CommsDecoder accepts a response frame for a different request, returning wrong data or None when the supervisor socket desyncs #74009

Description

@hkc-8010

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?

  • Yes I am willing to submit a PR!

Code of Conduct

  • I agree to follow this project's Code of Conduct

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions