kaxil commented on code in PR #73235:
URL: https://github.com/apache/airflow/pull/73235#discussion_r4109116862


##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/hitl.py:
##########
@@ -61,34 +62,43 @@ def upsert_hitl_detail(
     1. If a HITLOperator task instance does not have a HITLDetail,
        a new HITLDetail is created without a response section.
     2. If a HITLOperator task instance has a HITLDetail but lacks a response,
-       the existing HITLDetail is returned.
+       the request part is refreshed from the payload and the HITLDetail is 
returned.
        This situation occurs when a task instance is cleared before a response 
is received.
     3. If a HITLOperator task instance has both a HITLDetail and a response 
section,
-       the existing response is removed, and the HITLDetail is returned.
+       the request part is refreshed, the existing response is removed, and 
the HITLDetail is returned.
        This happens when a task instance is cleared after a response has been 
received.
        This design ensures that each task instance has only one HITLDetail.
     """
-    hitl_detail_model = 
session.scalar(select(HITLDetail).where(HITLDetail.ti_id == task_instance_id))
+    request_part = {
+        "options": payload.options,
+        "subject": payload.subject,
+        "body": payload.body,
+        "defaults": payload.defaults,
+        "multiple": payload.multiple,
+        "params": payload.params,
+        "assignees": [user.model_dump() for user in payload.assigned_users],
+        "created_at": timezone.utcnow(),
+    }
+    # Same lock order as the park transition and the Core API response path 
(TaskInstance, then
+    # hitl_detail), so a response landing between the read and the flush 
cannot survive the rewrite.
+    session.get(TI, task_instance_id, with_for_update={"of": TI})

Review Comment:
   These locks close the gap I raised last round, but the Core API response 
handler I pointed to as the reference doesn't re-read the row after it locks. 
It eager-loads `hitl_detail` at `core_api/routes/public/hitl.py:160`, throws 
away the result of its locked select at `:180`, and checks assignees (`:204`) 
and options (`:214`) against the object it loaded at `:160`. So if a reviewer 
submits while this upsert holds the locks, their answer is checked against the 
previous attempt's `assigned_users` and `options` and then written onto the 
refreshed request at `:228`, where the park transition resumes straight onto 
it. That stale read was harmless while these columns never changed after 
insert. Would it work to replace `:180-185` with `hitl_detail_model = 
session.scalar(select(HITLDetailModel).where(HITLDetailModel.ti_id == 
task_instance.id).with_for_update(of=HITLDetailModel).execution_options(populate_existing=True))`?



-- 
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