diff --git a/robusta_krr/core/abstract/metrics.py b/robusta_krr/core/abstract/metrics.py index 41881719..29327eff 100644 --- a/robusta_krr/core/abstract/metrics.py +++ b/robusta_krr/core/abstract/metrics.py @@ -16,5 +16,5 @@ class BaseMetric(ABC): @abstractmethod async def load_data( - self, object: K8sObjectData, period: datetime.timedelta, step: datetime.timedelta + self, object: K8sObjectData, period: datetime.timedelta, step: datetime.timedelta, end_time: datetime.datetime ) -> PodsTimeData: ... diff --git a/robusta_krr/core/abstract/strategies.py b/robusta_krr/core/abstract/strategies.py index 2ff675b1..0ddfac62 100644 --- a/robusta_krr/core/abstract/strategies.py +++ b/robusta_krr/core/abstract/strategies.py @@ -1,5 +1,6 @@ from __future__ import annotations +import time import abc import datetime from textwrap import dedent @@ -48,6 +49,7 @@ class StrategySettings(pd.BaseModel): 24 * 7 * 2, ge=1, description="The duration of the history data to use (in hours)." ) timeframe_duration: float = pd.Field(1.25, gt=0, description="The step for the history data (in minutes).") + end_epoch: float = pd.Field(time.time(), ge=0.0, description="The Unix epoch timestamp marking the end of the query window.") @property def history_timedelta(self) -> datetime.timedelta: @@ -57,6 +59,10 @@ def history_timedelta(self) -> datetime.timedelta: def timeframe_timedelta(self) -> datetime.timedelta: return datetime.timedelta(minutes=self.timeframe_duration) + @property + def end_datetime(self) -> datetime.datetime: + return datetime.datetime.fromtimestamp(self.end_epoch, tz=datetime.timezone.utc) + def history_range_enough(self, history_range: tuple[datetime.timedelta, datetime.timedelta]) -> bool: """Override this function to check if the history range is enough for the strategy.""" diff --git a/robusta_krr/core/integrations/prometheus/loader.py b/robusta_krr/core/integrations/prometheus/loader.py index 87c39829..d84bb486 100644 --- a/robusta_krr/core/integrations/prometheus/loader.py +++ b/robusta_krr/core/integrations/prometheus/loader.py @@ -107,6 +107,7 @@ async def gather_data( object: K8sObjectData, strategy: BaseStrategy, period: datetime.timedelta, + end_time: datetime.datetime, *, step: datetime.timedelta = datetime.timedelta(minutes=30), ) -> MetricsPodData: @@ -117,6 +118,7 @@ async def gather_data( object (K8sObjectData): The Kubernetes object. resource (ResourceType): The resource type. period (datetime.timedelta): The time period for which to gather data. + end_time (datetime.datetime): The timestamp marking the end of the query window. step (datetime.timedelta, optional): The time step between data points. Defaults to 30 minutes. Returns: @@ -124,6 +126,6 @@ async def gather_data( """ return { - MetricLoader.__name__: await self.loader.gather_data(object, MetricLoader, period, step) + MetricLoader.__name__: await self.loader.gather_data(object, MetricLoader, period, end_time, step) for MetricLoader in strategy.metrics } diff --git a/robusta_krr/core/integrations/prometheus/metrics/base.py b/robusta_krr/core/integrations/prometheus/metrics/base.py index 347e6b93..c29752ff 100644 --- a/robusta_krr/core/integrations/prometheus/metrics/base.py +++ b/robusta_krr/core/integrations/prometheus/metrics/base.py @@ -154,7 +154,7 @@ async def query_prometheus(self, data: PrometheusMetricData) -> list[PrometheusS return await loop.run_in_executor(self.executor, lambda: self._query_prometheus_sync(data)) async def load_data( - self, object: K8sObjectData, period: datetime.timedelta, step: datetime.timedelta + self, object: K8sObjectData, period: datetime.timedelta, step: datetime.timedelta, end_time: datetime.datetime ) -> PodsTimeData: """ Asynchronous method that loads metric data for a specific object. @@ -163,6 +163,7 @@ async def load_data( object (K8sObjectData): The object for which metrics need to be loaded. period (datetime.timedelta): The time period for which metrics need to be loaded. step (datetime.timedelta): The time interval between successive metric values. + end_time (datetime.datetime): The timestamp marking the end of the query window. Returns: ResourceHistoryData: An instance of the ResourceHistoryData class representing the loaded metrics. @@ -172,14 +173,13 @@ async def load_data( duration_str = self._step_to_string(period) query = self.get_query(object, duration_str, step_str) - end_time = datetime.datetime.utcnow().replace(second=0, microsecond=0) start_time = end_time - period # Here if we split the object into multiple sub-objects, we query each sub-object recursively. if self.pods_batch_size is not None and object.pods_count > self.pods_batch_size: results = await asyncio.gather( *[ - self.load_data(splitted_object, period, step) + self.load_data(splitted_object, period, step, end_time) for splitted_object in object.split_into_batches(self.pods_batch_size) ] ) diff --git a/robusta_krr/core/integrations/prometheus/metrics_service/base_metric_service.py b/robusta_krr/core/integrations/prometheus/metrics_service/base_metric_service.py index 6e352388..28fe0f80 100644 --- a/robusta_krr/core/integrations/prometheus/metrics_service/base_metric_service.py +++ b/robusta_krr/core/integrations/prometheus/metrics_service/base_metric_service.py @@ -43,6 +43,7 @@ async def gather_data( object: K8sObjectData, LoaderClass: type[PrometheusMetric], period: datetime.timedelta, + end_time: datetime.datetime, step: datetime.timedelta = datetime.timedelta(minutes=30), ) -> PodsTimeData: ... diff --git a/robusta_krr/core/integrations/prometheus/metrics_service/prometheus_metrics_service.py b/robusta_krr/core/integrations/prometheus/metrics_service/prometheus_metrics_service.py index bc19dcab..3e72fbf3 100644 --- a/robusta_krr/core/integrations/prometheus/metrics_service/prometheus_metrics_service.py +++ b/robusta_krr/core/integrations/prometheus/metrics_service/prometheus_metrics_service.py @@ -204,6 +204,7 @@ async def gather_data( object: K8sObjectData, LoaderClass: type[PrometheusMetric], period: timedelta, + end_time: datetime, step: timedelta = timedelta(minutes=30), ) -> PodsTimeData: """ @@ -212,7 +213,7 @@ async def gather_data( logger.debug(f"Gathering {LoaderClass.__name__} metric for {object}") try: metric_loader = LoaderClass(self.get_prometheus(), self.name(), self.executor) - data = await metric_loader.load_data(object, period, step) + data = await metric_loader.load_data(object, period, step, end_time) except Exception: logger.exception("Failed to gather resource history data for %s", object) data = {} diff --git a/robusta_krr/core/runner.py b/robusta_krr/core/runner.py index 23a63fd4..bd895765 100644 --- a/robusta_krr/core/runner.py +++ b/robusta_krr/core/runner.py @@ -348,6 +348,7 @@ async def _calculate_object_recommendations(self, object: K8sObjectData) -> Opti object, self._strategy, self._strategy.settings.history_timedelta, + self._strategy.settings.end_datetime, step=self._strategy.settings.timeframe_timedelta, )