Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Fix dag processor callback cleanup for versioned bundle files by hkc-8010 · Pull Request #66484 · apache/airflow · GitHub
Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix dag processor callback cleanup for versioned bundle files by hkc-8010 · Pull Request #66484 · apache/airflow · GitHub
Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix dag processor callback cleanup for versioned bundle files by hkc-8010 · Pull Request #66484 · apache/airflow · GitHub
Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' Fix dag processor callback cleanup for versioned bundle files by hkc-8010 · Pull Request #66484 · apache/airflow · GitHub
Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix dag processor callback cleanup for versioned bundle files by hkc-8010 · Pull Request #66484 · apache/airflow · GitHub
Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix dag processor callback cleanup for versioned bundle files by hkc-8010 · Pull Request #66484 · apache/airflow · GitHub
Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); Fix dag processor callback cleanup for versioned bundle files by hkc-8010 · Pull Request #66484 · apache/airflow · GitHub
Skip to content

Fix dag processor callback cleanup for versioned bundle files - #66484

Merged
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning
May 13, 2026
Merged

Fix dag processor callback cleanup for versioned bundle files#66484
ephraimbuddy merged 5 commits into
apache:mainfrom
hkc-8010:fix/dag-processor-versioned-callback-orphaning

Conversation

@hkc-8010

@hkc-8010hkc-8010 commented May 6, 2026

Copy link
Copy Markdown
Contributor

Fix dag-processor presence/orphan checks so versioned bundle files are not treated as removed or new solely because the queued/tracked DagFileInfo carries a bundle_version.

This is distinct from #66301 / #66474. Those changes addressed stale serialized DAG metadata when only the bundle version changed. This fix is later in the callback pipeline: scheduler-emitted DAG-level and task-level callback requests are fetched and queued correctly, but manager-side presence checks must not purge, orphan, or re-queue the same file just because the scanned DAG file is represented without bundle_version.

The change keeps DagFileInfo equality and hashing version-aware, but adds a presence_key on DagFileInfo and uses that for manager-side “is this file already present / already represented?” checks. This is intentionally narrower than changing global DagFileInfo equality with bundle_version=field(compare=False), because callback requests are tied to a specific bundle version and we do not want to collapse queue/process identity semantics outside these presence checks.

That is now applied in:

  • purge_removed_files_from_queue
  • terminate_orphan_processes
  • remove_orphaned_file_stats
  • _add_new_files_to_queue
  • prepare_file_queue
  • _sort_by_mtime / processed_recently

Tests added:

  • preserve a versioned queued file, processor, and stats entry when the same unversioned file is still present
  • purge/kill/remove those entries when the file is truly absent
  • avoid re-adding a scanned unversioned file when the same path is already represented in queue/processors/stats under a versioned key
  • avoid duplicate queue/stats behavior in prepare_file_queue() and modified-time mode when versioned entries already exist

Validation:

  • pytest airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend postgres --python 3.10 --downgrade-pendulum --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q
  • breeze testing core-tests --backend sqlite --python 3.10 --force-lowest-dependencies --db-reset -- airflow-core/tests/unit/dag_processing/test_manager.py -q

closes: #66483


Was generative AI tooling used to co-author this PR?
  • Yes (Codex)

Generated-by: Codex following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Diagnosis is right and the fix is in the right shape — separating queue identity from file presence is the correct move. Three things I'd want addressed before this lands:

  1. The same identity mismatch lives at manager.py:1267 in _add_new_files_to_queue. Both _file_stats and _processors accumulate versioned entries (set at L1112 / L1254). The next scan's unversioned DagFileInfo doesn't match a versioned stat, so the file is re-queued as "new" and re-parsed every cycle until the version key happens to align. Same bug family, untouched here. If you don't want to fix it in this PR, please open a follow-up issue and link it in the body.

  2. Why does bundle_version need to stay in DagFileInfo identity at all?_callback_to_execute[file_info] is a list — multiple callbacks for the same path already coexist there without needing distinct keys. The only place versioned identity matters is _processors, but L1248 (if file in self._processors: continue) currently lets two different versions each spawn a processor for the same path, which is arguably worse than collapsing them. A one-line bundle_version: str | None = field(compare=False) change would fix the bug without needing the _present_file_key machinery. Worth a sentence in the PR body explaining why the surgical fix was preferred.

  3. Scope is broader than the title suggests. Scheduler-emitted TaskCallbackRequest for normally-versioned bundles flows through the same queue and hits the same orphan check. Worth saying "DAG-level and task-level callbacks" in the body so reviewers understand the actual blast radius.


Additional inline notes (anchoring failed via API; included here)

On manager.py line 1023 — the new deque(...) comprehension in purge_removed_files_from_queue:

All three orphan-check methods recompute the same present_keys set from the same input. handle_removed_files is the only caller — compute it once there and pass keys down:

defhandle_removed_files(self, known_files):
present_keys= {(f.bundle_name, f.rel_path) forvinknown_files.values() forfinv}
self.purge_removed_files_from_queue(present_keys)
self.terminate_orphan_processes(present_keys)
self.remove_orphaned_file_stats(present_keys)

Minor at typical sizes, but the signature change is free and removes the duplication.

On test_manager.py line 453 — processor.kill.assert_not_called():

The three new tests only cover the "preserved" direction. None assert that an entry is correctly purged when the file is genuinely gone — a regression that made _present_file_key collapse all keys to a constant would still pass the suite.

Please add the negative case for each of the three methods. For terminate_orphan_processes specifically:

deftest_terminate_orphan_processes_kills_processor_when_file_is_truly_absent(self):
manager=DagFileProcessorManager(max_runs=1)
versioned_file=DagFileInfo(
bundle_name="testing",
rel_path=Path("callbacks.py"),
bundle_path=TEST_DAGS_FOLDER,
bundle_version="v1",
)
processor=MagicMock()
manager._processors[versioned_file] =processormanager.terminate_orphan_processes(present=set())
assertmanager._processors== {}
processor.kill.assert_called_once()

Without this, assert_not_called() doesn't actually pin the behavior — it just confirms kill wasn't invoked, which is also true if kill is never reachable.


Drafted-by: Claude Code (Opus 4.7); reviewed by @ephraimbuddy before posting

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010hkc-8010 changed the title Fix dag processor orphan cleanup for versioned callback filesFix dag processor callback cleanup for versioned bundle filesMay 6, 2026
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Thanks, this was really helpful. I pushed a follow-up that addresses the points you called out:

  • moved the dual-identity concept onto DagFileInfo as presence_key
  • compute present_keys once in handle_removed_files() and pass it through the orphan-cleanup helpers
  • fixed the same presence/equality mismatch in _add_new_files_to_queue()
  • also made prepare_file_queue(), processed_recently(), and _sort_by_mtime() presence-aware so scanned unversioned files do not get re-queued or depend on duplicate unversioned stats when versioned entries already exist
  • added negative tests for the orphan-cleanup methods, alongside the preserve-direction coverage

I kept bundle_version in DagFileInfo equality/hashing and updated the PR body to explain why. The goal here was to keep callback/process identity version-aware while narrowing the fix to manager-side “present/already represented” checks.

I also updated the PR description to call out the broader scope: DAG-level and task-level callbacks both flow through this path.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
@hkc-8010

Copy link
Copy Markdown
ContributorAuthor

By the way, the customer tested the PR as a patch on their current Airflow version and confirmed that the callbacks resumed functioning without any problems.

Comment threadairflow-core/src/airflow/dag_processing/manager.py Outdated
Comment threadairflow-core/src/airflow/dag_processing/manager.py

@ephraimbuddyephraimbuddy left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Comment threadairflow-core/src/airflow/dag_processing/manager.py
@ephraimbuddy
ephraimbuddy merged commit 6927283 into apache:mainMay 13, 2026
79 checks passed
@eladkaleladkal added this to the Airflow 3.2.3 milestone May 21, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Dag processor orphan cleanup drops versioned callback files before DAG callbacks execute

4 participants

@hkc-8010@ephraimbuddy@vatsrahul1001@eladkal