made all triggers async
This commit is contained in:
@@ -13,7 +13,6 @@ from cpu_governor_auto_adjust.logger import getLogger
|
||||
class TriggerTuple(NamedTuple):
|
||||
name: str
|
||||
loglevel: str
|
||||
type: str
|
||||
interval_in_seconds: int
|
||||
governor: str
|
||||
custom_config: OrderedDict[str, Any]
|
||||
@@ -84,7 +83,6 @@ class Config:
|
||||
for trigger_elem in _trigger_elements:
|
||||
name = self._get_single_text_value_from_xpath('name', trigger_elem)
|
||||
loglevel = self._get_single_text_value_from_xpath('logLevel', trigger_elem)
|
||||
type = self._get_single_text_value_from_xpath('type', trigger_elem)
|
||||
interval_in_seconds = self._get_single_text_value_from_xpath('intervalInSeconds', trigger_elem)
|
||||
governor = self._get_single_text_value_from_xpath('governor', trigger_elem)
|
||||
try:
|
||||
@@ -96,7 +94,6 @@ class Config:
|
||||
_new_trigger = TriggerTuple(
|
||||
name=name,
|
||||
loglevel=loglevel,
|
||||
type=type,
|
||||
interval_in_seconds=int(interval_in_seconds),
|
||||
governor=governor,
|
||||
custom_config=custom_config
|
||||
|
||||
@@ -49,23 +49,9 @@ class TriggerScheduler(AppClass):
|
||||
self.log.debug("callback trigger %s, state: %s", _trigger.name, _trigger.trigger_state)
|
||||
await asyncio.sleep(_trigger.config.interval_in_seconds)
|
||||
|
||||
async def sync_trigger(self, _trigger: Trigger) -> None:
|
||||
"""Run a trigger with a specific name at a given interval."""
|
||||
while True:
|
||||
start = monotonic()
|
||||
_trigger.run()
|
||||
end = monotonic()
|
||||
self.log.debug("run once trigger %s, state: %s, duration: %.3f ms", _trigger.name, _trigger.trigger_state, (end - start) * 1000)
|
||||
await asyncio.sleep(_trigger.config.interval_in_seconds)
|
||||
|
||||
async def start_trigger(self, _trigger: Trigger) -> None:
|
||||
"""Start a new trigger."""
|
||||
if _trigger.config.type == "async":
|
||||
self.loop.create_task(self.async_trigger(_trigger))
|
||||
elif _trigger.config.type == "sync":
|
||||
self.loop.create_task(self.sync_trigger(_trigger))
|
||||
else:
|
||||
raise ValueError(f"unknown trigger type: {_trigger.config.type}")
|
||||
self.loop.create_task(self.async_trigger(_trigger))
|
||||
self.running_triggers.append(_trigger)
|
||||
|
||||
async def stop_triggers(self) -> None:
|
||||
|
||||
@@ -2,6 +2,7 @@ from cpu_governor_auto_adjust.trigger import Trigger
|
||||
from cpu_governor_auto_adjust.config import Config
|
||||
from functools import cached_property
|
||||
from os import getloadavg
|
||||
import asyncio
|
||||
|
||||
|
||||
class CpuLoadTrigger(Trigger):
|
||||
@@ -54,29 +55,34 @@ class CpuLoadTrigger(Trigger):
|
||||
def current_load_average_over_under_threshold(self) -> bool:
|
||||
return any(load <= threshold for load, threshold in zip(self.current_load, self.low_threshold))
|
||||
|
||||
def run(self) -> None:
|
||||
current_active = self.active
|
||||
new_high_load_active = self.current_load_average_over_high_threshold
|
||||
new_low_load_active = self.current_load_average_over_under_threshold
|
||||
|
||||
if not current_active and new_high_load_active:
|
||||
self.log.info(
|
||||
"activating trigger, load: %s, high threshold: %s, governor: %s",
|
||||
self.current_load, self.high_threshold, self.governor.name
|
||||
)
|
||||
self.active = True
|
||||
async def check_load(self) -> None:
|
||||
while True:
|
||||
current_active = self.active
|
||||
new_high_load_active = self.current_load_average_over_high_threshold
|
||||
new_low_load_active = self.current_load_average_over_under_threshold
|
||||
|
||||
elif current_active and new_high_load_active:
|
||||
self.log.debug(
|
||||
"trigger is already active, load: %s, high threshold: %s",
|
||||
self.current_load, self.high_threshold
|
||||
)
|
||||
|
||||
elif current_active and new_low_load_active:
|
||||
self.log.info(
|
||||
"deactivating trigger, load: %s, low threshold: %s, governor: %s",
|
||||
self.current_load, self.low_threshold, self.governor.name
|
||||
)
|
||||
self.active = False
|
||||
else:
|
||||
self.log.debug("trigger is already inactive, load: %s, threshold: %s", self.current_load, self.high_threshold)
|
||||
if not current_active and new_high_load_active:
|
||||
self.log.info(
|
||||
"activating trigger, load: %s, high threshold: %s, governor: %s",
|
||||
self.current_load, self.high_threshold, self.governor.name
|
||||
)
|
||||
self.active = True
|
||||
|
||||
elif current_active and new_high_load_active:
|
||||
self.log.debug(
|
||||
"trigger is already active, load: %s, high threshold: %s",
|
||||
self.current_load, self.high_threshold
|
||||
)
|
||||
|
||||
elif current_active and new_low_load_active:
|
||||
self.log.info(
|
||||
"deactivating trigger, load: %s, low threshold: %s, governor: %s",
|
||||
self.current_load, self.low_threshold, self.governor.name
|
||||
)
|
||||
self.active = False
|
||||
else:
|
||||
self.log.debug("trigger is already inactive, load: %s, threshold: %s", self.current_load, self.high_threshold)
|
||||
await asyncio.sleep(1)
|
||||
|
||||
async def async_run(self) -> None:
|
||||
await self.check_load()
|
||||
|
||||
@@ -1,12 +1,15 @@
|
||||
from cpu_governor_auto_adjust.trigger import Trigger
|
||||
from cpu_governor_auto_adjust.config import Config
|
||||
import random
|
||||
import asyncio
|
||||
|
||||
class TestTrigger1(Trigger):
|
||||
def __init__(self, _config: Config) -> None:
|
||||
super().__init__(_config)
|
||||
|
||||
def run(self) -> None:
|
||||
choices = [False, True]
|
||||
self.log.debug("run check code of %s", self.__class__.__name__)
|
||||
self.active = random.choice(choices)
|
||||
async def async_run(self) -> None:
|
||||
while True:
|
||||
choices = [False, True]
|
||||
self.log.debug("run check code of %s", self.__class__.__name__)
|
||||
self.active = random.choice(choices)
|
||||
await asyncio.sleep(1)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
from cpu_governor_auto_adjust.trigger import Trigger
|
||||
from cpu_governor_auto_adjust.config import Config
|
||||
from datetime import datetime, time
|
||||
import asyncio
|
||||
|
||||
|
||||
class TimeTrigger(Trigger):
|
||||
@@ -14,11 +15,16 @@ class TimeTrigger(Trigger):
|
||||
def end_time(self) -> time:
|
||||
return datetime.strptime(self.config.custom_config['endTime'], '%H:%M').time()
|
||||
|
||||
def run(self) -> None:
|
||||
active_current = self.active
|
||||
active_new = self.start_time() <= datetime.now().time() <= self.end_time()
|
||||
if not active_current and active_new:
|
||||
self.log.info("activating trigger, start time: %s, end time: %s, governor: %s", self.start_time(), self.end_time(), self.governor.name)
|
||||
elif active_current and not active_new:
|
||||
self.log.info("deactivating trigger, start time: %s, end time: %s, governor: %s", self.start_time(), self.end_time(), self.governor.name)
|
||||
self.active = active_new
|
||||
async def check_time(self) -> None:
|
||||
while True:
|
||||
active_current = self.active
|
||||
active_new = self.start_time() <= datetime.now().time() <= self.end_time()
|
||||
if not active_current and active_new:
|
||||
self.log.info("activating trigger, start time: %s, end time: %s, governor: %s", self.start_time(), self.end_time(), self.governor.name)
|
||||
elif active_current and not active_new:
|
||||
self.log.info("deactivating trigger, start time: %s, end time: %s, governor: %s", self.start_time(), self.end_time(), self.governor.name)
|
||||
self.active = active_new
|
||||
await asyncio.sleep(1)
|
||||
|
||||
async def async_run(self) -> None:
|
||||
await self.check_time()
|
||||
|
||||
Reference in New Issue
Block a user