Skip to content
Draft
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
1 change: 1 addition & 0 deletions src/clusterfuzz/_internal/base/feature_flags.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ class FeatureFlags(Enum):
GCP_BATCH_JOBS_FREQUENCY = 'gcp_batch_jobs_frequency'

PREPROCESS_QUEUE_SIZE_LIMIT = 'preprocess_queue_size_limit'
CHROME_FUZZ_TARGET_CPUS = 'chrome_fuzz_target_cpus'

SWARMING_REMOTE_EXECUTION = 'swarming_remote_execution'
# TODO(b/516630567): Set this value based off dev & stage metrics and tests.
Expand Down
128 changes: 126 additions & 2 deletions src/clusterfuzz/_internal/cron/schedule_fuzz.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,24 +16,30 @@
from abc import ABC
from abc import abstractmethod
import collections
import json
import random
import time
import urllib.request

import google.auth.transport.requests
from google.cloud import monitoring_v3

from clusterfuzz._internal import swarming
from clusterfuzz._internal.base import memoize
from clusterfuzz._internal.base import tasks
from clusterfuzz._internal.base import utils
from clusterfuzz._internal.base.feature_flags import FeatureFlags
from clusterfuzz._internal.base.tasks import pub_sub_task_queue
from clusterfuzz._internal.base.tasks.pub_sub_task_queue import \
SWARMING_PREPROCESS_QUEUE
from clusterfuzz._internal.base.tasks.pub_sub_task_queue import PREPROCESS_QUEUE
from clusterfuzz._internal.base.tasks.pub_sub_task_queue import PubSubTaskQueue
from clusterfuzz._internal.config import local_config
from clusterfuzz._internal.datastore import data_types
from clusterfuzz._internal.datastore import ndb_utils
from clusterfuzz._internal.google_cloud_utils import credentials
from clusterfuzz._internal.metrics import logs
from clusterfuzz._internal.system import environment


@memoize.wrap(memoize.InMemory(60))
Expand Down Expand Up @@ -228,8 +234,75 @@ def _get_swarming_jobs():
]


def _remaining_queue_capacity(queue: PubSubTaskQueue) -> int:
"""Returns the remaining capacity of the given queue."""
CPUS_PER_FUZZ_TASK = 2
DEFAULT_BATCH_REGIONS = ['us-central1', 'us-west1']


def _get_batch_project_id() -> str:
"""Returns the GCP project ID where Batch jobs are executed."""
batch_project = environment.get_value('BATCH_CLOUD_PROJECT_ID')
if batch_project:
return batch_project
k8s_project = environment.get_value('K8S_PROJECT')
if k8s_project:
return k8s_project
return 'google.com:clusterfuzz'


def _get_batch_regions() -> list[str]:
"""Returns the configured regions for GCP Batch."""
try:
batch_config = local_config.BatchConfig()
subconfigs = batch_config.get('subconfigs', {})
regions = {
subconfig.get('region')
for subconfig in subconfigs.values()
if subconfig.get('region')
}
if regions:
return sorted(list(regions))
except Exception:
pass
return DEFAULT_BATCH_REGIONS


@memoize.wrap(memoize.InMemory(60))
def get_batch_active_jobs_count(project: str, region: str) -> int:
"""Queries active batch job counts for a region using the countByState API."""
creds = credentials.get_default()[0]
if not creds.valid:
creds.refresh(google.auth.transport.requests.Request())

headers = {
'Authorization': f'Bearer {creds.token}',
'Content-Type': 'application/json',
}
url = (f'https://batch.googleapis.com/v1alpha/projects/{project}/locations/'
f'{region}/jobs:countByState?states=RUNNING&states=SCHEDULED&states=QUEUED')
req = urllib.request.Request(url, headers=headers)
with urllib.request.urlopen(req, timeout=5) as response:
if response.status != 200:
logs.error(f'Batch countByState returned status {response.status} for {region}.')
return 0
data = json.loads(response.read())
job_counts = data.get('jobCounts', {})
return sum(int(count) for count in job_counts.values())


def get_total_batch_active_jobs(project: str, regions: list[str]) -> int | None:
"""Aggregates active batch jobs across all configured regions."""
try:
total = 0
for region in regions:
total += get_batch_active_jobs_count(project, region)
return total
except Exception as error:
logs.error(f'Failed to retrieve batch active jobs: {error}')
return None


def _remaining_queue_capacity_by_queue_size(queue: PubSubTaskQueue) -> int:
"""Returns the remaining capacity based on queue unacked message count."""
project = utils.get_application_id()
creds = credentials.get_default()[0]
preprocess_queue_size = get_queue_size(creds, project, queue.name)
Expand All @@ -243,6 +316,57 @@ def _remaining_queue_capacity(queue: PubSubTaskQueue) -> int:
return num_tasks


def _remaining_chrome_cpu_capacity(queue: PubSubTaskQueue) -> int:
"""Calculates remaining task capacity based on active CPU usage and target CPU limit."""
flag = FeatureFlags.CHROME_FUZZ_TARGET_CPUS
if not flag.content:
logs.warning('CHROME_FUZZ_TARGET_CPUS flag is enabled but empty. Using queue size.')
return _remaining_queue_capacity_by_queue_size(queue)

try:
target_cpus = int(flag.content)
except (ValueError, TypeError):
logs.error(f'Invalid CHROME_FUZZ_TARGET_CPUS value: {flag.content}. Using queue size.')
return _remaining_queue_capacity_by_queue_size(queue)

project = _get_batch_project_id()
regions = _get_batch_regions()
active_jobs = get_total_batch_active_jobs(project, regions)
if active_jobs is None:
logs.warning('Failed to query batch jobs. Falling back to queue size.')
return _remaining_queue_capacity_by_queue_size(queue)

app_id = utils.get_application_id()
creds = credentials.get_default()[0]
preprocess_size = get_queue_size(creds, app_id, queue.name)
utask_main_size = tasks.get_utask_main_queue_size(pub_sub_task_queue.UTASK_MAIN_QUEUE.name)

active_batch_cpus = active_jobs * CPUS_PER_FUZZ_TASK
in_flight_pipeline_cpus = (preprocess_size + utask_main_size) * CPUS_PER_FUZZ_TASK
total_occupied_cpus = active_batch_cpus + in_flight_pipeline_cpus

logs.info(
f'Chrome CPU Target: {target_cpus}. Active Batch CPUs: {active_batch_cpus} ({active_jobs} jobs). '
f'In-Flight Pipeline CPUs: {in_flight_pipeline_cpus}. Total Occupied: {total_occupied_cpus}.'
)

if total_occupied_cpus >= target_cpus:
logs.info('Target CPU capacity reached or exceeded. Not scheduling tasks.')
return 0

deficit_cpus = target_cpus - total_occupied_cpus
needed_tasks = deficit_cpus // CPUS_PER_FUZZ_TASK
return min(needed_tasks, queue.get_max_target_size())


def _remaining_queue_capacity(queue: PubSubTaskQueue) -> int:
"""Returns the remaining capacity of the given queue."""
if queue == PREPROCESS_QUEUE and FeatureFlags.CHROME_FUZZ_TARGET_CPUS.enabled:
return _remaining_chrome_cpu_capacity(queue)

return _remaining_queue_capacity_by_queue_size(queue)


def _fill_queue(queue: PubSubTaskQueue, provider: BaseFuzzTaskProvider):
"""Fills the given queue with tasks from the provider."""
start = time.time()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

from clusterfuzz._internal.cron import schedule_fuzz
from clusterfuzz._internal.datastore import data_types
from clusterfuzz._internal.tests.test_libs import helpers as test_helpers
from clusterfuzz._internal.tests.test_libs import test_utils

# pylint: disable=protected-access
Expand Down Expand Up @@ -233,3 +234,55 @@ def test_os_version_job_without_version(self):
self._setup_chrome_entities()
task = self._run_and_get_task()
self.assertIsNone(task.extra_info.get('base_os_version'))


@test_utils.with_cloud_emulators('datastore')
class ChromeCpuCapacityTest(unittest.TestCase):
"""Tests for CPU-aware scheduling capacity in schedule_fuzz."""

def setUp(self):
test_helpers.patch(self, [
'clusterfuzz._internal.cron.schedule_fuzz.get_total_batch_active_jobs',
'clusterfuzz._internal.cron.schedule_fuzz.get_queue_size',
'clusterfuzz._internal.base.tasks.get_utask_main_queue_size',
])

def test_cpu_capacity_schedules_deficit(self):
"""Tests that deficit tasks are scheduled when below CPU target."""
data_types.FeatureFlag(
id='chrome_fuzz_target_cpus', enabled=True, value='1000').put()
self.mock.get_total_batch_active_jobs.return_value = 300
self.mock.get_queue_size.return_value = 50
self.mock.get_utask_main_queue_size.return_value = 50

capacity = schedule_fuzz._remaining_chrome_cpu_capacity(
schedule_fuzz.PREPROCESS_QUEUE)
self.assertEqual(capacity, 100)

def test_cpu_capacity_zero_when_at_or_above_target(self):
"""Tests that no tasks are scheduled when CPU target is met or exceeded."""
data_types.FeatureFlag(
id='chrome_fuzz_target_cpus', enabled=True, value='1000').put()
self.mock.get_total_batch_active_jobs.return_value = 500
self.mock.get_queue_size.return_value = 0
self.mock.get_utask_main_queue_size.return_value = 0

capacity = schedule_fuzz._remaining_chrome_cpu_capacity(
schedule_fuzz.PREPROCESS_QUEUE)
self.assertEqual(capacity, 0)

def test_fallback_when_batch_query_fails(self):
"""Tests that queue size fallback is used if batch query fails."""
data_types.FeatureFlag(
id='chrome_fuzz_target_cpus', enabled=True, value='1000').put()
self.mock.get_total_batch_active_jobs.return_value = None
self.mock.get_queue_size.return_value = 200

test_helpers.patch(self, [
'clusterfuzz._internal.base.tasks.pub_sub_task_queue.PubSubTaskQueue.get_max_target_size'
])
self.mock.get_max_target_size.return_value = 1000

capacity = schedule_fuzz._remaining_chrome_cpu_capacity(
schedule_fuzz.PREPROCESS_QUEUE)
self.assertEqual(capacity, 800)
Loading