import time from logging import getLogger from defence360agent.internals.feature_flags import ( MESSAGE_LOSS_OBSERVABILITY_FLAG, is_enabled, ) from defence360agent.model.instance import db from defence360agent.model.messages_to_send import MessageToSend logger = getLogger(__name__) class PersistentMessagesQueue: """ The queue to store messages sent to the server if it is unavailable. - stores more recent data; if a limit is exceeded, older messages are deleted. - no duplicate messages are sent NOTE: it is worth remembering that when writing a large number of messages, the amount of memory used may increase by the size of the sqlite cache (this may not be immediately obvious). https://www.sqlite.org/pragma.html#pragma_cache_size """ def __init__(self, buffer_limit=20, storage_limit=1000, model=None): self._buffer_limit = buffer_limit self._storage_limit = storage_limit self._buffer = [] # [(timestamp, message),...] self._model = model or MessageToSend self.dropped_total = 0 self._evicted = 0 def pop_evicted(self) -> int: """Evictions since the last call, then reset (delta for metrics).""" evicted, self._evicted = self._evicted, 0 return evicted def push_buffer_to_storage(self) -> None: if self._buffer: with db.atomic(): # buffer may contain older messages than db, # so remove oldest items after insert self._model.insert_many(self._buffer) need_to_remove = self.storage_size - self._storage_limit if need_to_remove > 0: # keep only the most recent messages removed = self._model.delete_old(need_to_remove) # This is the last point at which the messages exist, so it # is the only place their loss can be reported. self.dropped_total += removed if is_enabled(MESSAGE_LOSS_OBSERVABILITY_FLAG): self._evicted += removed logger.warning( "Persistent message queue overflow: dropped %d oldest" " message(s), storage_limit=%d, dropped_total=%d", removed, self._storage_limit, self.dropped_total, ) self._buffer = [] def pop_all(self) -> list: items = [] with db.atomic(): items += list( self._model.select( self._model.timestamp, self._model.message ).tuples() ) self._model.delete().execute() items += self._buffer self._buffer = [] return sorted(items) # older first def peek_stored(self) -> list: """Return stored rows as (id, timestamp, message) oldest-first without deleting (buffer is neither flushed nor included).""" return list(self._model.get_all_ordered().tuples()) def drain_buffer(self) -> list: """Return and clear the in-memory buffer as (timestamp, message).""" items, self._buffer = self._buffer, [] return items def delete(self, ids: list) -> None: if ids: with db.atomic(): self._model.delete_in(ids) def update_message(self, message_id: int, message: bytes) -> None: with db.atomic(): self._model.set_message(message_id, message) def empty(self) -> bool: return self.qsize() == 0 def qsize(self) -> int: return self.storage_size + len(self._buffer) @property def buffer_size(self) -> int: return len(self._buffer) @property def storage_size(self) -> int: return self._model.select().count() def put(self, message: bytes, timestamp=None): if timestamp is None: timestamp = time.time() self._buffer.append((timestamp, message)) if self.buffer_size >= self._buffer_limit: self.push_buffer_to_storage() def put_many(self, messages: list[tuple[float, bytes]]) -> None: self._buffer.extend(messages) if self.buffer_size >= self._buffer_limit: self.push_buffer_to_storage()