kaxil commented on code in PR #73235:
URL: https://github.com/apache/airflow/pull/73235#discussion_r4139368334
##########
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:
Small wording nit: a clear never writes these columns.
`prepare_db_for_next_try` gives the TI a new id with `uuid7()`, and this row
follows it through the `onupdate="CASCADE"` FK, so if a clear wins the race
this select finds no row and `.one()` raises (a 500, same as before this PR).
The writer that actually rewrites options and assignees under us is the
Execution API upsert when the re-run parks again, so maybe say "validated
against the request committed by a concurrent upsert from the re-run".
##########
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:
`options` and `multiple` are the same in the stored row and in this request,
and nothing asserts on assignees, so taking any of those three out of
`request_part` would still pass. Also, `"assignees"` isn't a field on
`HITLDetailRequest` (it's `assigned_users`, and the model ignores unknown
keys), so that key is silently dropped from the POST. Seeding the existing row
with a different assignee, a different options list and `multiple=True`, then
asserting all three after the upsert, would pin the columns the Core handler
validates against, including the one it authorizes on.
##########
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:
This `populate_existing` is also what keeps `locked_ti.state` fresh for the
resume gate below. `session.get(TI, ..., with_for_update=...)` on L172 skips
the identity-map shortcut and emits the locking SELECT, but without
`populate_existing` it does not overwrite attributes on the already-loaded
`task_instance`. Its `state` is current only because this select eager-joins
`task_instance` and repopulates it. Dropping either the joined load or this
option would bring back a stale-state check at L239 with nothing failing.
Passing `populate_existing=True` to the `session.get` on L172 would make that
explicit rather than a side effect.
--
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]