TillerBurr opened a new issue, #73858:
URL: https://github.com/apache/airflow/issues/73858

   ### Under which category would you file this issue?
   
   Airflow Core
   
   ### Apache Airflow version
   
   3.1.7 and current main
   
   ### What happened and how to reproduce it?
   
   When DAGs are shipped as a zip in the local dags-folder bundle, the dag 
processor sometimes SIGKILLs a DAG callback while it is still running. The 
callback never completes. No file changed and no deploy ran. The trigger is the 
periodic bundle refresh.
   
   The kill only happens when a callback is still running at the moment a 
refresh runs. 
   
   This is related to #66483, but it's a different bug. #66483 was a mismatch 
on `bundle_version`, and #66484 fixed it by comparing `presence_key = 
(bundle_name, rel_path)`. Here the mismatch is on `rel_path` itself, so 
`presence_key` doesn't match either, and main is still affected.
   
   ### Reproduction:
   1. Put a DAG with a DAG-level on_success_callback that sleeps about 60 s 
into a zip in the dags folder.
   2. Keep the default bundle refresh interval, so the refresh runs while the 
callback sleeps.
   3. Trigger the DAG and let it succeed.
   4. Watch the dag processor log. You should see Stopping processor for 
zip/dag.py followed by exit_code=\<Negsignal.SIGKILL: -9\>, and the callback's 
side effect never happens.
   
   The sequence we observed in production:
   
   1. A DAG run succeeds and the scheduler creates a DagCallbackRequest for 
on_success_callback.
   2. The dag processor queues the callback. Its DagFileInfo is keyed by the 
DAG's path inside the zip: rel_path=my_dags.zip/my_dag.py.
   3. About 11 s later the timed bundle refresh runs (LocalDagBundle.refresh() 
is a no-op, but the refresh still rescans). _find_files_in_bundle() returns 
top-level paths only, including my_dags.zip.
   4. terminate_orphan_processes() finds that my_dags.zip/my_dag.py isn't in 
the scanned set and kills the processor:
   WARNING - Stopping processor for my_dags.zip/my_dag.py
   INFO - Process exited pid=133594 exit_code=<Negsignal.SIGKILL: -9> 
signal_sent=SIGKILL
   
   The callback, a Slack notification, was never sent.
   
   Minimal example showing keys don't match:
   
   ```python
   from pathlib import Path
   from unittest import mock
   
   import airflow
   from airflow.dag_processing.manager import DagFileInfo, 
DagFileProcessorManager
   
   bundle = dict(bundle_name="dags-folder", bundle_path=Path("/opt/dags"))
   callback_file = DagFileInfo(rel_path=Path("my_dags.zip/my_dag.py"), 
**bundle)  # key from _add_callback_to_queue
   scanned_file = DagFileInfo(rel_path=Path("my_dags.zip"), **bundle)  # what 
the bundle refresh finds
   
   print(callback_file.presence_key == scanned_file.prescence_key) # False
   
   # Get killed as an orphan
   manager = DagFileProcessorManager(max_runs=1)
   processor = mock.MagicMock()
   manager._processors = {callback_file: processor}
   
   manager.terminate_orphan_processes(present={scanned_file})
   
   print("airflow", airflow.__version__)
   print("callback processor killed:", processor.kill.called, 
processor.kill.call_args)
   ```
   
   
   ### What you think should happen instead?
   
   A callback for a DAG inside a zip should not be treated as orphaned while 
its containing zip is still present in the bundle scan. File presence for 
`my_dags.zip/my_dag.py` should be satisfied by `my_dags.zip` being present.
   
   ### Operating System
   
   _No response_
   
   ### Deployment
   
   None
   
   ### Apache Airflow Provider(s)
   
   _No response_
   
   ### Versions of Apache Airflow Providers
   
   _No response_
   
   ### Official Helm Chart version
   
   Not Applicable
   
   ### Kubernetes Version
   
   _No response_
   
   ### Helm Chart configuration
   
   _No response_
   
   ### Docker Image customizations
   
   _No response_
   
   ### Anything else?
   
   _No response_
   
   ### 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](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
   


-- 
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