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()