From 8b5efc075f81097fcb943a82d4590aa77b095a3f Mon Sep 17 00:00:00 2001 From: Paulo Borges Date: Tue, 11 Aug 2026 14:54:12 +0000 Subject: [PATCH] [Corpus] Refactor and add wildcard support for ff --- src/clusterfuzz/_internal/base/concurrency.py | 4 +- .../bot/tasks/utasks/corpus_pruning_task.py | 44 +++++++++++++------ 2 files changed, 34 insertions(+), 14 deletions(-) diff --git a/src/clusterfuzz/_internal/base/concurrency.py b/src/clusterfuzz/_internal/base/concurrency.py index 3edb6463bc8..9934011153f 100644 --- a/src/clusterfuzz/_internal/base/concurrency.py +++ b/src/clusterfuzz/_internal/base/concurrency.py @@ -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) diff --git a/src/clusterfuzz/_internal/bot/tasks/utasks/corpus_pruning_task.py b/src/clusterfuzz/_internal/bot/tasks/utasks/corpus_pruning_task.py index bd803079963..50d872162bb 100644 --- a/src/clusterfuzz/_internal/bot/tasks/utasks/corpus_pruning_task.py +++ b/src/clusterfuzz/_internal/bot/tasks/utasks/corpus_pruning_task.py @@ -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) @@ -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