hkc-8010 opened a new pull request, #74010: URL: https://github.com/apache/airflow/pull/74010
`CommsDecoder.send()` and `asend()` read the next frame off the supervisor socket and decoded it without checking that it answered the request they had just sent. The supervisor already echoes the request id into `_ResponseFrame.id`; nothing on the task side looked at it. If the stream gets out of step, every later call silently receives the reply to an earlier request. In #74009 that came from the OpenLineage provider's fork mode: the forked child shares the socket, gets killed by the listener with a request in flight, and its orphan reply (an empty ack) is then returned as the result of the parent's next `send()`. That surfaced as `'NoneType' object has no attribute 'start_date'` in `get_first_reschedule_date()` and `'NoneType' object has no attribute 'iter_asset_event_results'` in asset sensors. A `None` is the lucky case; the same off-by-one can hand a caller another request's `ConnectionResult` or `XComResult`. This PR passes the request id from `send()` / `asend()` (including the `ResendLoggingFD` path) into `_from_frame()`, which raises a `RuntimeError` naming both ids when they differ. Responses on this channel are strictly in request order because `send()` holds `_thread_lock` for the whole round trip, so a mismatch is never a reordering worth waiting out. What it deliberately does not change: - `_get_response()` still reads without an id check. It is also used for the unsolicited startup frames (`StartupDetails`, `DagFileParseRequest`), which arrive with `request_id=0` and have no matching request. - `TriggerCommsDecoder` already matches responses to pending futures by `frame.id`, so its reader loop calls `_from_frame()` without an id and is unaffected. - I checked every supervisor-side `send_msg` (task, dag processor, triggerer, callback supervisor, and `_send_new_log_fd`), and all of them echo the request id they were given. Known limitation: a forked child copies the parent's `id_counter`, so its first request id can equal the parent's next one. That exact collision still slips through. In the incident above the child had already made several requests before it was killed, so the check catches it. Tests: two new tests in `test_comms.py` queue a stray empty ack with a different id and assert `send()` / `asend()` raise. Both fail on `main` with `DID NOT RAISE`, because `send()` returns `None` there. Ran `test_comms.py` and `test_supervisor.py` under Breeze on Python 3.10 (345 passed), plus `dag_processing/test_processor.py`, `jobs/test_triggerer_job.py` and `cli/commands/test_dag_command.py`, since those drive a real decoder (368 passed). closes: #74009 --- ##### Was generative AI tooling used to co-author this PR? - [X] Yes (please specify the tool below) Generated-by: Claude Code following [the guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions) -- 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]
