@@ -99,7 +99,7 @@ async def process_connector_document(
9999 # Step 1: Upload file to Langflow
100100 logger .debug ("Uploading file to Langflow" , filename = document .filename )
101101 content = document .content
102-
102+
103103 # Clean filename and ensure we don't add a double extension
104104 processed_filename = clean_connector_filename (document .filename , document .mimetype )
105105
@@ -114,19 +114,30 @@ async def process_connector_document(
114114 if self .session_manager :
115115 try :
116116 from config .settings import get_index_name
117- opensearch_client = self .session_manager .get_user_opensearch_client (owner_user_id , jwt_token )
117+
118+ opensearch_client = self .session_manager .get_user_opensearch_client (
119+ owner_user_id , jwt_token
120+ )
118121 delete_body = {"query" : {"term" : {"filename" : processed_filename }}}
119- delete_result = await opensearch_client .delete_by_query (index = get_index_name (), body = delete_body )
122+ delete_result = await opensearch_client .delete_by_query (
123+ index = get_index_name (), body = delete_body
124+ )
120125 deleted_count = delete_result .get ("deleted" , 0 )
121- logger .info ("Deleted existing chunks before re-ingestion" , filename = processed_filename , deleted_count = deleted_count )
126+ logger .info (
127+ "Deleted existing chunks before re-ingestion" ,
128+ filename = processed_filename ,
129+ deleted_count = deleted_count ,
130+ )
122131 except Exception as delete_err :
123- logger .warning ("Failed to delete existing chunks before re-ingestion" , filename = processed_filename , error = str (delete_err ))
132+ logger .warning (
133+ "Failed to delete existing chunks before re-ingestion" ,
134+ filename = processed_filename ,
135+ error = str (delete_err ),
136+ )
124137
125138 langflow_file_id = None # Initialize to track if upload succeeded
126139 try :
127- upload_result = await self .langflow_service .upload_user_file (
128- file_tuple , jwt_token
129- )
140+ upload_result = await self .langflow_service .upload_user_file (file_tuple , jwt_token )
130141 langflow_file_id = upload_result ["id" ]
131142 langflow_file_path = upload_result ["path" ]
132143
@@ -137,9 +148,7 @@ async def process_connector_document(
137148 )
138149
139150 # Step 2: Run ingestion flow with the uploaded file
140- logger .debug (
141- "Running Langflow ingestion flow" , file_path = langflow_file_path
142- )
151+ logger .debug ("Running Langflow ingestion flow" , file_path = langflow_file_path )
143152
144153 connector_tweak_settings = None
145154 if isinstance (ingest_settings , dict ):
@@ -222,7 +231,6 @@ async def process_connector_document(
222231 )
223232 raise
224233
225-
226234 async def sync_connector_files (
227235 self ,
228236 connection_id : str ,
@@ -246,9 +254,7 @@ async def sync_connector_files(
246254
247255 connector = await self .get_connector (connection_id )
248256 if not connector :
249- raise ValueError (
250- f"Connection '{ connection_id } ' not found or not authenticated"
251- )
257+ raise ValueError (f"Connection '{ connection_id } ' not found or not authenticated" )
252258
253259 logger .debug ("Got connector" , authenticated = connector .is_authenticated )
254260
@@ -264,13 +270,9 @@ async def sync_connector_files(
264270
265271 while True :
266272 # List files from connector with limit
267- logger .debug (
268- "Calling list_files" , page_size = page_size , page_token = page_token
269- )
273+ logger .debug ("Calling list_files" , page_size = page_size , page_token = page_token )
270274 file_list = await connector .list_files (page_token , limit = page_size )
271- logger .debug (
272- "Got files from connector" , file_count = len (file_list .get ("files" , []))
273- )
275+ logger .debug ("Got files from connector" , file_count = len (file_list .get ("files" , [])))
274276 files = file_list ["files" ]
275277
276278 if not files :
@@ -333,7 +335,7 @@ async def sync_specific_files(
333335 """
334336 Sync specific files by their IDs using Langflow processing.
335337 Automatically expands folders to their contents.
336-
338+
337339 Args:
338340 connection_id: The connection ID
339341 user_id: The user ID
@@ -351,9 +353,7 @@ async def sync_specific_files(
351353
352354 connector = await self .get_connector (connection_id )
353355 if not connector :
354- raise ValueError (
355- f"Connection '{ connection_id } ' not found or not authenticated"
356- )
356+ raise ValueError (f"Connection '{ connection_id } ' not found or not authenticated" )
357357
358358 if not connector .is_authenticated :
359359 raise ValueError (f"Connection '{ connection_id } ' not authenticated" )
@@ -368,7 +368,7 @@ async def sync_specific_files(
368368
369369 # If file_infos provided, cache them in the connector for later use
370370 # This allows get_file_content to use download URLs directly
371- if file_infos and hasattr (connector , ' set_file_infos' ):
371+ if file_infos and hasattr (connector , " set_file_infos" ):
372372 connector .set_file_infos (file_infos )
373373 logger .info (f"Cached { len (file_infos )} file infos with download URLs in connector" )
374374
@@ -435,9 +435,7 @@ async def sync_specific_files(
435435 original_filenames = {}
436436 if file_infos :
437437 original_filenames = {
438- f ["id" ]: clean_connector_filename (
439- f ["name" ], f .get ("mimeType" ) or f .get ("mimetype" )
440- )
438+ f ["id" ]: clean_connector_filename (f ["name" ], f .get ("mimeType" ) or f .get ("mimetype" ))
441439 for f in file_infos
442440 if "id" in f and "name" in f
443441 }
0 commit comments