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]

Reply via email to