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]

Reply via email to