From 3bbec5ff0c0e016a92ca3d521a8f42d5a8d2e7ca Mon Sep 17 00:00:00 2001 From: Tomas Kopecek Date: Mar 01 2024 10:13:32 +0000 Subject: [PATCH 1/7] Basic host data for scheduler Related: https://pagure.io/koji/issue/4030 --- diff --git a/koji/daemon.py b/koji/daemon.py index 2b8e700..14f2368 100644 --- a/koji/daemon.py +++ b/koji/daemon.py @@ -1025,9 +1025,21 @@ class TaskManager(object): else: self.logger.info("Lingering task %r (pid %r)" % (id, pid)) + def _get_host_data(self): + data = { + 'methods': list(self.handlers.keys()), + 'maxjobs': self.options.maxjobs, + # TODO: now it would be duplicated by updateHost + #'ready': self.ready, + #'task_load': self.task_load, + #cpu_load, free_mem, free_disk, ... + } + return data + def getNextTask(self): self.ready = self.readyForTask() self.session.host.updateHost(self.task_load, self.ready) + self.session.host.setHostData(self._get_host_data()) if not self.ready: self.logger.info("Not ready for task") return False From ed63d144237b088e1ac0bdba01671147227caa32 Mon Sep 17 00:00:00 2001 From: Tomas Kopecek Date: Mar 01 2024 10:24:21 +0000 Subject: [PATCH 2/7] json conversion --- diff --git a/koji/daemon.py b/koji/daemon.py index 14f2368..1fab163 100644 --- a/koji/daemon.py +++ b/koji/daemon.py @@ -24,6 +24,7 @@ from __future__ import absolute_import, division import errno import hashlib +import json import logging import os import re @@ -1039,7 +1040,7 @@ class TaskManager(object): def getNextTask(self): self.ready = self.readyForTask() self.session.host.updateHost(self.task_load, self.ready) - self.session.host.setHostData(self._get_host_data()) + self.session.host.setHostData(json.dumps(self._get_host_data())) if not self.ready: self.logger.info("Not ready for task") return False From 0ca6c6464ebf860f51547974dd62022b3db438a7 Mon Sep 17 00:00:00 2001 From: Tomas Kopecek Date: Mar 04 2024 14:49:18 +0000 Subject: [PATCH 3/7] Use also builder's maxjobs setting in scheduler Related: https://pagure.io/koji/issue/4038 --- diff --git a/kojihub/scheduler.py b/kojihub/scheduler.py index b35d4d6..f3c26e0 100644 --- a/kojihub/scheduler.py +++ b/kojihub/scheduler.py @@ -298,6 +298,7 @@ class TaskScheduler(object): h_refused = refusals.get(task['task_id'], {}) for host in self.hosts_by_bin.get(task['_bin'], []): if (host['ready'] and host['_ntasks'] < self.maxjobs and + host['_ntasks'] < host['data']['maxjobs'] and host['capacity'] - host['_load'] > min_avail and host['id'] not in h_refused): task['_hosts'].append(host) @@ -324,7 +325,8 @@ class TaskScheduler(object): [(h['name'], "%(_rank).2f" % h) for h in task['_hosts']]) for host in task['_hosts']: if (host['capacity'] - host['_load'] > min_avail and - host['_ntasks'] < self.maxjobs): + host['_ntasks'] < self.maxjobs and + host['_ntasks'] < host['data']['maxjobs']): # add run entry self.assign(task, host) # update our totals and rank @@ -535,6 +537,7 @@ class TaskScheduler(object): ('host.ready', 'ready'), ('host_config.arches', 'arches'), ('host_config.capacity', 'capacity'), + ('scheduler_host_data.data', 'data'), ) fields, aliases = zip(*fields) @@ -547,7 +550,8 @@ class TaskScheduler(object): 'host_config.active IS TRUE', ], joins=[ - 'host_config ON host.id = host_config.host_id' + 'LEFT JOIN host_config ON host.id = host_config.host_id', + 'LEFT JOIN scheduler_host_data ON host.id = scheduler_host_data.host_id', ] ) From d1a1595a33eec72c766de7406de43047f2044b4c Mon Sep 17 00:00:00 2001 From: Mike McLean Date: Mar 04 2024 18:52:37 +0000 Subject: [PATCH 4/7] handle missing host data, use self.maxjobs as default --- diff --git a/kojihub/scheduler.py b/kojihub/scheduler.py index f3c26e0..12ae989 100644 --- a/kojihub/scheduler.py +++ b/kojihub/scheduler.py @@ -191,7 +191,6 @@ class TaskScheduler(object): self.active_tasks = [] self.free_tasks = [] - # TODO these things need proper config self.maxjobs = context.opts['MaxJobs'] self.capacity_overcommit = context.opts['CapacityOvercommit'] self.ready_timeout = context.opts['ReadyTimeout'] @@ -282,6 +281,11 @@ class TaskScheduler(object): host.setdefault('_load', 0.0) host.setdefault('_ntasks', 0) host.setdefault('_demand', 0.0) + # host data might be unset + hostdata = host['data'] + if hostdata is None: + hostdata = {} + host.setdefault('_maxjobs', hostdata.get('maxjobs') or self.maxjobs) # temporary test code logger.info(f'Host: {host}') ldiff = host['task_load'] - host['_load'] @@ -297,8 +301,8 @@ class TaskScheduler(object): min_avail = min(0, task['weight'] - self.capacity_overcommit) h_refused = refusals.get(task['task_id'], {}) for host in self.hosts_by_bin.get(task['_bin'], []): - if (host['ready'] and host['_ntasks'] < self.maxjobs and - host['_ntasks'] < host['data']['maxjobs'] and + if (host['ready'] and + host['_ntasks'] < host['_maxjobs'] and host['capacity'] - host['_load'] > min_avail and host['id'] not in h_refused): task['_hosts'].append(host) @@ -325,8 +329,7 @@ class TaskScheduler(object): [(h['name'], "%(_rank).2f" % h) for h in task['_hosts']]) for host in task['_hosts']: if (host['capacity'] - host['_load'] > min_avail and - host['_ntasks'] < self.maxjobs and - host['_ntasks'] < host['data']['maxjobs']): + host['_ntasks'] < host['_maxjobs']): # add run entry self.assign(task, host) # update our totals and rank From 93fc5525a0d6002fc0900889039705b1af58ad79 Mon Sep 17 00:00:00 2001 From: Mike McLean Date: Mar 04 2024 18:53:17 +0000 Subject: [PATCH 5/7] unit test --- diff --git a/tests/test_hub/test_scheduler.py b/tests/test_hub/test_scheduler.py index 3ad6464..984727f 100644 --- a/tests/test_hub/test_scheduler.py +++ b/tests/test_hub/test_scheduler.py @@ -120,6 +120,133 @@ class TestScheduler(BaseTest): # TODO +class TestDoSchedule(BaseTest): + + def setUp(self): + super(TestDoSchedule, self).setUp() + self.sched = scheduler.TaskScheduler() + self.sched.get_refusals = mock.MagicMock() + self.sched._get_hosts = mock.MagicMock() + self.assigns = [] + self.sched.assign = mock.MagicMock(side_effect=self.my_assign) + + def my_assign(self, task, host): + self.assigns.append((task,host)) + + def test_no_hosts_no_tasks(self): + self.sched._get_hosts.return_value = [] + self.sched.get_hosts() + self.sched.do_schedule() + + self.sched.assign.assert_not_called() + self.assertEqual(len(self.queries), 0) + self.assertEqual(len(self.inserts), 0) + self.assertEqual(len(self.updates), 0) + + def mktask(self, **kw): + data = kw.copy() + data.setdefault('host_id', None) + data.setdefault('waiting', False) + data.setdefault('weight', 1.0) + data.setdefault('channel_id', 1) + data.setdefault('arch', 'noarch') + data.setdefault('_bin', '%(channel_id)s:%(arch)s' % data) + data.setdefault('task_id', mock.MagicMock()) # ?? + return data + + def mkhost(self, **kw): + data = kw.copy() + data.setdefault('task_load', 0.0) + data.setdefault('ready', True) + data.setdefault('data', None) + data.setdefault('capacity', 10.0) + data.setdefault('channels', [1]) + data.setdefault('arches', 'x86_64') + data.setdefault('id', mock.MagicMock()) # ?? + data.setdefault('name', f'Host {data["id"]}') # ?? + return data + + def test_no_hosts_free_tasks(self): + self.sched._get_hosts.return_value = [] + self.sched.get_hosts() + self.sched.free_tasks = [self.mktask(task_id=n) for n in range(5)] + + self.sched.do_schedule() + + self.sched.assign.assert_not_called() + self.assertEqual(len(self.queries), 0) + self.assertEqual(len(self.inserts), 0) + self.assertEqual(len(self.updates), 0) + + def test_no_tasks_avail_hosts(self): + self.sched._get_hosts.return_value = [self.mkhost(id=n) for n in range(5)] + self.sched.get_hosts() + self.sched.free_tasks = [] + + self.sched.do_schedule() + + self.sched.assign.assert_not_called() + self.assertEqual(len(self.queries), 0) + self.assertEqual(len(self.inserts), 0) + self.assertEqual(len(self.updates), 0) + + def test_no_tasks_avail_hosts(self): + self.sched._get_hosts.return_value = [self.mkhost(id=n) for n in range(5)] + self.sched.get_hosts() + self.sched.free_tasks = [] + + self.sched.do_schedule() + + self.sched.assign.assert_not_called() + self.assertEqual(len(self.queries), 0) + self.assertEqual(len(self.inserts), 0) + self.assertEqual(len(self.updates), 0) + + def test_free_tasks_avail_hosts(self): + self.sched._get_hosts.return_value = [self.mkhost(id=n) for n in range(5)] + self.sched.get_hosts() + self.sched.free_tasks = [self.mktask(task_id=n) for n in range(5)] + + self.sched.do_schedule() + + # in this simple case, tasks should be evenly assigned + self.assertEqual(len(self.assigns), 5) + t_assigned = [] + h_used = [] + for task, host in self.assigns: + t_assigned.append(task['task_id']) + h_used.append(host['id']) + t_assigned.sort() + h_used.sort() + self.assertEqual(t_assigned, list(range(5))) + self.assertEqual(h_used, list(range(5))) + + def test_active_tasks(self): + self.context.opts['CapacityOvercommit'] = 1.0 + hosts = [self.mkhost(id=n, capacity=2.0) for n in range(5)] + active = [self.mktask(task_id=n, host_id=n, weight=4.0) for n in range(3)] + # so, first three hosts have a task of weight=4, more than capacity+overcommit + free = [self.mktask(task_id=n, weight=2.0) for n in range(3,5)] + # two free tasks with weight=2 + self.sched._get_hosts.return_value = hosts + self.sched.get_hosts() + self.sched.free_tasks = free + self.sched.active_tasks = active + + self.sched.do_schedule() + + # we expect that the two free tasks will be assigned evenly to the hosts with free capacity + self.assertEqual(len(self.assigns), 2) + t_assigned = [] + h_used = [] + for task, host in self.assigns: + t_assigned.append(task['task_id']) + h_used.append(host['id']) + t_assigned.sort() + h_used.sort() + self.assertEqual(t_assigned, list(range(3,5))) + self.assertEqual(h_used, list(range(3,5))) + class TestCheckActiveRuns(BaseTest): def setUp(self): From 6840c535c60eb7d4e63e414e080f0a07b227369c Mon Sep 17 00:00:00 2001 From: Mike McLean Date: Mar 04 2024 20:38:43 +0000 Subject: [PATCH 6/7] streamline hub calls --- diff --git a/koji/daemon.py b/koji/daemon.py index 1fab163..2fe583c 100644 --- a/koji/daemon.py +++ b/koji/daemon.py @@ -1039,8 +1039,7 @@ class TaskManager(object): def getNextTask(self): self.ready = self.readyForTask() - self.session.host.updateHost(self.task_load, self.ready) - self.session.host.setHostData(json.dumps(self._get_host_data())) + self.session.host.updateHost(self.task_load, self.ready, data=self._get_host_data()) if not self.ready: self.logger.info("Not ready for task") return False diff --git a/kojihub/kojihub.py b/kojihub/kojihub.py index 1c75a09..3f30d07 100644 --- a/kojihub/kojihub.py +++ b/kojihub/kojihub.py @@ -14703,10 +14703,19 @@ class HostExports(object): host.verify() return host.id - def updateHost(self, task_load, ready): + def updateHost(self, task_load, ready, data=None): + """Update host data + + :param float task_load: current task load + :param bool ready: whether the host is ready to take a task + :param dict data: data for the scheduler + + """ host = Host() host.verify() host.updateHost(task_load, ready) + if data is not None: + scheduler.set_host_data(host.id, data) def getLoadData(self): host = Host() @@ -14760,21 +14769,21 @@ class HostExports(object): return task.setWeight(weight) def setHostData(self, hostdata): - """Builder will update all its resources + """Provide host data for the scheduler - Initial implementation contains: - - available task methods - - maxjobs - - host readiness + :param dict hostdata: host data + + For backwards compatibility, we also accept hostdata as a string containing a + json-encoded dictionary. """ host = Host() host.verify() - upsert = UpsertProcessor( - table='scheduler_host_data', - keys=['host_id'], - data={'host_id': host.id, 'data': hostdata}, - ) - upsert.execute() + if isinstance(hostdata, str): + # for backwards compatibility + data = json.loads(hostdata) + else: + data = hostdata + scheduler.set_host_data(host.id, hostdata) def getTasks(self): host = Host() diff --git a/kojihub/scheduler.py b/kojihub/scheduler.py index 12ae989..511fcd8 100644 --- a/kojihub/scheduler.py +++ b/kojihub/scheduler.py @@ -153,6 +153,17 @@ def get_host_data(hostID=None): return query.execute() +def set_host_data(hostID, data): + if not isinstance(data, dict): + raise koji.ParameterError('Host data should be a dictionary') + upsert = UpsertProcessor( + table='scheduler_host_data', + keys=['host_id'], + data={'host_id': hostID, 'data': json.dumps(data)}, + ) + upsert.execute() + + class TaskRunsQuery(QueryView): tables = ['scheduler_task_runs'] From 7b0ad45d3f2d4297f0b8f4cdceee2baf1508ea66 Mon Sep 17 00:00:00 2001 From: Mike McLean Date: Mar 04 2024 20:43:30 +0000 Subject: [PATCH 7/7] flake8, fix typo --- diff --git a/koji/daemon.py b/koji/daemon.py index 2fe583c..fe2d292 100644 --- a/koji/daemon.py +++ b/koji/daemon.py @@ -24,7 +24,6 @@ from __future__ import absolute_import, division import errno import hashlib -import json import logging import os import re @@ -1031,9 +1030,9 @@ class TaskManager(object): 'methods': list(self.handlers.keys()), 'maxjobs': self.options.maxjobs, # TODO: now it would be duplicated by updateHost - #'ready': self.ready, - #'task_load': self.task_load, - #cpu_load, free_mem, free_disk, ... + # 'ready': self.ready, + # 'task_load': self.task_load, + # cpu_load, free_mem, free_disk, ... } return data diff --git a/kojihub/kojihub.py b/kojihub/kojihub.py index 3f30d07..d64a702 100644 --- a/kojihub/kojihub.py +++ b/kojihub/kojihub.py @@ -14780,9 +14780,7 @@ class HostExports(object): host.verify() if isinstance(hostdata, str): # for backwards compatibility - data = json.loads(hostdata) - else: - data = hostdata + hostdata = json.loads(hostdata) scheduler.set_host_data(host.id, hostdata) def getTasks(self):