Skip to content

Commit 0e260a7

Browse files
committed
#824 Fixed problem with aggregating items
1 parent 0417cf5 commit 0e260a7

3 files changed

Lines changed: 13 additions & 11 deletions

File tree

src/backend_api/app/api/routes/unidentified_item.py

Lines changed: 1 addition & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -116,7 +116,6 @@ async def get_non_aggregated(
116116

117117
@router.post(
118118
"/add_aggregated/",
119-
response_model=list[schemas.UnidentifiedItem],
120119
dependencies=[Depends(get_current_active_superuser)],
121120
)
122121
async def add_aggregated(
@@ -129,8 +128,4 @@ async def add_aggregated(
129128
And inserts the new items
130129
"""
131130

132-
deleted_non_aggregated = await CRUD_unidentifiedItem.add_aggregated(
133-
db=db, aggregated_objs=aggregated_objs
134-
)
135-
136-
return deleted_non_aggregated
131+
await CRUD_unidentifiedItem.add_aggregated(db=db, aggregated_objs=aggregated_objs)

src/backend_api/app/crud/extensions/crud_unidentifiedItem.py

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -35,12 +35,10 @@ async def get_non_aggregated(self, db: Session) -> list[UnidentifiedItem]:
3535

3636
async def add_aggregated(
3737
self, db: Session, aggregated_objs: list[UnidentifiedItemCreate]
38-
) -> list[UnidentifiedItem]:
38+
):
3939
stmt = delete(model_UnidentifiedItem).where(
4040
model_UnidentifiedItem.aggregated.isnot(True),
4141
)
42-
non_aggregated_objs = db.execute(stmt)
42+
db.execute(stmt)
4343

4444
await self.create(db, obj_in=aggregated_objs, return_nothing=True)
45-
46-
return non_aggregated_objs

src/backend_data_retrieval/data_retrieval_app/external_data_retrieval/data_retrieval/poe_api_handler.py

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -302,6 +302,12 @@ async def _send_n_recursion_requests(
302302
logger.exception(
303303
f"The following exception occured during {self._send_n_recursion_requests.__name__}: {e}"
304304
)
305+
raise
306+
except ProgramTooSlowException:
307+
pass
308+
except ProgramFinished:
309+
pass
310+
finally:
305311
if waiting_for_next_id_lock.locked():
306312
logger.info("Released lock after crash")
307313
if headers is not None:
@@ -321,7 +327,6 @@ async def _send_n_recursion_requests(
321327
)
322328

323329
waiting_for_next_id_lock.release()
324-
raise
325330

326331
async def _follow_stream(
327332
self,
@@ -355,6 +360,10 @@ async def _follow_stream(
355360
stashes = []
356361
stashes_ready_event.set()
357362
await asyncio.sleep(1)
363+
except ProgramTooSlowException:
364+
pass
365+
except ProgramFinished:
366+
pass
358367
except Exception as e:
359368
logger.info(
360369
f"The following exception occured during {self._follow_stream}: {e}"

0 commit comments

Comments
 (0)