Skip to content
Open
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
2 changes: 1 addition & 1 deletion robusta_krr/core/abstract/metrics.py
Original file line number Diff line number Diff line change
Expand Up @@ -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: ...
6 changes: 6 additions & 0 deletions robusta_krr/core/abstract/strategies.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
from __future__ import annotations

import time
import abc
import datetime
from textwrap import dedent
Expand Down Expand Up @@ -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:
Expand All @@ -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."""

Expand Down
4 changes: 3 additions & 1 deletion robusta_krr/core/integrations/prometheus/loader.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -117,13 +118,14 @@ 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:
ResourceHistoryData: The gathered resource history 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
}
6 changes: 3 additions & 3 deletions robusta_krr/core/integrations/prometheus/metrics/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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.
Expand All @@ -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)
]
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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: ...

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,7 @@ async def gather_data(
object: K8sObjectData,
LoaderClass: type[PrometheusMetric],
period: timedelta,
end_time: datetime,
step: timedelta = timedelta(minutes=30),
) -> PodsTimeData:
"""
Expand All @@ -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 = {}
Expand Down
1 change: 1 addition & 0 deletions robusta_krr/core/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
)

Expand Down