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]