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]