Skip to content

Commit cc73b3f

Browse files
committed
merge
2 parents 33b36f2 + c32061f commit cc73b3f

4 files changed

Lines changed: 29 additions & 37 deletions

File tree

src/connectors/google_drive_acl.py

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -100,10 +100,7 @@ async def _get_cloud_identity_group_roles(
100100
roles: list[str] = []
101101
seen: set[str] = set()
102102
page_token: str | None = None
103-
query = (
104-
f"member_key_id == '{_cel_string(user_email)}' && "
105-
f"'{GOOGLE_GROUP_LABEL}' in labels"
106-
)
103+
query = f"member_key_id == '{_cel_string(user_email)}' && '{GOOGLE_GROUP_LABEL}' in labels"
107104

108105
try:
109106
while True:

src/connectors/langflow_connector_service.py

Lines changed: 27 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -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
}

src/connectors/onedrive/connector.py

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -488,9 +488,7 @@ async def _extract_onedrive_acl(self, file_id: str, file_metadata: Dict) -> Docu
488488

489489
# Granted to identities (can include users and groups)
490490
identities = (
491-
perm.get("grantedToIdentitiesV2")
492-
or perm.get("grantedToIdentities")
493-
or []
491+
perm.get("grantedToIdentitiesV2") or perm.get("grantedToIdentities") or []
494492
)
495493
if identities:
496494
for identity in identities:

src/session_manager.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -254,7 +254,6 @@ def create_opensearch_jwt_token(
254254
"""Create a short-lived OpenSearch JWT with current connector groups."""
255255
if ttl_seconds is None:
256256
from config.settings import get_opensearch_jwt_ttl_seconds
257-
258257
ttl_seconds = get_opensearch_jwt_ttl_seconds()
259258
return self._create_signed_jwt_token(
260259
user,

0 commit comments

Comments
 (0)