Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion src/clusterfuzz/_internal/base/concurrency.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,9 @@ def make_pool(pool_size=POOL_SIZE, max_pool_size=None, use_threads=False):
if max_pool_size is not None:
pool_size = min(pool_size, max_pool_size)

# Don't use processes on Windows and unittests to avoid hangs.
# Use threads if explicitly requested (e.g. for I/O bound tasks like GCS ops
# to avoid multiprocessing memory overhead), or on Windows/unittests to
# avoid hangs.
if (use_threads or environment.get_value('PY_UNITTESTS') or
environment.platform() == 'WINDOWS'):
yield futures.ThreadPoolExecutor(pool_size)
Expand Down
44 changes: 31 additions & 13 deletions src/clusterfuzz/_internal/bot/tasks/utasks/corpus_pruning_task.py
Original file line number Diff line number Diff line change
Expand Up @@ -1155,6 +1155,35 @@ def _create_backup_urls(fuzz_target: data_types.FuzzTarget,
corpus_pruning_task_input.dated_backup_signed_url = dated_backup_signed_url


def _should_use_threaded_ops(fuzzer_name: str) -> bool:
"""Returns True if threaded storage ops should be used."""
threaded_ops_flag = (
feature_flags.FeatureFlags.STORAGE_THREADED_OPS_FUZZ_TARGETS)
if not threaded_ops_flag.enabled:
return False

threaded_ops_targets = threaded_ops_flag.string_value.strip()
if not threaded_ops_targets:
return False

allowed_targets = [
t.strip().lower() for t in threaded_ops_targets.split(',') if t.strip()
]

if '*' in allowed_targets:
# TODO(paulovlb): Remove this once fixed.
logs.info('[Corpus Fix] Enabled threaded storage ops for ALL targets.')
return True

if fuzzer_name.lower() in allowed_targets:
# TODO(paulovlb): Remove this once fixed.
logs.info(f'[Corpus Fix] Enabled threaded storage ops for target: '
f'{fuzzer_name}')
return True

return False


def _utask_preprocess(fuzzer_name, job_type, uworker_env):
"""Runs preprocessing for corpus pruning task."""
fuzz_target = data_handler.get_fuzz_target(fuzzer_name)
Expand Down Expand Up @@ -1217,19 +1246,8 @@ def _utask_preprocess(fuzzer_name, job_type, uworker_env):
if uworker_env is None:
uworker_env = {}

threaded_ops_flag = (
feature_flags.FeatureFlags.STORAGE_THREADED_OPS_FUZZ_TARGETS)
if threaded_ops_flag.enabled:
threaded_ops_targets = threaded_ops_flag.string_value
if threaded_ops_targets:
allowed_targets = [
t.strip() for t in threaded_ops_targets.split(',') if t.strip()
]
if fuzzer_name in allowed_targets:
uworker_env['USE_THREADED_STORAGE_OPS'] = 'True'
# TODO(paulovlb): Remove this once fixed.
logs.info(f'[Corpus Fix] Enabled threaded storage ops for target: '
f'{fuzzer_name}')
if _should_use_threaded_ops(fuzzer_name):
uworker_env['USE_THREADED_STORAGE_OPS'] = 'True'

logs.info('done preprocess')
return uworker_msg_pb2.Input( # pylint: disable=no-member
Expand Down
Loading