Skip to content

Extract shared request handlers from supervisor _handle_request methods #65570

Description

@ferruzzi

We have four WatchedSubprocess subclasses that each implement _handle_request with large if/elif chains dispatching comms messages. Many of these handlers are duplicated inline across multiple supervisors:

  • ActivitySubprocess (task-sdk, supervisor.py)
  • CallbackSubprocess (task-sdk, callback_supervisor.py)
  • TriggerRunnerSupervisor (airflow-core, triggerer_job_runner.py)
  • DagFileProcessorProcess (airflow-core, processor.py)

A request_handlers.py module already exists in task-sdk/src/airflow/sdk/execution_time/ with shared handlers for GetConnection, GetVariable, GetAssetByName, GetAssetByUri, and MaskSecret. The Callback and Activity supervisors use these, but the Triggerer and DFP still have inline duplicates.

The Task:

Extract the remaining duplicated handler logic into shared functions in request_handlers.py and have all four supervisors call them.

Note: The implementations across supervisors are similar but not identical. Some include isinstance() guards on the response; some pass different parameter subsets; some set dump_opts while others don't. When extracting, don't just copy one supervisor's version verbatim. Compare all implementations of each handler and produce a single version that incorporates the best practices from each (e.g. proper response type checking, consistent dump_opts, full parameter forwarding).

Here is my suggestion, but you may come up with a better set of rules once you start working on it:

  1. Use the existing shared handlers as a guide, but don't be afraid to modify them if you come up with a new standard.
  2. Each new handler should be a standalone function with this signature:
    def handle_<message_type>(client: Client, msg: <MessageType>) -> tuple[BaseModel | None, dict[str, bool]]:
  3. Always guard the response with isinstance before converting it to a comms result model. The API client can return ErrorResponse on failure; when that happens, pass it through unchanged.
  4. Always return dump_opts as the second tuple element. Use {"exclude_unset": True} for result models that wrap API responses to avoid serializing None fields. Add {"by_alias": True} when the result model uses field aliases (currently only ConnectionResult). For simple pass-through responses or error responses, return {}.
  5. Forward all message fields to the client call. Some supervisors currently pass fewer parameters than others for the same message type (e.g. the Triggerer's GetXCom omits include_prior_dates). The shared handler should forward every field that the message carries and let the client and server handle defaults for optional fields.
  6. Always mask secrets. If the response contains sensitive data (passwords, tokens, variable values), call mask_secret() before returning. Look at handle_get_connection and handle_get_variable for examples.
  7. Handlers that don't return a response such as PutVariable and DeleteXCom can use a simpler signature. return (None, {}) or return None directly like handle_mask_secret does. Pick whichever is cleaner for the specific case, but be consistent within a batch.
  8. Keep handlers stateless. They should not touch subprocess internals (exit codes, terminal states, etc.). If a message type requires updating subprocess state, it belongs inline in that supervisor's _handle_request, not in a shared handler.

Phase 1: High-overlap handlers (three or more supervisors share this logic)

  • GetXCom - Activity, Triggerer, DFP
  • PutVariable - Activity, Triggerer, DFP
  • DeleteVariable - Activity, Triggerer, DFP
  • GetTICount - Activity, Triggerer, DFP
  • GetTaskStates - Activity, Triggerer, DFP
  • GetPreviousTI - Activity, Triggerer, DFP

Phase 2: Medium-overlap handlers (two supervisors share these)

  • SetXCom - Activity, Triggerer
  • DeleteXCom - Activity, Triggerer
  • GetDRCount - Activity, Triggerer
  • GetDagRunState - Activity, Triggerer
  • GetPreviousDagRun - Activity, DFP
  • GetPrevSuccessfulDagRun - Activity, DFP
  • GetXComCount - Activity, DFP
  • GetXComSequenceItem - Activity, DFP
  • GetXComSequenceSlice - Activity, DFP

Phase 2.5: Callback Supervisor

Comms channels were added to the Callback supervisor in #65269 and all of the comms channels it has as of the time I am writing this are already in shared helpers, but have a look and make sure that is still true while you are doing this.

Phase 3: Migrate DagFileProcessorProcess onto shared handlers

The DFP currently has fully inline versions of GetConnection, GetVariable, and MaskSecret even though shared handlers already exist. One thing to note: the DFP's inline GetConnection handler skips the mask_secret() calls on password and extra that the shared handler performs, which I personally feel is a bug but may require discussion. Maybe there's a reason I'm not aware of.

Out of Scope:

Supervisor-specific messages (e.g. TaskState, DeferTask, TriggerStateChanges, DagFileParsingResult) should either stay inline or move to a supervisor-specific helper/utils module (I vote leave alone, personally) since they interact with internal subprocess state.

Activity

  1. added theissue type on Apr 20, 2026
  2. removed theissue type on Apr 20, 2026
  3. leeyspaul commented on Apr 20, 2026

    @leeyspaul
    Contributor

    Interesting, I would like to take this one on! I'll be working on it, and report back with further questions or a PR.

  4. ferruzzi commented on Apr 20, 2026

    @ferruzzi
    ContributorAuthor

    MOST of it should be copypasta to reduce redundant code into helpers, but it may take some discussion, feel free to ask for advice in the community slack server if anything looks odd.

  5. leeyspaul commented on Apr 20, 2026

    @leeyspaul
    Contributor

    Hey @ferruzzi , I'm currently ramping up on context of the mentioned code but I'm not seeing a request_handlers.py under the mentioned task-sdk/src/airflow/sdk/execution_time/. Is this something we want created here or is it already present on another branch/PR?

  6. ferruzzi commented on Apr 20, 2026

    @ferruzzi
    ContributorAuthor

    Got ahead of myself. It will be added when #65269 merges. Sorry about that 😅

  7. leeyspaul commented on Apr 20, 2026

    @leeyspaul
    Contributor

    Thank you, that helps! I'll assume #65269 lands first and implement this on top of that structure after rebasing.

  8. ferruzzi commented on Apr 21, 2026

    @ferruzzi
    ContributorAuthor

    Yeah, #65269 should theoretically merge in the next day or three, so you can pull that down and use that as your upstream if you want, or just assume it'll be rebased in later. Whatever fits your workflow best.

  9. leeyspaul commented on Apr 24, 2026

    @leeyspaul
    Contributor

    @ferruzzi this PR: #65624 is now ready for review 👀

  10. ferruzzi commented on Apr 24, 2026

    @ferruzzi
    ContributorAuthor

    I just created a follow-up Issue for after this one is done. I'm going to be away on vacation, but if you or someone else wants to tackle it while I'm away, it's the next logical step after this gets merged.

  11. leeyspaul commented on Apr 24, 2026

    @leeyspaul
    Contributor

    Yeah, I'm happy to take it up as it builds on the previous. Have a good vacation!

  12. Quantum0uasar commented on May 21, 2026

    @Quantum0uasar

    Hi! I'd love to help with this refactoring effort.

    My understanding of the task:

    • Extract duplicated _handle_request handler logic from ActivitySubprocess, CallbackSubprocess, TriggerRunnerSupervisor, and DagFileProcessorProcess into shared handlers in request_handlers.py
    • Follow the pattern already established in task-sdk/src/airflow/sdk/execution_time/request_handlers.py

    I'll start by:

    1. Auditing the if/elif chains in each supervisor's _handle_request to identify common patterns
    2. Creating appropriate handler functions or a dispatch table in a shared module
    3. Updating the supervisors to call the shared handlers

    Would this be best done incrementally (one supervisor at a time) or all at once? Happy to open a draft PR with the first supervisor extracted to get early feedback!

  13. ferruzzi commented on May 21, 2026

    @ferruzzi
    ContributorAuthor

    @leeyspaul - with your PR merged, this is completed, right?

  14. leeyspaul commented on May 21, 2026

    @leeyspaul
    Contributor

    Yes, I believe so!

  15. ferruzzi commented on May 21, 2026

    @ferruzzi
    ContributorAuthor

    Thanks for the work! Closing.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions