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

Description

@hkc-8010

Apache Airflow version

3.1.8 and current main

What happened?

We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

The observed sequence is:

  1. Scheduler creates a DagCallbackRequest
  2. Dag processor fetches it from the DB
  3. Dag processor queues it for file processing
  4. Dag processor starts or prepares callback processing for the file
  5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

What you think should happen instead?

If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

In other words, present/orphan checks in the dag processor should treat file presence as:

  • bundle_name
  • rel_path

and should not require exact DagFileInfo equality including bundle_version.

How to reproduce

A minimal pattern is:

  1. Use a DAG bundle implementation with supports_versioning=True
  2. Have a DAG with a DAG-level on_failure_callback
  3. Trigger a failing DAG run so a DagCallbackRequest is created
  4. In the dag processor:
    • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
    • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
  5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
  6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

Operating System

Linux

Deployment

Other

Anything else?

The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

  • _add_callback_to_queue() creates:

    file_info=DagFileInfo(
    rel_path=Path(request.filepath),
    bundle_path=bundle.path,
    bundle_name=request.bundle_name,
    bundle_version=request.bundle_version,
    )
  • _refresh_dag_bundles() builds known_files with:

    DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
  • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

A small local proof of the mismatch is:

frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
rel_path=Path("dags/example.py"),
bundle_name="main",
bundle_path=Path("/tmp/bundle"),
bundle_version="v1",
)
present=DagFileInfo(
rel_path=Path("dags/example.py"),
bundle_name="main",
bundle_path=Path("/tmp/bundle"),
)
assertversioned!=present

The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

  • bundle_name
  • rel_path

That preserves distinct callback entries per version while preventing false orphan cleanup.

Are you willing to submit PR?

Yes

Metadata

Metadata

Assignees

No one assigned

    Labels

    area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions

      , '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" + '
      
      Skip to content

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

      Description

      @hkc-8010

      Apache Airflow version

      3.1.8 and current main

      What happened?

      We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

      The observed sequence is:

      1. Scheduler creates a DagCallbackRequest
      2. Dag processor fetches it from the DB
      3. Dag processor queues it for file processing
      4. Dag processor starts or prepares callback processing for the file
      5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

      This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

      That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

      What you think should happen instead?

      If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

      In other words, present/orphan checks in the dag processor should treat file presence as:

      • bundle_name
      • rel_path

      and should not require exact DagFileInfo equality including bundle_version.

      How to reproduce

      A minimal pattern is:

      1. Use a DAG bundle implementation with supports_versioning=True
      2. Have a DAG with a DAG-level on_failure_callback
      3. Trigger a failing DAG run so a DagCallbackRequest is created
      4. In the dag processor:
        • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
        • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
      5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
      6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

      Operating System

      Linux

      Deployment

      Other

      Anything else?

      The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

      • _add_callback_to_queue() creates:

        file_info=DagFileInfo(
        rel_path=Path(request.filepath),
        bundle_path=bundle.path,
        bundle_name=request.bundle_name,
        bundle_version=request.bundle_version,
        )
      • _refresh_dag_bundles() builds known_files with:

        DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
      • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

      Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

      That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

      A small local proof of the mismatch is:

      frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
      rel_path=Path("dags/example.py"),
      bundle_name="main",
      bundle_path=Path("/tmp/bundle"),
      bundle_version="v1",
      )
      present=DagFileInfo(
      rel_path=Path("dags/example.py"),
      bundle_name="main",
      bundle_path=Path("/tmp/bundle"),
      )
      assertversioned!=present

      The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

      • bundle_name
      • rel_path

      That preserves distinct callback entries per version while preventing false orphan cleanup.

      Are you willing to submit PR?

      Yes

      Metadata

      Metadata

      Assignees

      No one assigned

        Labels

        area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

        Type

        No type

        Projects

        No projects

          Milestone

          No milestone

          Relationships

          None yet

          Development

          No branches or pull requests

          Issue actions

          , '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('^' + ".*" + '
          Skip to content

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

          Description

          @hkc-8010

          Apache Airflow version

          3.1.8 and current main

          What happened?

          We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

          The observed sequence is:

          1. Scheduler creates a DagCallbackRequest
          2. Dag processor fetches it from the DB
          3. Dag processor queues it for file processing
          4. Dag processor starts or prepares callback processing for the file
          5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

          This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

          That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

          What you think should happen instead?

          If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

          In other words, present/orphan checks in the dag processor should treat file presence as:

          • bundle_name
          • rel_path

          and should not require exact DagFileInfo equality including bundle_version.

          How to reproduce

          A minimal pattern is:

          1. Use a DAG bundle implementation with supports_versioning=True
          2. Have a DAG with a DAG-level on_failure_callback
          3. Trigger a failing DAG run so a DagCallbackRequest is created
          4. In the dag processor:
            • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
            • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
          5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
          6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

          Operating System

          Linux

          Deployment

          Other

          Anything else?

          The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

          • _add_callback_to_queue() creates:

            file_info=DagFileInfo(
            rel_path=Path(request.filepath),
            bundle_path=bundle.path,
            bundle_name=request.bundle_name,
            bundle_version=request.bundle_version,
            )
          • _refresh_dag_bundles() builds known_files with:

            DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
          • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

          Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

          That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

          A small local proof of the mismatch is:

          frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
          rel_path=Path("dags/example.py"),
          bundle_name="main",
          bundle_path=Path("/tmp/bundle"),
          bundle_version="v1",
          )
          present=DagFileInfo(
          rel_path=Path("dags/example.py"),
          bundle_name="main",
          bundle_path=Path("/tmp/bundle"),
          )
          assertversioned!=present

          The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

          • bundle_name
          • rel_path

          That preserves distinct callback entries per version while preventing false orphan cleanup.

          Are you willing to submit PR?

          Yes

          Metadata

          Metadata

          Assignees

          No one assigned

            Labels

            area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

            Type

            No type

            Projects

            No projects

              Milestone

              No milestone

              Relationships

              None yet

              Development

              No branches or pull requests

              Issue actions

              , '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('^' + ".*" + '
              Skip to content

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

              Description

              @hkc-8010

              Apache Airflow version

              3.1.8 and current main

              What happened?

              We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

              The observed sequence is:

              1. Scheduler creates a DagCallbackRequest
              2. Dag processor fetches it from the DB
              3. Dag processor queues it for file processing
              4. Dag processor starts or prepares callback processing for the file
              5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

              This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

              That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

              What you think should happen instead?

              If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

              In other words, present/orphan checks in the dag processor should treat file presence as:

              • bundle_name
              • rel_path

              and should not require exact DagFileInfo equality including bundle_version.

              How to reproduce

              A minimal pattern is:

              1. Use a DAG bundle implementation with supports_versioning=True
              2. Have a DAG with a DAG-level on_failure_callback
              3. Trigger a failing DAG run so a DagCallbackRequest is created
              4. In the dag processor:
                • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
                • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
              5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
              6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

              Operating System

              Linux

              Deployment

              Other

              Anything else?

              The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

              • _add_callback_to_queue() creates:

                file_info=DagFileInfo(
                rel_path=Path(request.filepath),
                bundle_path=bundle.path,
                bundle_name=request.bundle_name,
                bundle_version=request.bundle_version,
                )
              • _refresh_dag_bundles() builds known_files with:

                DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
              • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

              Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

              That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

              A small local proof of the mismatch is:

              frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
              rel_path=Path("dags/example.py"),
              bundle_name="main",
              bundle_path=Path("/tmp/bundle"),
              bundle_version="v1",
              )
              present=DagFileInfo(
              rel_path=Path("dags/example.py"),
              bundle_name="main",
              bundle_path=Path("/tmp/bundle"),
              )
              assertversioned!=present

              The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

              • bundle_name
              • rel_path

              That preserves distinct callback entries per version while preventing false orphan cleanup.

              Are you willing to submit PR?

              Yes

              Metadata

              Metadata

              Assignees

              No one assigned

                Labels

                area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

                Type

                No type

                Projects

                No projects

                  Milestone

                  No milestone

                  Relationships

                  None yet

                  Development

                  No branches or pull requests

                  Issue actions

                  , '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" + '
                  Skip to content

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

                  Description

                  @hkc-8010

                  Apache Airflow version

                  3.1.8 and current main

                  What happened?

                  We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

                  The observed sequence is:

                  1. Scheduler creates a DagCallbackRequest
                  2. Dag processor fetches it from the DB
                  3. Dag processor queues it for file processing
                  4. Dag processor starts or prepares callback processing for the file
                  5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

                  This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

                  That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

                  What you think should happen instead?

                  If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

                  In other words, present/orphan checks in the dag processor should treat file presence as:

                  • bundle_name
                  • rel_path

                  and should not require exact DagFileInfo equality including bundle_version.

                  How to reproduce

                  A minimal pattern is:

                  1. Use a DAG bundle implementation with supports_versioning=True
                  2. Have a DAG with a DAG-level on_failure_callback
                  3. Trigger a failing DAG run so a DagCallbackRequest is created
                  4. In the dag processor:
                    • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
                    • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
                  5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
                  6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

                  Operating System

                  Linux

                  Deployment

                  Other

                  Anything else?

                  The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

                  • _add_callback_to_queue() creates:

                    file_info=DagFileInfo(
                    rel_path=Path(request.filepath),
                    bundle_path=bundle.path,
                    bundle_name=request.bundle_name,
                    bundle_version=request.bundle_version,
                    )
                  • _refresh_dag_bundles() builds known_files with:

                    DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
                  • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

                  Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

                  That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

                  A small local proof of the mismatch is:

                  frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
                  rel_path=Path("dags/example.py"),
                  bundle_name="main",
                  bundle_path=Path("/tmp/bundle"),
                  bundle_version="v1",
                  )
                  present=DagFileInfo(
                  rel_path=Path("dags/example.py"),
                  bundle_name="main",
                  bundle_path=Path("/tmp/bundle"),
                  )
                  assertversioned!=present

                  The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

                  • bundle_name
                  • rel_path

                  That preserves distinct callback entries per version while preventing false orphan cleanup.

                  Are you willing to submit PR?

                  Yes

                  Metadata

                  Metadata

                  Assignees

                  No one assigned

                    Labels

                    area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

                    Type

                    No type

                    Projects

                    No projects

                      Milestone

                      No milestone

                      Relationships

                      None yet

                      Development

                      No branches or pull requests

                      Issue actions

                      , '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('^' + ".*" + '
                      Skip to content

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

                      Description

                      @hkc-8010

                      Apache Airflow version

                      3.1.8 and current main

                      What happened?

                      We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

                      The observed sequence is:

                      1. Scheduler creates a DagCallbackRequest
                      2. Dag processor fetches it from the DB
                      3. Dag processor queues it for file processing
                      4. Dag processor starts or prepares callback processing for the file
                      5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

                      This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

                      That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

                      What you think should happen instead?

                      If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

                      In other words, present/orphan checks in the dag processor should treat file presence as:

                      • bundle_name
                      • rel_path

                      and should not require exact DagFileInfo equality including bundle_version.

                      How to reproduce

                      A minimal pattern is:

                      1. Use a DAG bundle implementation with supports_versioning=True
                      2. Have a DAG with a DAG-level on_failure_callback
                      3. Trigger a failing DAG run so a DagCallbackRequest is created
                      4. In the dag processor:
                        • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
                        • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
                      5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
                      6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

                      Operating System

                      Linux

                      Deployment

                      Other

                      Anything else?

                      The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

                      • _add_callback_to_queue() creates:

                        file_info=DagFileInfo(
                        rel_path=Path(request.filepath),
                        bundle_path=bundle.path,
                        bundle_name=request.bundle_name,
                        bundle_version=request.bundle_version,
                        )
                      • _refresh_dag_bundles() builds known_files with:

                        DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
                      • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

                      Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

                      That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

                      A small local proof of the mismatch is:

                      frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
                      rel_path=Path("dags/example.py"),
                      bundle_name="main",
                      bundle_path=Path("/tmp/bundle"),
                      bundle_version="v1",
                      )
                      present=DagFileInfo(
                      rel_path=Path("dags/example.py"),
                      bundle_name="main",
                      bundle_path=Path("/tmp/bundle"),
                      )
                      assertversioned!=present

                      The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

                      • bundle_name
                      • rel_path

                      That preserves distinct callback entries per version while preventing false orphan cleanup.

                      Are you willing to submit PR?

                      Yes

                      Metadata

                      Metadata

                      Assignees

                      No one assigned

                        Labels

                        area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

                        Type

                        No type

                        Projects

                        No projects

                          Milestone

                          No milestone

                          Relationships

                          None yet

                          Development

                          No branches or pull requests

                          Issue actions

                          , '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('^' + ".*" + '
                          Skip to content

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

                          Description

                          @hkc-8010

                          Apache Airflow version

                          3.1.8 and current main

                          What happened?

                          We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

                          The observed sequence is:

                          1. Scheduler creates a DagCallbackRequest
                          2. Dag processor fetches it from the DB
                          3. Dag processor queues it for file processing
                          4. Dag processor starts or prepares callback processing for the file
                          5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

                          This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

                          That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

                          What you think should happen instead?

                          If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

                          In other words, present/orphan checks in the dag processor should treat file presence as:

                          • bundle_name
                          • rel_path

                          and should not require exact DagFileInfo equality including bundle_version.

                          How to reproduce

                          A minimal pattern is:

                          1. Use a DAG bundle implementation with supports_versioning=True
                          2. Have a DAG with a DAG-level on_failure_callback
                          3. Trigger a failing DAG run so a DagCallbackRequest is created
                          4. In the dag processor:
                            • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
                            • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
                          5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
                          6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

                          Operating System

                          Linux

                          Deployment

                          Other

                          Anything else?

                          The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

                          • _add_callback_to_queue() creates:

                            file_info=DagFileInfo(
                            rel_path=Path(request.filepath),
                            bundle_path=bundle.path,
                            bundle_name=request.bundle_name,
                            bundle_version=request.bundle_version,
                            )
                          • _refresh_dag_bundles() builds known_files with:

                            DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
                          • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

                          Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

                          That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

                          A small local proof of the mismatch is:

                          frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
                          rel_path=Path("dags/example.py"),
                          bundle_name="main",
                          bundle_path=Path("/tmp/bundle"),
                          bundle_version="v1",
                          )
                          present=DagFileInfo(
                          rel_path=Path("dags/example.py"),
                          bundle_name="main",
                          bundle_path=Path("/tmp/bundle"),
                          )
                          assertversioned!=present

                          The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

                          • bundle_name
                          • rel_path

                          That preserves distinct callback entries per version while preventing false orphan cleanup.

                          Are you willing to submit PR?

                          Yes

                          Metadata

                          Metadata

                          Assignees

                          No one assigned

                            Labels

                            area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

                            Type

                            No type

                            Projects

                            No projects

                              Milestone

                              No milestone

                              Relationships

                              None yet

                              Development

                              No branches or pull requests

                              Issue actions

                              , '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); } })(); })();
                              Skip to content

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

                              Description

                              @hkc-8010

                              Apache Airflow version

                              3.1.8 and current main

                              What happened?

                              We found a case where DAG-level callbacks for versioned DAG bundles are fetched and queued by the dag processor, but then get treated as orphaned before callback execution.

                              The observed sequence is:

                              1. Scheduler creates a DagCallbackRequest
                              2. Dag processor fetches it from the DB
                              3. Dag processor queues it for file processing
                              4. Dag processor starts or prepares callback processing for the file
                              5. The same callback file is later treated as "not present" / orphaned and is purged or the processor is stopped before _execute_dag_callbacks runs

                              This is distinct from the stale serialized DAG metadata problem in #66301 and the fix in #66474.

                              That earlier issue was about SerializedDagModel not refreshing when only bundle version changed. Here, the failure is later in the dag processor manager's callback/orphan cleanup path.

                              What you think should happen instead?

                              If a callback request is queued for bundle_name=X, rel_path=Y, and the same file is still present in the current bundle scan, it should not be treated as removed solely because the callback queue entry includes bundle_version while the scanned "present files" entry does not.

                              In other words, present/orphan checks in the dag processor should treat file presence as:

                              • bundle_name
                              • rel_path

                              and should not require exact DagFileInfo equality including bundle_version.

                              How to reproduce

                              A minimal pattern is:

                              1. Use a DAG bundle implementation with supports_versioning=True
                              2. Have a DAG with a DAG-level on_failure_callback
                              3. Trigger a failing DAG run so a DagCallbackRequest is created
                              4. In the dag processor:
                                • _add_callback_to_queue() creates a DagFileInfo with bundle_version=<version>
                                • _refresh_dag_bundles() builds known_files using DagFileInfo(rel_path=..., bundle_name=..., bundle_path=...) with no bundle_version
                              5. Let purge_removed_files_from_queue() / terminate_orphan_processes() run
                              6. Observe that the versioned callback file can be treated as absent even though the same bundle_name + rel_path is still present in the scanned files

                              Operating System

                              Linux

                              Deployment

                              Other

                              Anything else?

                              The relevant code path in airflow-core/src/airflow/dag_processing/manager.py is:

                              • _add_callback_to_queue() creates:

                                file_info=DagFileInfo(
                                rel_path=Path(request.filepath),
                                bundle_path=bundle.path,
                                bundle_name=request.bundle_name,
                                bundle_version=request.bundle_version,
                                )
                              • _refresh_dag_bundles() builds known_files with:

                                DagFileInfo(rel_path=p, bundle_name=bundle.name, bundle_path=bundle.path)
                              • purge_removed_files_from_queue() and terminate_orphan_processes() then use full DagFileInfo equality against the present set

                              Because DagFileInfo includes bundle_version in equality/hashing, a versioned callback entry and an unversioned scanned entry for the same file do not compare equal.

                              That means the callback queue entry or active processor can be treated as orphaned even though the file is still present.

                              A small local proof of the mismatch is:

                              frompathlibimportPathfromairflow.dag_processing.managerimportDagFileInfoversioned=DagFileInfo(
                              rel_path=Path("dags/example.py"),
                              bundle_name="main",
                              bundle_path=Path("/tmp/bundle"),
                              bundle_version="v1",
                              )
                              present=DagFileInfo(
                              rel_path=Path("dags/example.py"),
                              bundle_name="main",
                              bundle_path=Path("/tmp/bundle"),
                              )
                              assertversioned!=present

                              The fix I tested locally is to keep callback identity version-aware, but make present/orphan checks compare only a stable file-presence key of:

                              • bundle_name
                              • rel_path

                              That preserves distinct callback entries per version while preventing false orphan cleanup.

                              Are you willing to submit PR?

                              Yes

                              Metadata

                              Metadata

                              Assignees

                              No one assigned

                                Labels

                                area:corekind:bugThis is a clearly a bugpriority:highHigh priority bug that should be patched quickly but does not require immediate new release

                                Type

                                No type

                                Projects

                                No projects

                                  Milestone

                                  No milestone

                                  Relationships

                                  None yet

                                  Development

                                  No branches or pull requests

                                  Issue actions