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


##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/hitl.py:
##########
@@ -175,14 +175,15 @@ def update_hitl_detail(
     )
     # Lock the hitl_detail row (FOR UPDATE OF hitl_detail). of= scopes the 
lock to hitl_detail, which
     # eager-joins task_instance (lazy="joined"); a bare with_for_update() 
would emit FOR UPDATE against
-    # the nullable side of that outer join, which Postgres rejects. The 
joinedloaded relationship object
-    # reused below is the same identity-mapped row, now locked for this 
transaction.
-    session.execute(
+    # the nullable side of that outer join, which Postgres rejects. 
populate_existing re-reads the
+    # joinedloaded row under the lock, so assignees and options are validated 
against the state committed
+    # by a concurrent clear that held the lock, not the snapshot taken before 
locking.

Review Comment:
   Reworded to point at the concurrent upsert from the re-run.



##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/hitl.py:
##########
@@ -175,14 +175,15 @@ def update_hitl_detail(
     )
     # Lock the hitl_detail row (FOR UPDATE OF hitl_detail). of= scopes the 
lock to hitl_detail, which
     # eager-joins task_instance (lazy="joined"); a bare with_for_update() 
would emit FOR UPDATE against
-    # the nullable side of that outer join, which Postgres rejects. The 
joinedloaded relationship object
-    # reused below is the same identity-mapped row, now locked for this 
transaction.
-    session.execute(
+    # the nullable side of that outer join, which Postgres rejects. 
populate_existing re-reads the
+    # joinedloaded row under the lock, so assignees and options are validated 
against the state committed
+    # by a concurrent clear that held the lock, not the snapshot taken before 
locking.
+    hitl_detail_model = session.scalars(
         select(HITLDetailModel)
         .where(HITLDetailModel.ti_id == task_instance.id)
         .with_for_update(of=HITLDetailModel)
-    )
-    hitl_detail_model = task_instance.hitl_detail
+        .execution_options(populate_existing=True)

Review Comment:
   Added `populate_existing=True` to the locked `session.get` so the TI state 
refresh no longer relies on the joined load.



##########
airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_hitl.py:
##########
@@ -113,26 +115,48 @@ def test_upsert_hitl_detail(
     session.commit()
 
     if existing_hitl_detail_args:
-        session.add(HITLDetail(ti_id=ti.id, **existing_hitl_detail_args))
+        session.add(
+            HITLDetail(
+                ti_id=ti.id,
+                created_at=convert_to_utc(datetime(2025, 7, 1, 0, 0, 0)),
+                **existing_hitl_detail_args,
+            )
+        )
         session.commit()
 
+    request_kwargs = {
+        **default_hitl_detail_request_kwargs,
+        "subject": "Regenerated subject",
+        "body": "regenerated body",
+        "defaults": ["Reject"],
+        "params": {"input_1": 3},

Review Comment:
   The request now sends new options, `multiple=True` and `assigned_users` (the 
old `assignees` key was dropped), and the test asserts all three after the 
upsert.



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