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]