@@ -146,11 +146,11 @@ async def store_by_id(client, key, value):
146146 return True
147147
148148
149- def search_by_id (client , doc_id ):
149+ async def search_by_id (client , doc_id ):
150150 if logflag :
151151 logger .info (f"[ search by id ] searching docs of { doc_id } " )
152152 try :
153- results = client .load_document (doc_id )
153+ results = await client .load_document (doc_id )
154154 if logflag :
155155 logger .info (f"[ search by id ] search success of { doc_id } : { results } " )
156156 return results
@@ -174,9 +174,10 @@ def drop_index(index_name, redis_url=REDIS_URL):
174174 return True
175175
176176
177- def delete_by_id (client , id ):
177+ async def delete_by_id (client , id ):
178178 try :
179- assert client .delete_document (id )
179+ res = await client .delete_document (id )
180+ assert res
180181 if logflag :
181182 logger .info (f"[ delete by id ] delete id success: { id } " )
182183 except Exception as e :
@@ -286,7 +287,7 @@ def __init__(self, name: str, description: str, config: dict = None):
286287 self .client = redis .Redis (connection_pool = redis_pool )
287288 self .data_index_client , self .key_index_client = asyncio .run (self ._initialize_client ())
288289 self .embedder = asyncio .run (self ._initialize_embedder ())
289- health_status = self .check_health ()
290+ health_status = asyncio . run ( self .check_health () )
290291 if not health_status :
291292 logger .error ("OpeaRedisDataprep health check failed." )
292293
@@ -390,7 +391,8 @@ async def ingest_files(
390391 # check whether the file already exists
391392 key_ids = None
392393 try :
393- key_ids = search_by_id (self .key_index_client , doc_id ).key_ids
394+ result = await search_by_id (self .key_index_client , doc_id )
395+ key_ids = result .key_ids
394396 if logflag :
395397 logger .info (f"[ redis ingest] File { file .filename } already exists." )
396398 except Exception as e :
@@ -435,7 +437,8 @@ async def ingest_files(
435437 # check whether the link file already exists
436438 key_ids = None
437439 try :
438- key_ids = search_by_id (self .key_index_client , doc_id ).key_ids
440+ result = await search_by_id (self .key_index_client , doc_id )
441+ key_ids = result .key_ids
439442 if logflag :
440443 logger .info (f"[ redis ingest] Link { link } already exists." )
441444 except Exception as e :
@@ -565,7 +568,8 @@ async def delete_files(self, file_path: str = Body(..., embed=True)):
565568
566569 # determine whether this file exists in db KEY_INDEX_NAME
567570 try :
568- key_ids = search_by_id (self .key_index_client , doc_id ).key_ids
571+ result = await search_by_id (self .key_index_client , doc_id )
572+ key_ids = result .key_ids
569573 except Exception as e :
570574 if logflag :
571575 logger .info (f"[ redis delete ] { e } , File { file_path } does not exists." )
@@ -576,7 +580,8 @@ async def delete_files(self, file_path: str = Body(..., embed=True)):
576580
577581 # delete file keys id in db KEY_INDEX_NAME
578582 try :
579- assert delete_by_id (self .key_index_client , doc_id )
583+ res = await delete_by_id (self .key_index_client , doc_id )
584+ assert res
580585 except Exception as e :
581586 if logflag :
582587 logger .info (f"[ redis delete ] { e } . File { file_path } delete failed for db { KEY_INDEX_NAME } ." )
@@ -586,7 +591,7 @@ async def delete_files(self, file_path: str = Body(..., embed=True)):
586591 for file_id in file_ids :
587592 # determine whether this file exists in db INDEX_NAME
588593 try :
589- search_by_id (self .data_index_client , file_id )
594+ await search_by_id (self .data_index_client , file_id )
590595 except Exception as e :
591596 if logflag :
592597 logger .info (f"[ redis delete ] { e } . File { file_path } does not exists." )
@@ -596,7 +601,8 @@ async def delete_files(self, file_path: str = Body(..., embed=True)):
596601
597602 # delete file content
598603 try :
599- assert delete_by_id (self .data_index_client , file_id )
604+ res = await delete_by_id (self .data_index_client , file_id )
605+ assert res
600606 except Exception as e :
601607 if logflag :
602608 logger .info (f"[ redis delete ] { e } . File { file_path } delete failed for db { INDEX_NAME } " )
0 commit comments