diff --git a/hyrex/dispatcher/sqlc/set_orphaned_task_execution_to_lost_and_retry.py b/hyrex/dispatcher/sqlc/set_orphaned_task_execution_to_lost_and_retry.py index 6c05533..c51af0b 100644 --- a/hyrex/dispatcher/sqlc/set_orphaned_task_execution_to_lost_and_retry.py +++ b/hyrex/dispatcher/sqlc/set_orphaned_task_execution_to_lost_and_retry.py @@ -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 """ diff --git a/hyrex/worker/executor/time_series_averager.py b/hyrex/worker/executor/time_series_averager.py index b7b5411..1f5c681 100644 --- a/hyrex/worker/executor/time_series_averager.py +++ b/hyrex/worker/executor/time_series_averager.py @@ -1,4 +1,5 @@ import time +from collections import deque from pydantic import BaseModel @@ -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 @@ -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, + ) diff --git a/hyrex/worker/root_process.py b/hyrex/worker/root_process.py index e9bb81e..be6cb04 100644 --- a/hyrex/worker/root_process.py +++ b/hyrex/worker/root_process.py @@ -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...") @@ -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.")