kaxil commented on PR #63491: URL: https://github.com/apache/airflow/pull/63491#issuecomment-5766223566
Approval stands, nothing below is gating. ### What's fixed The `_LegacyWorkloadFlag` descriptor closes all three sub-claims from my last round. `BaseExecutor.supports_callbacks` now reads as a real `False` at class level instead of a truthy `property` object, `class Restricted(LocalExecutor): supports_connection_test = False` drops `TEST_CONNECTION` from the inherited frozenset, and instance assignment reaches `__set__` instead of writing a shadowing attribute. `test_workload_families_track_every_workload_type` is the registry guard I was after, and it's a better answer than the hand-maintained checklist it replaced. ### EdgeExecutor queues callbacks but doesn't declare them `EdgeExecutor`'s entire capability declaration is `supports_multi_team: bool = True` at [edge_executor.py:79](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py#L79), so it inherits the base default `frozenset({WorkloadType.EXECUTE_TASK})`. It does handle callbacks: `queue_workload` at [:122](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py#L122) branches on `is_callback_execute(workload)` at [:129](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py#L129) and writes them to `EdgeJobModel` under `EXECUTE_CALLBACK_TAG`. Celery declares both types at [celery_executor.py:53](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/celery/src/airflow/provid ers/celery/executors/celery_executor.py#L53). So after this PR the registry and the implementation disagree, and `EdgeExecutor.supports_callbacks` resolves through the descriptor to `False` for an executor that queues callbacks. Nothing breaks today, and it's pre-existing in the narrow sense that merge-base Edge had no `supports_callbacks = True` either, so that bool read `False` there too. Edge escapes the gate because it overrides `queue_workload` wholesale instead of going through the base, and there is no capability check on the callback dispatch path. What makes it worth fixing in this PR rather than after it: this PR promotes that attribute from an unused bool into the single source of truth, and adds two scheduler sites that already trust it, [:4158](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/jobs/scheduler_job_runner.py#L4158) and [:4229](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/jobs/scheduler_job_runner.py#L4229), both for `TEST_CONNECTION`. The symmetric callback gate is the obvious follow-up this design invites, and it would silently stop Edge callbacks. Four lines mirroring Celery, and both imports are already there at [:43](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py#L43) and [:49-50](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/edge3/src/airflow/ providers/edge3/executors/edge_executor.py#L49): ```python if AIRFLOW_V_3_4_PLUS: supported_workload_types: frozenset[WorkloadType] = frozenset( {WorkloadType.EXECUTE_TASK, WorkloadType.EXECUTE_CALLBACK} ) ``` ### The event-buffer raise is on me Last round I said `process_executor_events` was the one of the three sites that failed quietly, and [the change at scheduler_job_runner.py:1446](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/jobs/scheduler_job_runner.py#L1446) is what I asked for. The mechanism I suggested was wrong. `get_event_buffer` drains destructively ([base_executor.py:640](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/executors/base_executor.py#L640)), so a raise inside the loop discards the remainder of a batch that nothing can re-read, and `_process_executor_events` re-raises ([:1359](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/jobs/scheduler_job_runner.py#L1359)) into a `for executor in self.executors` loop with no per-executor boundary ([:1961](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/jobs/scheduler_job_runner.py#L1961)). Loud was right, fatal on a drained batch was not. I could not find an in-tree producer of a foreign key, so this is latent rather than live. The docstring at [:1403](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/jobs/scheduler_job_runner.py#L1403) still promises "logs error and continues", which is now false for this branch, and the newsfragment does not mention the change. Log at error, count it, skip that one event and carry on would get what I was actually after. The added test asserts only that it raises, so a mixed valid/invalid batch proving the valid events still persist would be worth having alongside it. ### Which legacy declaration forms are meant to be supported? Two forms slip past `__init_subclass__`, and I'd rather ask than assert since neither occurs in tree. Assigning after the class body (`SomeExecutor.supports_callbacks = True` at module scope) lands in the subclass `__dict__` and shadows the descriptor, so it reads back `True` while `supported_workload_types` never changes, and `__init_subclass__` structurally cannot see it. A non-bool declaration is skipped too: the `isinstance(value, bool)` guard at [base_executor.py:256](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/airflow-core/src/airflow/executors/base_executor.py#L256) passes over a `@property` or a truthy sentinel silently, while the method-override loop ten lines below warns on mere presence in `cls.__dict__`. Whichever way you answer, the "still honored" list in the newsfragment is the place to say it. ### One-line consistency nit Six provider executors import `WorkloadType` from the module rather than the package: [celery](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/celery/src/airflow/providers/celery/executors/celery_executor.py#L51), [cncf-kubernetes](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py#L70), [ecs](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/amazon/src/airflow/providers/amazon/aws/executors/ecs/ecs_executor.py#L62), [batch](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/amazon/src/airflow/providers/amazon/aws/executors/batch/batch_executor.py#L46), [lambda](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/amazon/src/airflow/providers/amazon/aws/executors/aws_lambda/lambda_executor.py#L54) and [edge3](https://github.com/apache/airflow/blob/bd7b8d1ec4b4dbd29cd2dba2dcc9c063695d141d/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py#L50). Core imports it from `airflow.executors.workloads`, which lists it in `__all__`. Worth normalising. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
