Skip to content

Commit 853df1e

Browse files
committed
init commit to add support for service-level max queue memory size
1 parent 1074a59 commit 853df1e

10 files changed

Lines changed: 118 additions & 26 deletions

File tree

‎connectors/config.py‎

Lines changed: 19 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -9,12 +9,12 @@
99
from envyaml import EnvYAML
1010

1111
from connectors.logger import logger
12-
13-
DEFAULT_ELASTICSEARCH_MAX_RETRIES = 5
14-
DEFAULT_ELASTICSEARCH_RETRY_INTERVAL = 10
15-
16-
DEFAULT_MAX_FILE_SIZE = 10485760 # 10MB
17-
12+
from connectors.utils import (
13+
DEFAULT_MAX_FILE_SIZE,
14+
DEFAULT_QUEUE_MEM_SIZE,
15+
DEFAULT_ELASTICSEARCH_MAX_RETRIES,
16+
DEFAULT_ELASTICSEARCH_RETRY_INTERVAL
17+
)
1818

1919
def load_config(config_file):
2020
logger.info(f"Loading config from {config_file}")
@@ -67,7 +67,7 @@ def _default_config():
6767
"verify_certs": True,
6868
"bulk": {
6969
"queue_max_size": 1024,
70-
"queue_max_mem_size": 25,
70+
"queue_max_mem_size": 25 * 1024 * 1024, # 25 MB
7171
"queue_refresh_interval": 1,
7272
"queue_refresh_timeout": 600,
7373
"display_every": 100,
@@ -108,6 +108,7 @@ def _default_config():
108108
"max_concurrent_access_control_syncs": 1,
109109
"max_concurrent_scheduling_tasks": 4,
110110
"max_file_download_size": DEFAULT_MAX_FILE_SIZE,
111+
"queue_max_mem_size": DEFAULT_QUEUE_MEM_SIZE,
111112
"job_cleanup_interval": 300,
112113
"log_level": "INFO",
113114
},
@@ -220,19 +221,28 @@ class DataSourceFrameworkConfig:
220221
preventing them from requiring substantial changes to access new configs that may be added.
221222
"""
222223

223-
def __init__(self, max_file_size):
224+
def __init__(self, max_file_size, queue_max_mem_size):
224225
"""
225226
Should not be called directly. Use the Builder.
226227
"""
227228
self.max_file_size = max_file_size
229+
self.max_queue_mem_size = queue_max_mem_size
228230

229231
class Builder:
230232
def __init__(self):
231233
self.max_file_size = DEFAULT_MAX_FILE_SIZE
234+
self.max_queue_mem_size = DEFAULT_QUEUE_MEM_SIZE
232235

233236
def with_max_file_size(self, max_file_size):
234237
self.max_file_size = max_file_size
235238
return self
236239

240+
def with_max_queue_memory_size(self, queue_max_mem_size):
241+
self.max_queue_mem_size = queue_max_mem_size
242+
return self
243+
237244
def build(self):
238-
return DataSourceFrameworkConfig(self.max_file_size)
245+
return DataSourceFrameworkConfig(
246+
self.max_file_size,
247+
self.max_queue_mem_size
248+
)

‎connectors/es/sink.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1023,7 +1023,7 @@ async def async_bulk(
10231023

10241024
stream = MemQueue(
10251025
maxsize=queue_size,
1026-
maxmemsize=queue_mem_size * 1024 * 1024,
1026+
maxmemsize=queue_mem_size,
10271027
refresh_timeout=mem_queue_refresh_timeout,
10281028
refresh_interval=mem_queue_refresh_interval,
10291029
)

‎connectors/sources/box.py‎

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -475,7 +475,21 @@ async def _consumer(self):
475475
else:
476476
yield item
477477

478-
async def get_docs(self, filtering=None):
478+
def _update_queue_max_mem_size(self):
479+
"""
480+
Helper function to update the queue's max memory size
481+
if the configured value exceeds the constant defined in the data source
482+
"""
483+
max_mem_size_from_config = self.framework_config.max_queue_mem_size
484+
485+
if max_mem_size_from_config > QUEUE_MEM_SIZE:
486+
self.queue.update_maxmemsize(max_mem_size_from_config)
487+
logger.debug(f"MemQueue max memory size updated from {QUEUE_MEM_SIZE} to {max_mem_size_from_config}")
488+
489+
async def get_docs(self, filtering=None): # type: ignore
490+
# update queue max memory size if needed
491+
self._update_queue_max_mem_size()
492+
479493
seen_ids = set()
480494
root_folder = "0"
481495
if self.is_enterprise == BOX_ENTERPRISE:

‎connectors/sources/confluence.py‎

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -767,7 +767,7 @@ async def get_user(self):
767767
user=user
768768
)
769769

770-
async def get_access_control(self):
770+
async def get_access_control(self): # type: ignore
771771
"""Get access control documents for active Atlassian users.
772772
773773
This method fetches access control documents for active Atlassian users when document level security (DLS)
@@ -1277,14 +1277,28 @@ async def _consumer(self):
12771277
else:
12781278
yield item
12791279

1280-
async def get_docs(self, filtering=None):
1280+
def _update_queue_max_mem_size(self):
1281+
"""
1282+
Helper function to update the queue's max memory size
1283+
if the configured value exceeds the constant defined in the data source
1284+
"""
1285+
max_mem_size_from_config = self.framework_config.max_queue_mem_size
1286+
1287+
if max_mem_size_from_config > QUEUE_MEM_SIZE:
1288+
self.queue.update_maxmemsize(max_mem_size_from_config)
1289+
logger.debug(f"MemQueue max memory size updated from {QUEUE_MEM_SIZE} to {max_mem_size_from_config}")
1290+
1291+
async def get_docs(self, filtering=None): # type: ignore
12811292
"""Executes the logic to fetch Confluence content in async manner.
12821293
12831294
Args:
12841295
filtering (Filtering): Object of Class Filtering
12851296
Yields:
12861297
dictionary: dictionary containing meta-data of the content.
12871298
"""
1299+
# update queue max memory size if needed
1300+
self._update_queue_max_mem_size()
1301+
12881302
self._logger.info("Successfully connected to Confluence")
12891303
if filtering and filtering.has_advanced_rules():
12901304
advanced_rules = filtering.get_advanced_rules()

‎connectors/sources/jira.py‎

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -729,7 +729,7 @@ async def _issue_access_control(self, issue_key, project):
729729
return project_access_controls
730730
return list(access_control)
731731

732-
async def get_access_control(self):
732+
async def get_access_control(self): # type: ignore
733733
"""Get access control documents for active Atlassian users.
734734
735735
This method fetches access control documents for active Atlassian users when document level security (DLS)
@@ -1017,7 +1017,18 @@ async def _consumer(self):
10171017
else:
10181018
yield item
10191019

1020-
async def get_docs(self, filtering=None):
1020+
def _update_queue_max_mem_size(self):
1021+
"""
1022+
Helper function to update the queue's max memory size
1023+
if the configured value exceeds the constant defined in the data source
1024+
"""
1025+
max_mem_size_from_config = self.framework_config.max_queue_mem_size
1026+
1027+
if max_mem_size_from_config > QUEUE_MEM_SIZE:
1028+
self.queue.update_maxmemsize(max_mem_size_from_config)
1029+
logger.debug(f"MemQueue max memory size updated from {QUEUE_MEM_SIZE} to {max_mem_size_from_config}")
1030+
1031+
async def get_docs(self, filtering=None): # type: ignore
10211032
"""Executes the logic to fetch jira objects in async manner
10221033
10231034
Args:
@@ -1026,6 +1037,9 @@ async def get_docs(self, filtering=None):
10261037
Yields:
10271038
dictionary: dictionary containing meta-data of the files.
10281039
"""
1040+
# update queue max memory size if needed
1041+
self._update_queue_max_mem_size()
1042+
10291043
self.custom_fields = await anext(self.jira_client.get_jira_fields())
10301044

10311045
if filtering and filtering.has_advanced_rules():

‎connectors/sources/microsoft_teams.py‎

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -452,7 +452,7 @@ async def _get(self, absolute_url, use_token=True):
452452
url=absolute_url,
453453
) as resp:
454454
yield resp
455-
except aiohttp.client_exceptions.ClientOSError:
455+
except aiohttp.client_exceptions.ClientOSError: # type: ignore
456456
self._logger.error(
457457
"Graph API dropped the connection. It might indicate, that connector makes too many requests - decrease concurrency settings, otherwise Graph API can block this app."
458458
)
@@ -1267,7 +1267,18 @@ async def calendars_producer(self, user):
12671267
)
12681268
await self.queue.put(EndSignal.CALENDAR_TASK_FINISHED)
12691269

1270-
async def get_docs(self, filtering=None):
1270+
def _update_queue_max_mem_size(self):
1271+
"""
1272+
Helper function to update the queue's max memory size
1273+
if the configured value exceeds the constant defined in the data source
1274+
"""
1275+
max_mem_size_from_config = self.framework_config.max_queue_mem_size
1276+
1277+
if max_mem_size_from_config > QUEUE_MEM_SIZE:
1278+
self.queue.update_maxmemsize(max_mem_size_from_config)
1279+
logger.debug(f"MemQueue max memory size updated from {QUEUE_MEM_SIZE} to {max_mem_size_from_config}")
1280+
1281+
async def get_docs(self, filtering=None): # type: ignore
12711282
"""Executes the logic to fetch Microsoft Teams objects in async manner
12721283
12731284
Args:
@@ -1276,6 +1287,9 @@ async def get_docs(self, filtering=None):
12761287
Yields:
12771288
dictionary: dictionary containing meta-data of the files.
12781289
"""
1290+
# update queue max memory size if needed
1291+
self._update_queue_max_mem_size()
1292+
12791293
async for chats in self.client.get_user_chats():
12801294
for chat in chats:
12811295
await self.fetchers.put(partial(self.user_chat_producer, chat))

‎connectors/sources/outlook.py‎

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -50,8 +50,6 @@
5050
RETRIES = 3
5151
RETRY_INTERVAL = 2
5252

53-
QUEUE_MEM_SIZE = 5 * 1024 * 1024 # Size in Megabytes
54-
5553
OUTLOOK_SERVER = "outlook_server"
5654
OUTLOOK_CLOUD = "outlook_cloud"
5755
API_SCOPE = "https://graph.microsoft.com/.default"
@@ -821,7 +819,7 @@ def _dls_enabled(self):
821819

822820
return self.configuration["use_document_level_security"]
823821

824-
async def get_access_control(self):
822+
async def get_access_control(self): # type: ignore
825823
if not self._dls_enabled():
826824
self._logger.warning("DLS is not enabled. Skipping")
827825
return
@@ -1064,7 +1062,7 @@ async def ping(self):
10641062
await self.client.ping()
10651063
self._logger.info("Successfully connected to Outlook")
10661064

1067-
async def get_docs(self, filtering=None):
1065+
async def get_docs(self, filtering=None): # type: ignore
10681066
"""Executes the logic to fetch outlook objects in async manner
10691067
10701068
Args:

‎connectors/sources/servicenow.py‎

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -384,7 +384,7 @@ async def validate(self, advanced_rules):
384384

385385
async def _remote_validation(self, advanced_rules):
386386
try:
387-
ServiceNowAdvancedRulesValidator.SCHEMA(advanced_rules)
387+
ServiceNowAdvancedRulesValidator.SCHEMA(advanced_rules) # type: ignore
388388
except fastjsonschema.JsonSchemaValueException as e:
389389
return SyncRuleValidationResult(
390390
rule_id=SyncRuleValidationResult.ADVANCED_RULES,
@@ -567,7 +567,7 @@ async def _fetch_users_by_roles(self, role):
567567
):
568568
yield user
569569

570-
async def get_access_control(self):
570+
async def get_access_control(self): # type: ignore
571571
if not self._dls_enabled():
572572
self._logger.warning("DLS is not enabled. Skipping")
573573
return
@@ -858,7 +858,18 @@ async def _consumer(self):
858858
else:
859859
yield item
860860

861-
async def get_docs(self, filtering=None):
861+
def _update_queue_max_mem_size(self):
862+
"""
863+
Helper function to update the queue's max memory size
864+
if the configured value exceeds the constant defined in the data source
865+
"""
866+
max_mem_size_from_config = self.framework_config.max_queue_mem_size
867+
868+
if max_mem_size_from_config > QUEUE_MEM_SIZE:
869+
self.queue.update_maxmemsize(max_mem_size_from_config)
870+
logger.debug(f"MemQueue max memory size updated from {QUEUE_MEM_SIZE} to {max_mem_size_from_config}")
871+
872+
async def get_docs(self, filtering=None): # type: ignore
862873
"""Get documents from ServiceNow.
863874
864875
Args:
@@ -867,6 +878,8 @@ async def get_docs(self, filtering=None):
867878
Yields:
868879
dict: Documents from ServiceNow.
869880
"""
881+
# update queue max memory size if needed
882+
self._update_queue_max_mem_size()
870883

871884
self._logger.info("Fetching ServiceNow data")
872885
if filtering and filtering.has_advanced_rules():

‎connectors/sync_job_runner.py‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -218,9 +218,14 @@ async def execute(self):
218218
await self.data_provider.close()
219219

220220
def _data_source_framework_config(self):
221-
builder = DataSourceFrameworkConfig.Builder().with_max_file_size(
221+
builder = DataSourceFrameworkConfig.Builder()
222+
223+
builder.with_max_file_size(
222224
self.service_config.get("max_file_download_size")
223225
)
226+
builder.with_max_queue_memory_size(
227+
self.service_config.get("queue_max_mem_size")
228+
)
224229
return builder.build()
225230

226231
async def _execute_access_control_sync_job(self, job_type, bulk_options):

‎connectors/utils.py‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,10 +36,13 @@
3636
DEFAULT_CHUNK_SIZE = 500
3737
DEFAULT_QUEUE_SIZE = 1024
3838
DEFAULT_DISPLAY_EVERY = 100
39-
DEFAULT_QUEUE_MEM_SIZE = 5
39+
DEFAULT_QUEUE_MEM_SIZE = 5242880 # 5 MB
4040
DEFAULT_CHUNK_MEM_SIZE = 25
4141
DEFAULT_MAX_CONCURRENCY = 5
4242
DEFAULT_CONCURRENT_DOWNLOADS = 10
43+
DEFAULT_MAX_FILE_SIZE = 10485760 # 10 MB
44+
DEFAULT_ELASTICSEARCH_MAX_RETRIES = 5
45+
DEFAULT_ELASTICSEARCH_RETRY_INTERVAL = 10
4346

4447
# Regular expression pattern to match a basic email format (no whitespace, valid domain)
4548
EMAIL_REGEX_PATTERN = r"^\S+@\S+\.\S+$"
@@ -297,9 +300,16 @@ def __init__(
297300
self._current_memsize = 0
298301
self.refresh_timeout = refresh_timeout
299302

303+
logger.debug(f"MemQueue created with maxsize={maxsize}, maxmemsize={maxmemsize}, "
304+
f"refresh_interval={refresh_interval}, refresh_timeout={refresh_timeout}")
305+
300306
def qmemsize(self):
301307
return self._current_memsize
302308

309+
def update_maxmemsize(self, new_maxmemsize):
310+
current_maxmemsize = self.maxmemsize
311+
self.maxmemsize = new_maxmemsize
312+
303313
def _get(self):
304314
item_size, item = self._queue.popleft() # pyright: ignore
305315
self._current_memsize -= item_size

0 commit comments

Comments
 (0)