Skip to content

Commit a0d2de8

Browse files
committed
update latest chnage to async func
1 parent 669305c commit a0d2de8

1 file changed

Lines changed: 4 additions & 15 deletions

File tree

‎airflow-core/src/airflow/api_fastapi/execution_api/routes/asset_state_store.py‎

Lines changed: 4 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,6 @@
5353
NULL_UUID = UUID(int=0)
5454

5555

56-
<<<<<<< HEAD
5756
class _TIWriterFields(NamedTuple):
5857
dag_id: str
5958
run_id: str
@@ -65,14 +64,9 @@ class _TIWriterFields(NamedTuple):
6564
try_number: int
6665

6766

68-
def _fetch_ti_writer_fields(token: TIToken, session: SessionDep) -> _TIWriterFields:
69-
"""Return exact writer attribution for the execution identified by the token."""
70-
row = session.execute(
71-
=======
7267
async def _fetch_ti_writer_fields(token: TIToken, session: AsyncSessionDep) -> _TIWriterFields:
73-
"""Return (dag_id, run_id, task_id, map_index) for the TI identified by the token."""
74-
result = await session.execute(
75-
>>>>>>> fcf3abbd73 (convert all asset-state-store endpoint to an async endpoint)
68+
"""Return exact writer attribution for the execution identified by the token."""
69+
result_row = await session.execute(
7670
select(
7771
TaskInstance.dag_id,
7872
TaskInstance.run_id,
@@ -84,7 +78,7 @@ async def _fetch_ti_writer_fields(token: TIToken, session: AsyncSessionDep) -> _
8478
TaskInstance.try_number,
8579
).where(TaskInstance.id == token.id)
8680
)
87-
row = result.one_or_none()
81+
row = result_row.one_or_none()
8882
if row is None:
8983
raise HTTPException(
9084
status_code=status.HTTP_404_NOT_FOUND,
@@ -168,12 +162,7 @@ async def _put_asset_state_store(
168162
session=session,
169163
)
170164
else:
171-
<<<<<<< HEAD
172-
writer = _fetch_ti_writer_fields(token, session)
173-
=======
174-
ti_fields = await _fetch_ti_writer_fields(token, session)
175-
dag_id, run_id, task_id, map_index = ti_fields
176-
>>>>>>> fcf3abbd73 (convert all asset-state-store endpoint to an async endpoint)
165+
writer = await _fetch_ti_writer_fields(token, session)
177166

178167
await backend.aset_asset_state_store(
179168
scope,

0 commit comments

Comments
 (0)