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
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@
NOW()
FROM lost_tasks
WHERE attempt_number < max_retries
ON CONFLICT (task_name, idempotency_key) WHERE idempotency_key IS NOT NULL DO NOTHING
"""


Expand Down
20 changes: 15 additions & 5 deletions hyrex/worker/executor/time_series_averager.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import time
from collections import deque

from pydantic import BaseModel

Expand All @@ -13,9 +14,17 @@ class MinuteAverage(BaseModel):
average: float


# Bound the per-averager ring buffer. The executor poll loop submits to these
# averagers on every iteration and previously never pruned, which leaked ~40
# MB/min per worker under an idle busy-poll. 10k entries is enough to hold
# several minutes of stats at realistic submit rates while capping worst-case
# memory at ~2 MB per averager.
_MAX_DATA_POINTS = 10_000


class TimeSeriesAverager:
def __init__(self):
self.data_points: list[DataPoint] = []
self.data_points: deque[DataPoint] = deque(maxlen=_MAX_DATA_POINTS)

def _get_minute_timestamp(self, timestamp: int) -> int:
# Round down to nearest minute
Expand Down Expand Up @@ -71,9 +80,10 @@ def get_current_minute_average(self) -> MinuteAverage:
return result

def clear(self) -> None:
self.data_points = []
self.data_points.clear()

def prune_data_older_than(self, timestamp: int) -> None:
self.data_points = [
point for point in self.data_points if point.timestamp >= timestamp
]
self.data_points = deque(
(point for point in self.data_points if point.timestamp >= timestamp),
maxlen=_MAX_DATA_POINTS,
)
7 changes: 4 additions & 3 deletions hyrex/worker/root_process.py
Original file line number Diff line number Diff line change
Expand Up @@ -227,7 +227,9 @@ def send_heartbeats(self):
)

def run(self):
self.message_listener_thread = threading.Thread(target=self._message_listener)
self.message_listener_thread = threading.Thread(
target=self._message_listener, daemon=True
)
self.message_listener_thread.start()
self.logger.info("Incoming message queue now active...")

Expand Down Expand Up @@ -332,8 +334,7 @@ def stop(self):
self.message_listener_thread.join(timeout=5.0)
if self.message_listener_thread.is_alive():
self.logger.warning("Message listener thread did not exit cleanly within timeout.")
# Force terminate the thread by setting it as daemon and exiting
# Python will clean it up on process exit
# Thread is a daemon, so the interpreter will terminate it on process exit.
else:
self.logger.info("Message listener thread closed successfully.")

Expand Down