Refactor common.logging to include both syslog handlers
This commit is contained in:
@@ -1,540 +0,0 @@
|
|||||||
import collections
|
|
||||||
import contextlib
|
|
||||||
from dataclasses import dataclass, field
|
|
||||||
import logging
|
|
||||||
import os
|
|
||||||
import queue
|
|
||||||
import socket
|
|
||||||
import ssl
|
|
||||||
import threading
|
|
||||||
import time
|
|
||||||
|
|
||||||
from typing import Optional, Tuple
|
|
||||||
|
|
||||||
from hat.syslog import common, encoder
|
|
||||||
from hat.syslog.handler import (
|
|
||||||
SyslogHandler as hat_syslog_handler_SyslogHandler,
|
|
||||||
_ThreadState,
|
|
||||||
_create_dropped_msg,
|
|
||||||
_record_to_msg,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
_SOCKET_FACTORIES = {
|
|
||||||
common.CommType.TCP: lambda state, ctx: _create_tcp_socket(state),
|
|
||||||
common.CommType.TLS: lambda state, ctx: _create_tcp_socket(state, ctx),
|
|
||||||
common.CommType.UDP: lambda state, ctx: _create_udp_socket(state),
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True)
|
|
||||||
class RetryItem:
|
|
||||||
"""Encapsulates a structured syslog entry in the retry transport pipeline."""
|
|
||||||
synthetic: bool
|
|
||||||
msg: common.Msg
|
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
|
||||||
class ThreadScoreboard:
|
|
||||||
"""Tracks precision execution lifecycles and diagnostic markers for background workers."""
|
|
||||||
start: Tuple[float, int] = field(default_factory=lambda: (time.time(), time.monotonic_ns()))
|
|
||||||
alive: Optional[Tuple[float, int]] = None
|
|
||||||
initialized: Optional[Tuple[float, int]] = None
|
|
||||||
previous_start: Optional[Tuple[float, int]] = None
|
|
||||||
|
|
||||||
|
|
||||||
def _create_tcp_socket(state, ctx=None):
|
|
||||||
"""Establishes an optimized TCP or wrapped TLS stream transport connection."""
|
|
||||||
s = socket.create_connection((state.host, state.port), timeout=5.0)
|
|
||||||
s.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
|
|
||||||
if ctx:
|
|
||||||
s = ctx.wrap_socket(s)
|
|
||||||
return s
|
|
||||||
|
|
||||||
|
|
||||||
def _create_udp_socket(state):
|
|
||||||
"""Establishes an un-bonded UDP datagram socket connection endpoint."""
|
|
||||||
s = socket.socket(type=socket.SOCK_DGRAM)
|
|
||||||
s.connect((state.host, state.port))
|
|
||||||
return s
|
|
||||||
|
|
||||||
|
|
||||||
def _item_completed(retry_queue, core_queue, item):
|
|
||||||
"""
|
|
||||||
Drains the processed item from the retry queue and issues a task acknowledgment
|
|
||||||
to the synchronized queue if the item was not synthetically created.
|
|
||||||
"""
|
|
||||||
# SUCCESS: Remove the successfully sent item from the chronological pipeline
|
|
||||||
retry_queue.popleft()
|
|
||||||
if not item.synthetic:
|
|
||||||
core_queue.task_done()
|
|
||||||
|
|
||||||
|
|
||||||
def _logging_handler_thread(state, logger=logger):
|
|
||||||
"""
|
|
||||||
Worker thread that drains messages and guarantees strict transport-level delivery
|
|
||||||
by checking stateless item origins packed inside RetryItem containers.
|
|
||||||
|
|
||||||
Independent worker thread routine responsible for draining log messages
|
|
||||||
from the synchronized queue and streaming them to the remote hat-syslog endpoint.
|
|
||||||
|
|
||||||
Uses an internal thread-local retry queue to maintain precise chronological
|
|
||||||
ordering of log messages during transport connection failures.
|
|
||||||
|
|
||||||
Args:
|
|
||||||
state (_ThreadState): Thread-safe tracking state containing connection params.
|
|
||||||
logger (logging.Logger): Diagnostic logger instance for transport anomalies.
|
|
||||||
"""
|
|
||||||
ctx = None
|
|
||||||
if common.CommType.TLS == state.comm_type:
|
|
||||||
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT)
|
|
||||||
ctx.check_hostname = False
|
|
||||||
ctx.verify_mode = ssl.VerifyMode.CERT_NONE
|
|
||||||
|
|
||||||
# Chronological staging queue dedicated to maintaining strict FIFO order during outages
|
|
||||||
retry_queue = collections.deque()
|
|
||||||
|
|
||||||
# Loop persists on shutdown until retry_queue is fully empty.
|
|
||||||
while retry_queue or not state.closed.is_set():
|
|
||||||
# connect to the endpoint
|
|
||||||
s = None
|
|
||||||
try:
|
|
||||||
factory = _SOCKET_FACTORIES[state.comm_type]
|
|
||||||
s = factory(state, ctx)
|
|
||||||
except KeyError:
|
|
||||||
raise NotImplementedError(f'Unsupported comm_type: {state.comm_type}')
|
|
||||||
except Exception:
|
|
||||||
if s is not None:
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
s.close()
|
|
||||||
time.sleep(state.reconnect_delay)
|
|
||||||
continue
|
|
||||||
|
|
||||||
# Connection successfully established; enter message transmission loop
|
|
||||||
# while True optimizes throughput and relies on socket exceptions to break the loop context
|
|
||||||
try:
|
|
||||||
captured_drops = ()
|
|
||||||
drop_payload = None
|
|
||||||
item = None
|
|
||||||
msg = None
|
|
||||||
msg_bytes = None
|
|
||||||
|
|
||||||
while True:
|
|
||||||
# Harvest and process dropped counts
|
|
||||||
# SAFELY EXTRACT AND PROCESS ALL TRACKED OVERFLOWS CHRONOLOGICALLY
|
|
||||||
try:
|
|
||||||
with state.cv:
|
|
||||||
if 1 < len(state.dropped) or 0 < state.dropped[0]:
|
|
||||||
# Freeze and copy the entire sequence array out of shared memory
|
|
||||||
captured_drops = tuple(state.dropped)
|
|
||||||
# Re-instantiate the tracked array to a clean initial state instantly
|
|
||||||
state.dropped.clear()
|
|
||||||
state.dropped.append(0)
|
|
||||||
|
|
||||||
# Convert captured thresholds into synthetic inline logs inside our local staging worker
|
|
||||||
for chunked_count in captured_drops:
|
|
||||||
if 0 < chunked_count:
|
|
||||||
drop_payload = _create_dropped_msg(
|
|
||||||
chunked_count, '_logging_handler_thread', 0,
|
|
||||||
)
|
|
||||||
# Appending to the retry queue guarantees strict chronological reporting order
|
|
||||||
retry_queue.append(RetryItem(synthetic=True, msg=drop_payload))
|
|
||||||
finally:
|
|
||||||
captured_drops = ()
|
|
||||||
drop_payload = None
|
|
||||||
|
|
||||||
# Grab from the main queue if the local transport staging queue is empty
|
|
||||||
# If the retry queue is empty, block and wait for a fresh log message
|
|
||||||
if not retry_queue:
|
|
||||||
try:
|
|
||||||
msg = state.queue.get(timeout=state.reconnect_delay)
|
|
||||||
except queue.Empty:
|
|
||||||
if state.closed.is_set():
|
|
||||||
# The queue is empty and the handler has been explicitly closed.
|
|
||||||
# Break the transmission loop to allow the worker thread to exit cleanly.
|
|
||||||
break
|
|
||||||
continue
|
|
||||||
else:
|
|
||||||
retry_queue.append(RetryItem(synthetic=False, msg=msg))
|
|
||||||
finally:
|
|
||||||
msg = None
|
|
||||||
|
|
||||||
# Transmit head message and track task lifecycle states
|
|
||||||
try:
|
|
||||||
# Peek at the oldest message without popping it yet
|
|
||||||
item = retry_queue[0]
|
|
||||||
msg_bytes = encoder.msg_to_str(item.msg).encode()
|
|
||||||
|
|
||||||
if common.CommType.UDP == state.comm_type:
|
|
||||||
s.send(msg_bytes)
|
|
||||||
else:
|
|
||||||
s.send(f'{len(msg_bytes)} '.encode() + msg_bytes)
|
|
||||||
except (TypeError, ValueError):
|
|
||||||
# Message structure failure: Drop item immediately and preserve connection state
|
|
||||||
logger.exception('Dropping un-convertible poison-pill log message')
|
|
||||||
_item_completed(retry_queue, state.queue, item)
|
|
||||||
except UnicodeEncodeError:
|
|
||||||
# String binary processing failure: Drop item immediately and preserve connection state
|
|
||||||
logger.exception('Dropping un-encodable Unicode log message string')
|
|
||||||
_item_completed(retry_queue, state.queue, item)
|
|
||||||
except Exception:
|
|
||||||
# On socket break, tear down this loop context cleanly.
|
|
||||||
# The current message remains cleanly preserved at index 0 of retry_queue.
|
|
||||||
# Because task_done() is skipped, state.queue.join() will continue to block.
|
|
||||||
# Connection lost or infrastructure network failure: Break out loop to trigger socket reconnect
|
|
||||||
# Element stays safe at index 0 of the retry_queue cache.
|
|
||||||
# After reconnecting we will try it again.
|
|
||||||
break
|
|
||||||
else:
|
|
||||||
_item_completed(retry_queue, state.queue, item)
|
|
||||||
finally:
|
|
||||||
# Clear references to optimize memory tracking
|
|
||||||
item = None
|
|
||||||
msg_bytes = None
|
|
||||||
finally:
|
|
||||||
# close the connection to avoid leaking it
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
s.close()
|
|
||||||
|
|
||||||
|
|
||||||
class SyslogHandler(hat_syslog_handler_SyslogHandler):
|
|
||||||
"""
|
|
||||||
A process-safe wrapper for hat.syslog.handler.SyslogHandler.
|
|
||||||
|
|
||||||
Bypasses immutable NamedTuple state constraints on fork boundaries,
|
|
||||||
avoids thread-lock corruption from os.fork(), and short-circuits
|
|
||||||
the time-blocking flush/close loops during a Huey graceful shutdown.
|
|
||||||
|
|
||||||
A process-safe wrapper subclass for hat.syslog.handler.SyslogHandler.
|
|
||||||
|
|
||||||
Bypasses immutable state limitations on fork boundaries, handles thread-lock
|
|
||||||
corruption resulting from os.fork(), and avoids blocking task-execution loops
|
|
||||||
during a graceful worker process shutdown.
|
|
||||||
"""
|
|
||||||
def __init__(self, host, port, comm_type, queue_size=1024, reconnect_delay=5, *args, **kwargs):
|
|
||||||
"""Initializes the wrapper and neutralizes conflicting parent process states."""
|
|
||||||
super().__init__(host, port, common.CommType.UDP, queue_size, reconnect_delay, *args, **kwargs)
|
|
||||||
|
|
||||||
state = self._get_parent_attr('__state')
|
|
||||||
|
|
||||||
if state:
|
|
||||||
state.closed.set()
|
|
||||||
thread = self._get_parent_attr('__thread')
|
|
||||||
if thread and thread.is_alive():
|
|
||||||
with state.cv, contextlib.suppress(Exception):
|
|
||||||
state.cv.notify_all()
|
|
||||||
|
|
||||||
# Initialize native, thread-safe sync tracking structures
|
|
||||||
self.__state = _ThreadState(
|
|
||||||
host=host,
|
|
||||||
port=port,
|
|
||||||
comm_type=self._determine_comm_type(comm_type),
|
|
||||||
queue=queue.Queue(maxsize=queue_size),
|
|
||||||
queue_size=queue_size,
|
|
||||||
reconnect_delay=reconnect_delay,
|
|
||||||
cv=threading.Condition(),
|
|
||||||
closed=threading.Event(),
|
|
||||||
dropped=list((0,)),
|
|
||||||
)
|
|
||||||
|
|
||||||
self.__thread = None
|
|
||||||
self._initial_pid = os.getpid()
|
|
||||||
self._closing = threading.Event()
|
|
||||||
|
|
||||||
def _alive_thread(self):
|
|
||||||
"""Thread-safe validator confirming background worker availability."""
|
|
||||||
if not (self.__thread and self.__thread.is_alive()):
|
|
||||||
return False
|
|
||||||
|
|
||||||
with self.__state.cv:
|
|
||||||
if self.__thread and self.__thread.is_alive():
|
|
||||||
if hasattr(self.__thread, '_scoreboard') and self.__thread._scoreboard:
|
|
||||||
self.__thread._scoreboard.alive = (time.time(), time.monotonic_ns())
|
|
||||||
return True
|
|
||||||
return False
|
|
||||||
|
|
||||||
def _after_fork(self, current_pid=None):
|
|
||||||
"""
|
|
||||||
Intercepts the Unix process boundary skew. If a fork is identified,
|
|
||||||
it cleans old references and initializes fresh process-isolated primitives.
|
|
||||||
"""
|
|
||||||
|
|
||||||
# Detect if we crossed the Unix fork boundary into Huey's process worker.
|
|
||||||
if current_pid is None:
|
|
||||||
current_pid = os.getpid()
|
|
||||||
|
|
||||||
if current_pid == self._initial_pid:
|
|
||||||
# Return early when we have not forked.
|
|
||||||
return
|
|
||||||
|
|
||||||
self._initial_pid = current_pid
|
|
||||||
|
|
||||||
state = self.__state
|
|
||||||
new_state = _ThreadState(
|
|
||||||
host=state.host,
|
|
||||||
port=state.port,
|
|
||||||
comm_type=state.comm_type,
|
|
||||||
queue=queue.Queue(maxsize=state.queue_size),
|
|
||||||
queue_size=state.queue_size,
|
|
||||||
reconnect_delay=state.reconnect_delay,
|
|
||||||
cv=threading.Condition(),
|
|
||||||
closed=threading.Event(),
|
|
||||||
dropped=list((0,)),
|
|
||||||
)
|
|
||||||
self.__state = new_state
|
|
||||||
self.__thread = None
|
|
||||||
|
|
||||||
def _create_thread(self):
|
|
||||||
"""Aligns process states and builds a clean worker thread context."""
|
|
||||||
|
|
||||||
self._after_fork()
|
|
||||||
|
|
||||||
if self._closing.is_set() or self.__state.closed.is_set() or self._alive_thread():
|
|
||||||
return
|
|
||||||
|
|
||||||
initial = self.__thread is None
|
|
||||||
with self.__state.cv:
|
|
||||||
previous = self.__thread
|
|
||||||
|
|
||||||
self.__thread = threading.Thread(
|
|
||||||
target=_logging_handler_thread,
|
|
||||||
args=(self.__state,),
|
|
||||||
daemon=True,
|
|
||||||
)
|
|
||||||
|
|
||||||
scoreboard = ThreadScoreboard()
|
|
||||||
self.__thread._scoreboard = scoreboard
|
|
||||||
|
|
||||||
if initial:
|
|
||||||
scoreboard.alive = scoreboard.start
|
|
||||||
scoreboard.initialized = scoreboard.start
|
|
||||||
elif previous and hasattr(previous, '_scoreboard') and previous._scoreboard:
|
|
||||||
scoreboard.alive = previous._scoreboard.alive
|
|
||||||
scoreboard.initialized = previous._scoreboard.initialized
|
|
||||||
scoreboard.previous_start = previous._scoreboard.start
|
|
||||||
logger.debug(f'Created a replacement thread: {scoreboard=}')
|
|
||||||
|
|
||||||
self.__thread.start()
|
|
||||||
|
|
||||||
def _determine_comm_type(self, comm_type):
|
|
||||||
"""Maps and cross-checks string connection descriptions to Enumeration definitions."""
|
|
||||||
if isinstance(comm_type, str):
|
|
||||||
needle = comm_type
|
|
||||||
haystack = frozenset(common.CommType.__members__)
|
|
||||||
vary = lambda x: {
|
|
||||||
x, x.upper(),
|
|
||||||
x.casefold(), x.casefold().upper(),
|
|
||||||
x.lower(), x.lower().upper(),
|
|
||||||
}
|
|
||||||
try:
|
|
||||||
matched_elements = tuple(haystack.intersection(vary(needle)))
|
|
||||||
member = matched_elements[0]
|
|
||||||
return common.CommType[member]
|
|
||||||
except (IndexError, KeyError) as e:
|
|
||||||
raise ValueError(f'Specify a valid comm_type from this list: {list(haystack)}') from e
|
|
||||||
|
|
||||||
if not isinstance(comm_type, common.CommType):
|
|
||||||
raise ValueError('Invalid comm_type argument')
|
|
||||||
|
|
||||||
def _parent_class_name(self):
|
|
||||||
return hat_syslog_handler_SyslogHandler.__name__
|
|
||||||
|
|
||||||
def _mangled_name(self, attr_name):
|
|
||||||
return f'_{self._parent_class_name()}{attr_name}'
|
|
||||||
|
|
||||||
def _get_parent_attr(self, attr_name):
|
|
||||||
"""Computes and gets mangled attributes from the super class."""
|
|
||||||
return getattr(self, self._mangled_name(attr_name), None)
|
|
||||||
|
|
||||||
def _set_parent_attr(self, attr_name, value):
|
|
||||||
"""Computes and sets mangled attributes on the super class."""
|
|
||||||
setattr(self, self._mangled_name(attr_name), value)
|
|
||||||
|
|
||||||
def emit(self, record):
|
|
||||||
"""Enqueues new log records and guarantees active connection coverage."""
|
|
||||||
if self._closing.is_set():
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
logger.handle(record)
|
|
||||||
return
|
|
||||||
|
|
||||||
self._create_thread()
|
|
||||||
|
|
||||||
if not self._alive_thread():
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
logger.handle(record)
|
|
||||||
return
|
|
||||||
|
|
||||||
state = self.__state
|
|
||||||
if state.closed.is_set():
|
|
||||||
self._closing.set()
|
|
||||||
logger.warning('Closed in emit')
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
logger.handle(record)
|
|
||||||
return
|
|
||||||
|
|
||||||
msg = _record_to_msg(record)
|
|
||||||
|
|
||||||
try:
|
|
||||||
state.queue.put_nowait(msg)
|
|
||||||
except queue.Full:
|
|
||||||
# ACQUIRE LOCK ON MAIN THREAD BEFORE INCREMENTING COUNTER SLICES
|
|
||||||
with state.cv:
|
|
||||||
dropped_count = state.dropped[-1]
|
|
||||||
if 1_000_000 < dropped_count:
|
|
||||||
state.dropped.append(1)
|
|
||||||
else:
|
|
||||||
state.dropped[-1] = 1 + dropped_count
|
|
||||||
|
|
||||||
logger.warning(f'Dropped a log message in emit due to buffer overflow: {msg.msg!r}')
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
logger.handle(record)
|
|
||||||
|
|
||||||
def flush(self):
|
|
||||||
"""Blocks execution until the internal logging queue is empty."""
|
|
||||||
self._create_thread()
|
|
||||||
|
|
||||||
if not self._alive_thread():
|
|
||||||
return
|
|
||||||
|
|
||||||
state = self.__state
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
state.queue.join()
|
|
||||||
|
|
||||||
def close(self):
|
|
||||||
"""
|
|
||||||
Gracefully flushes the queue and terminates the background logging thread.
|
|
||||||
Cleans up the queue, flags the state as closed, and shuts down
|
|
||||||
the background thread without allowing new ones to be generated.
|
|
||||||
"""
|
|
||||||
|
|
||||||
# Align/verify the background worker thread state immediately
|
|
||||||
self._closing.clear()
|
|
||||||
self._create_thread()
|
|
||||||
state = self.__state
|
|
||||||
if state.closed.is_set():
|
|
||||||
# Only return early when the thread is alive
|
|
||||||
state.closed.clear()
|
|
||||||
self._create_thread()
|
|
||||||
self._closing.set()
|
|
||||||
|
|
||||||
# The native queue.join() blockade is now perfectly synchronized with the internal
|
|
||||||
# tracking flow loop. It will block until retry_queue is 100% empty.
|
|
||||||
if self._alive_thread():
|
|
||||||
logger.debug('Flushing logging queue in close')
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
state.queue.join()
|
|
||||||
|
|
||||||
# Immediately trip the closed flag to end the networking thread
|
|
||||||
self.__state.closed.set()
|
|
||||||
|
|
||||||
# =====================================================================
|
|
||||||
# DYNAMIC GRANDPARENT BYPASS VIA MRO
|
|
||||||
# Instructs Python to search for close() starting *after* our direct
|
|
||||||
# parent class type descriptor. This dynamically resolves grandfather
|
|
||||||
# dependencies while completely avoiding the parent class thread-joins.
|
|
||||||
# =====================================================================
|
|
||||||
super(hat_syslog_handler_SyslogHandler, self).close()
|
|
||||||
|
|
||||||
|
|
||||||
if '__main__' == __name__:
|
|
||||||
import unittest
|
|
||||||
|
|
||||||
class MockSyslogServer:
|
|
||||||
"""Stands up an isolated local background socket server to harvest transport streams."""
|
|
||||||
def __init__(self, host='127.0.0.1', port=0):
|
|
||||||
self.host = host
|
|
||||||
self.port = port
|
|
||||||
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
||||||
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
|
|
||||||
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
||||||
self.sock.bind((self.host, self.port))
|
|
||||||
self.port = self.sock.getsockname()[1]
|
|
||||||
|
|
||||||
self.received_messages = []
|
|
||||||
self.running = threading.Event()
|
|
||||||
self._thread = None
|
|
||||||
|
|
||||||
def start(self):
|
|
||||||
self.running.set()
|
|
||||||
self._thread = threading.Thread(target=self._listen_loop, daemon=True)
|
|
||||||
self._thread.start()
|
|
||||||
|
|
||||||
def stop(self):
|
|
||||||
self.running.clear()
|
|
||||||
s = None
|
|
||||||
try:
|
|
||||||
s = socket.create_connection((self.host, self.port), timeout=0.1)
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
finally:
|
|
||||||
if s is not None:
|
|
||||||
s.close()
|
|
||||||
if self._thread:
|
|
||||||
self._thread.join(timeout=1.0)
|
|
||||||
self.sock.close()
|
|
||||||
|
|
||||||
def _listen_loop(self):
|
|
||||||
self.sock.listen(1)
|
|
||||||
while self.running.is_set():
|
|
||||||
try:
|
|
||||||
conn, _ = self.sock.accept()
|
|
||||||
if not self.running.is_set():
|
|
||||||
conn.close()
|
|
||||||
break
|
|
||||||
|
|
||||||
with conn:
|
|
||||||
while self.running.is_set():
|
|
||||||
data = conn.recv(4096)
|
|
||||||
if not data:
|
|
||||||
break
|
|
||||||
self.received_messages.append(data.decode('utf-8'))
|
|
||||||
except Exception:
|
|
||||||
break
|
|
||||||
|
|
||||||
class TestSyslogHandlerIntegration(unittest.TestCase):
|
|
||||||
def setUp(self):
|
|
||||||
"""Initializes the background mock network collection service before running assertions."""
|
|
||||||
self.server = MockSyslogServer()
|
|
||||||
self.server.start()
|
|
||||||
|
|
||||||
self.handler = SyslogHandler(
|
|
||||||
host=self.server.host,
|
|
||||||
port=self.server.port,
|
|
||||||
comm_type='tcp',
|
|
||||||
queue_size=10,
|
|
||||||
reconnect_delay=1,
|
|
||||||
)
|
|
||||||
|
|
||||||
self.test_logger = logging.getLogger('integration_test')
|
|
||||||
self.test_logger.setLevel(logging.DEBUG)
|
|
||||||
self.test_logger.addHandler(self.handler)
|
|
||||||
|
|
||||||
def tearDown(self):
|
|
||||||
"""Cleans up the network service topology profiles safely upon validation teardown."""
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
self.handler.close()
|
|
||||||
self.server.stop()
|
|
||||||
|
|
||||||
def test_pipeline_delivery_and_flush(self):
|
|
||||||
"""Verifies that items are completely delivered down the wire before flush unblocks."""
|
|
||||||
self.test_logger.debug('Message A')
|
|
||||||
self.test_logger.info('Message B')
|
|
||||||
|
|
||||||
start_time = time.monotonic()
|
|
||||||
self.handler.flush()
|
|
||||||
elapsed = time.monotonic() - start_time
|
|
||||||
|
|
||||||
self.assertLess(elapsed, 2.0, 'The flush operations deadlocked the execution loop context')
|
|
||||||
self.assertTrue(any('Message A' in msg for msg in self.server.received_messages))
|
|
||||||
self.assertTrue(any('Message B' in msg for msg in self.server.received_messages))
|
|
||||||
|
|
||||||
def test_graceful_close_lifecycle(self):
|
|
||||||
"""Confirms that close drains remaining log states and tears down the worker thread."""
|
|
||||||
self.test_logger.info('Shutdown Message')
|
|
||||||
self.handler.close()
|
|
||||||
self.assertTrue(any('Shutdown Message' in msg for msg in self.server.received_messages))
|
|
||||||
|
|
||||||
unittest.main()
|
|
||||||
|
|
||||||
@@ -1,6 +1,7 @@
|
|||||||
import logging
|
import logging
|
||||||
from django.conf import settings
|
from django.conf import settings
|
||||||
##from .logging import default_handler, syslog_handler
|
##from .logging import default_handler
|
||||||
|
##from .logging.syslog.std import default_handler as syslog_handler
|
||||||
from .utils import getenv
|
from .utils import getenv
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
10
tubesync/common/logging/__init__.py
Normal file
10
tubesync/common/logging/__init__.py
Normal file
@@ -0,0 +1,10 @@
|
|||||||
|
from . import syslog
|
||||||
|
from ._default import default_formatter, default_handler
|
||||||
|
from ._filters import RemoveSpecificLogFilter
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
'default_formatter',
|
||||||
|
'default_handler',
|
||||||
|
'syslog',
|
||||||
|
'RemoveSpecificLogFilter',
|
||||||
|
]
|
||||||
12
tubesync/common/logging/_default.py
Normal file
12
tubesync/common/logging/_default.py
Normal file
@@ -0,0 +1,12 @@
|
|||||||
|
import logging
|
||||||
|
|
||||||
|
|
||||||
|
default_formatter = logging.Formatter(
|
||||||
|
'%(asctime)s [%(name)s/%(levelname)s] %(message)s'
|
||||||
|
)
|
||||||
|
|
||||||
|
default_handler = logging.StreamHandler()
|
||||||
|
default_handler.setFormatter(default_formatter)
|
||||||
|
default_handler.setLevel(logging.INFO)
|
||||||
|
|
||||||
|
__all__ = ['default_formatter', 'default_handler']
|
||||||
@@ -1,5 +1,4 @@
|
|||||||
import logging
|
import logging
|
||||||
from logging.handlers import SysLogHandler
|
|
||||||
|
|
||||||
|
|
||||||
class RemoveSpecificLogFilter(logging.Filter):
|
class RemoveSpecificLogFilter(logging.Filter):
|
||||||
@@ -69,21 +68,7 @@ class RemoveSpecificLogFilter(logging.Filter):
|
|||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
default_formatter = logging.Formatter(
|
__all__ = [
|
||||||
'%(asctime)s [%(name)s/%(levelname)s] %(message)s'
|
'RemoveSpecificLogFilter',
|
||||||
)
|
]
|
||||||
default_handler = logging.StreamHandler()
|
|
||||||
default_handler.setFormatter(default_formatter)
|
|
||||||
default_handler.setLevel(logging.INFO)
|
|
||||||
|
|
||||||
syslog_formatter = logging.Formatter(
|
|
||||||
'%(asctime)s %(name)s: %(message)s',
|
|
||||||
'%b %d %H:%M:%S',
|
|
||||||
)
|
|
||||||
syslog_handler = SysLogHandler(
|
|
||||||
address='/dev/log',
|
|
||||||
facility=SysLogHandler.LOG_LOCAL0,
|
|
||||||
)
|
|
||||||
syslog_handler.setFormatter(syslog_formatter)
|
|
||||||
syslog_handler.setLevel(logging.DEBUG)
|
|
||||||
|
|
||||||
2
tubesync/common/logging/syslog/__init__.py
Normal file
2
tubesync/common/logging/syslog/__init__.py
Normal file
@@ -0,0 +1,2 @@
|
|||||||
|
from . import hat as hat
|
||||||
|
from . import std as std
|
||||||
3
tubesync/common/logging/syslog/hat/__init__.py
Normal file
3
tubesync/common/logging/syslog/hat/__init__.py
Normal file
@@ -0,0 +1,3 @@
|
|||||||
|
from ._default import * # noqa: F403
|
||||||
|
from ._logger import logger as logger
|
||||||
|
logger = logger(__name__)
|
||||||
106
tubesync/common/logging/syslog/hat/__main__.py
Normal file
106
tubesync/common/logging/syslog/hat/__main__.py
Normal file
@@ -0,0 +1,106 @@
|
|||||||
|
import contextlib
|
||||||
|
import logging
|
||||||
|
import socket
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
|
import unittest
|
||||||
|
|
||||||
|
from ._default import handler
|
||||||
|
|
||||||
|
|
||||||
|
class MockSyslogServer:
|
||||||
|
"""Stands up an isolated local background socket server to harvest transport streams."""
|
||||||
|
def __init__(self, host='127.0.0.1', port=0):
|
||||||
|
self.host = host
|
||||||
|
self.port = port
|
||||||
|
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||||
|
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
|
||||||
|
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||||
|
self.sock.bind((self.host, self.port))
|
||||||
|
self.port = self.sock.getsockname()[1]
|
||||||
|
|
||||||
|
self.received_messages = []
|
||||||
|
self.running = threading.Event()
|
||||||
|
self._thread = None
|
||||||
|
|
||||||
|
def start(self):
|
||||||
|
self.running.set()
|
||||||
|
self._thread = threading.Thread(target=self._listen_loop, daemon=True)
|
||||||
|
self._thread.start()
|
||||||
|
|
||||||
|
def stop(self):
|
||||||
|
self.running.clear()
|
||||||
|
s = None
|
||||||
|
try:
|
||||||
|
s = socket.create_connection((self.host, self.port), timeout=0.1)
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
if s is not None:
|
||||||
|
s.close()
|
||||||
|
if self._thread:
|
||||||
|
self._thread.join(timeout=1.0)
|
||||||
|
self.sock.close()
|
||||||
|
|
||||||
|
def _listen_loop(self):
|
||||||
|
self.sock.listen(1)
|
||||||
|
while self.running.is_set():
|
||||||
|
try:
|
||||||
|
conn, _ = self.sock.accept()
|
||||||
|
if not self.running.is_set():
|
||||||
|
conn.close()
|
||||||
|
break
|
||||||
|
|
||||||
|
with conn:
|
||||||
|
while self.running.is_set():
|
||||||
|
data = conn.recv(4096)
|
||||||
|
if not data:
|
||||||
|
break
|
||||||
|
self.received_messages.append(data.decode('utf-8'))
|
||||||
|
except Exception:
|
||||||
|
break
|
||||||
|
|
||||||
|
class TestSyslogHandlerIntegration(unittest.TestCase):
|
||||||
|
def setUp(self):
|
||||||
|
"""Initializes the background mock network collection service before running assertions."""
|
||||||
|
self.server = MockSyslogServer()
|
||||||
|
self.server.start()
|
||||||
|
|
||||||
|
self.handler = handler(
|
||||||
|
host=self.server.host,
|
||||||
|
port=self.server.port,
|
||||||
|
comm_type='tcp',
|
||||||
|
queue_size=10,
|
||||||
|
reconnect_delay=1,
|
||||||
|
)
|
||||||
|
|
||||||
|
self.test_logger = logging.getLogger('integration_test')
|
||||||
|
self.test_logger.setLevel(logging.DEBUG)
|
||||||
|
self.test_logger.addHandler(self.handler)
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
"""Cleans up the network service topology profiles safely upon validation teardown."""
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
self.handler.close()
|
||||||
|
self.server.stop()
|
||||||
|
|
||||||
|
def test_pipeline_delivery_and_flush(self):
|
||||||
|
"""Verifies that items are completely delivered down the wire before flush unblocks."""
|
||||||
|
self.test_logger.debug('Message A')
|
||||||
|
self.test_logger.info('Message B')
|
||||||
|
|
||||||
|
start_time = time.monotonic()
|
||||||
|
self.handler.flush()
|
||||||
|
elapsed = time.monotonic() - start_time
|
||||||
|
|
||||||
|
self.assertLess(elapsed, 2.0, 'The flush operations deadlocked the execution loop context')
|
||||||
|
self.assertTrue(any('Message A' in msg for msg in self.server.received_messages))
|
||||||
|
self.assertTrue(any('Message B' in msg for msg in self.server.received_messages))
|
||||||
|
|
||||||
|
def test_graceful_close_lifecycle(self):
|
||||||
|
"""Confirms that close drains remaining log states and tears down the worker thread."""
|
||||||
|
self.test_logger.info('Shutdown Message')
|
||||||
|
self.handler.close()
|
||||||
|
self.assertTrue(any('Shutdown Message' in msg for msg in self.server.received_messages))
|
||||||
|
|
||||||
|
unittest.main()
|
||||||
469
tubesync/common/logging/syslog/hat/_default.py
Normal file
469
tubesync/common/logging/syslog/hat/_default.py
Normal file
@@ -0,0 +1,469 @@
|
|||||||
|
import collections
|
||||||
|
import contextlib
|
||||||
|
from dataclasses import dataclass, field
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import queue
|
||||||
|
import socket
|
||||||
|
import ssl
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
|
|
||||||
|
from typing import Optional, Tuple
|
||||||
|
|
||||||
|
from ._logger import logger
|
||||||
|
|
||||||
|
try:
|
||||||
|
from hat.syslog import common, encoder
|
||||||
|
from hat.syslog.handler import (
|
||||||
|
SyslogHandler as hat_syslog_handler_SyslogHandler,
|
||||||
|
_ThreadState,
|
||||||
|
_create_dropped_msg,
|
||||||
|
_record_to_msg,
|
||||||
|
)
|
||||||
|
except ImportError:
|
||||||
|
handler = None
|
||||||
|
else:
|
||||||
|
handler = True
|
||||||
|
|
||||||
|
default_formatter = logging.Formatter()
|
||||||
|
|
||||||
|
logger = logger()
|
||||||
|
|
||||||
|
__all__ = ['default_formatter', 'default_handler', 'handler']
|
||||||
|
|
||||||
|
if not handler:
|
||||||
|
# Create only enough for tests to fail instead of creating hard to diagnose errors
|
||||||
|
from ..std import default_handler as std_default_handler, handler
|
||||||
|
class MockSyslogHandler(handler):
|
||||||
|
def __init__(self, host, port, comm_type, queue_size, reconnect_delay, *args, **kwargs):
|
||||||
|
args = ()
|
||||||
|
kwargs = {}
|
||||||
|
kwargs['address'] = std_default_handler.address
|
||||||
|
kwargs['facility'] = std_default_handler.facility
|
||||||
|
super().__init__(*args, **kwargs)
|
||||||
|
handler = MockSyslogHandler
|
||||||
|
default_handler = handler('127.0.0.1', 6514, 'UDP', 1024, 1)
|
||||||
|
else:
|
||||||
|
_SOCKET_FACTORIES = {
|
||||||
|
common.CommType.TCP: lambda state, ctx: _create_tcp_socket(state),
|
||||||
|
common.CommType.TLS: lambda state, ctx: _create_tcp_socket(state, ctx),
|
||||||
|
common.CommType.UDP: lambda state, ctx: _create_udp_socket(state),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class RetryItem:
|
||||||
|
"""Encapsulates a structured syslog entry in the retry transport pipeline."""
|
||||||
|
synthetic: bool
|
||||||
|
msg: common.Msg
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class ThreadScoreboard:
|
||||||
|
"""Tracks precision execution lifecycles and diagnostic markers for background workers."""
|
||||||
|
start: Tuple[float, int] = field(default_factory=lambda: (time.time(), time.monotonic_ns()))
|
||||||
|
alive: Optional[Tuple[float, int]] = None
|
||||||
|
initialized: Optional[Tuple[float, int]] = None
|
||||||
|
previous_start: Optional[Tuple[float, int]] = None
|
||||||
|
|
||||||
|
|
||||||
|
def _create_tcp_socket(state, ctx=None):
|
||||||
|
"""Establishes an optimized TCP or wrapped TLS stream transport connection."""
|
||||||
|
s = socket.create_connection((state.host, state.port), timeout=5.0)
|
||||||
|
s.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
|
||||||
|
if ctx:
|
||||||
|
s = ctx.wrap_socket(s)
|
||||||
|
return s
|
||||||
|
|
||||||
|
|
||||||
|
def _create_udp_socket(state):
|
||||||
|
"""Establishes an un-bonded UDP datagram socket connection endpoint."""
|
||||||
|
s = socket.socket(type=socket.SOCK_DGRAM)
|
||||||
|
s.connect((state.host, state.port))
|
||||||
|
return s
|
||||||
|
|
||||||
|
|
||||||
|
def _item_completed(retry_queue, core_queue, item):
|
||||||
|
"""
|
||||||
|
Drains the processed item from the retry queue and issues a task acknowledgment
|
||||||
|
to the synchronized queue if the item was not synthetically created.
|
||||||
|
"""
|
||||||
|
# SUCCESS: Remove the successfully sent item from the chronological pipeline
|
||||||
|
retry_queue.popleft()
|
||||||
|
if not item.synthetic:
|
||||||
|
core_queue.task_done()
|
||||||
|
|
||||||
|
|
||||||
|
def _logging_handler_thread(state, logger=logger):
|
||||||
|
"""
|
||||||
|
Worker thread that drains messages and guarantees strict transport-level delivery
|
||||||
|
by checking stateless item origins packed inside RetryItem containers.
|
||||||
|
|
||||||
|
Independent worker thread routine responsible for draining log messages
|
||||||
|
from the synchronized queue and streaming them to the remote hat-syslog endpoint.
|
||||||
|
|
||||||
|
Uses an internal thread-local retry queue to maintain precise chronological
|
||||||
|
ordering of log messages during transport connection failures.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
state (_ThreadState): Thread-safe tracking state containing connection params.
|
||||||
|
logger (logging.Logger): Diagnostic logger instance for transport anomalies.
|
||||||
|
"""
|
||||||
|
ctx = None
|
||||||
|
if common.CommType.TLS == state.comm_type:
|
||||||
|
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT)
|
||||||
|
ctx.check_hostname = False
|
||||||
|
ctx.verify_mode = ssl.VerifyMode.CERT_NONE
|
||||||
|
|
||||||
|
# Chronological staging queue dedicated to maintaining strict FIFO order during outages
|
||||||
|
retry_queue = collections.deque()
|
||||||
|
|
||||||
|
# Loop persists on shutdown until retry_queue is fully empty.
|
||||||
|
while retry_queue or not state.closed.is_set():
|
||||||
|
# connect to the endpoint
|
||||||
|
s = None
|
||||||
|
try:
|
||||||
|
factory = _SOCKET_FACTORIES[state.comm_type]
|
||||||
|
s = factory(state, ctx)
|
||||||
|
except KeyError:
|
||||||
|
raise NotImplementedError(f'Unsupported comm_type: {state.comm_type}')
|
||||||
|
except Exception:
|
||||||
|
if s is not None:
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
s.close()
|
||||||
|
time.sleep(state.reconnect_delay)
|
||||||
|
continue
|
||||||
|
|
||||||
|
# Connection successfully established; enter message transmission loop
|
||||||
|
# while True optimizes throughput and relies on socket exceptions to break the loop context
|
||||||
|
try:
|
||||||
|
captured_drops = ()
|
||||||
|
drop_payload = None
|
||||||
|
item = None
|
||||||
|
msg = None
|
||||||
|
msg_bytes = None
|
||||||
|
|
||||||
|
while True:
|
||||||
|
# Harvest and process dropped counts
|
||||||
|
# SAFELY EXTRACT AND PROCESS ALL TRACKED OVERFLOWS CHRONOLOGICALLY
|
||||||
|
try:
|
||||||
|
with state.cv:
|
||||||
|
if 1 < len(state.dropped) or 0 < state.dropped[0]:
|
||||||
|
# Freeze and copy the entire sequence array out of shared memory
|
||||||
|
captured_drops = tuple(state.dropped)
|
||||||
|
# Re-instantiate the tracked array to a clean initial state instantly
|
||||||
|
state.dropped.clear()
|
||||||
|
state.dropped.append(0)
|
||||||
|
|
||||||
|
# Convert captured thresholds into synthetic inline logs inside our local staging worker
|
||||||
|
for chunked_count in captured_drops:
|
||||||
|
if 0 < chunked_count:
|
||||||
|
drop_payload = _create_dropped_msg(
|
||||||
|
chunked_count, '_logging_handler_thread', 0,
|
||||||
|
)
|
||||||
|
# Appending to the retry queue guarantees strict chronological reporting order
|
||||||
|
retry_queue.append(RetryItem(synthetic=True, msg=drop_payload))
|
||||||
|
finally:
|
||||||
|
captured_drops = ()
|
||||||
|
drop_payload = None
|
||||||
|
|
||||||
|
# Grab from the main queue if the local transport staging queue is empty
|
||||||
|
# If the retry queue is empty, block and wait for a fresh log message
|
||||||
|
if not retry_queue:
|
||||||
|
try:
|
||||||
|
msg = state.queue.get(timeout=state.reconnect_delay)
|
||||||
|
except queue.Empty:
|
||||||
|
if state.closed.is_set():
|
||||||
|
# The queue is empty and the handler has been explicitly closed.
|
||||||
|
# Break the transmission loop to allow the worker thread to exit cleanly.
|
||||||
|
break
|
||||||
|
continue
|
||||||
|
else:
|
||||||
|
retry_queue.append(RetryItem(synthetic=False, msg=msg))
|
||||||
|
finally:
|
||||||
|
msg = None
|
||||||
|
|
||||||
|
# Transmit head message and track task lifecycle states
|
||||||
|
try:
|
||||||
|
# Peek at the oldest message without popping it yet
|
||||||
|
item = retry_queue[0]
|
||||||
|
msg_bytes = encoder.msg_to_str(item.msg).encode()
|
||||||
|
|
||||||
|
if common.CommType.UDP == state.comm_type:
|
||||||
|
s.send(msg_bytes)
|
||||||
|
else:
|
||||||
|
s.send(f'{len(msg_bytes)} '.encode() + msg_bytes)
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
# Message structure failure: Drop item immediately and preserve connection state
|
||||||
|
logger.exception('Dropping un-convertible poison-pill log message')
|
||||||
|
_item_completed(retry_queue, state.queue, item)
|
||||||
|
except UnicodeEncodeError:
|
||||||
|
# String binary processing failure: Drop item immediately and preserve connection state
|
||||||
|
logger.exception('Dropping un-encodable Unicode log message string')
|
||||||
|
_item_completed(retry_queue, state.queue, item)
|
||||||
|
except Exception:
|
||||||
|
# On socket break, tear down this loop context cleanly.
|
||||||
|
# The current message remains cleanly preserved at index 0 of retry_queue.
|
||||||
|
# Because task_done() is skipped, state.queue.join() will continue to block.
|
||||||
|
# Connection lost or infrastructure network failure: Break out loop to trigger socket reconnect
|
||||||
|
# Element stays safe at index 0 of the retry_queue cache.
|
||||||
|
# After reconnecting we will try it again.
|
||||||
|
break
|
||||||
|
else:
|
||||||
|
_item_completed(retry_queue, state.queue, item)
|
||||||
|
finally:
|
||||||
|
# Clear references to optimize memory tracking
|
||||||
|
item = None
|
||||||
|
msg_bytes = None
|
||||||
|
finally:
|
||||||
|
# close the connection to avoid leaking it
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
s.close()
|
||||||
|
|
||||||
|
|
||||||
|
class SyslogHandler(hat_syslog_handler_SyslogHandler):
|
||||||
|
"""
|
||||||
|
A process-safe wrapper for hat.syslog.handler.SyslogHandler.
|
||||||
|
|
||||||
|
Bypasses immutable NamedTuple state constraints on fork boundaries,
|
||||||
|
avoids thread-lock corruption from os.fork(), and short-circuits
|
||||||
|
the time-blocking flush/close loops during a Huey graceful shutdown.
|
||||||
|
|
||||||
|
A process-safe wrapper subclass for hat.syslog.handler.SyslogHandler.
|
||||||
|
|
||||||
|
Bypasses immutable state limitations on fork boundaries, handles thread-lock
|
||||||
|
corruption resulting from os.fork(), and avoids blocking task-execution loops
|
||||||
|
during a graceful worker process shutdown.
|
||||||
|
"""
|
||||||
|
def __init__(self, host, port, comm_type, queue_size=1024, reconnect_delay=5, *args, **kwargs):
|
||||||
|
"""Initializes the wrapper and neutralizes conflicting parent process states."""
|
||||||
|
super().__init__(host, port, common.CommType.UDP, queue_size, reconnect_delay, *args, **kwargs)
|
||||||
|
|
||||||
|
state = self._get_parent_attr('__state')
|
||||||
|
|
||||||
|
if state:
|
||||||
|
state.closed.set()
|
||||||
|
thread = self._get_parent_attr('__thread')
|
||||||
|
if thread and thread.is_alive():
|
||||||
|
with state.cv, contextlib.suppress(Exception):
|
||||||
|
state.cv.notify_all()
|
||||||
|
|
||||||
|
# Initialize native, thread-safe sync tracking structures
|
||||||
|
self.__state = _ThreadState(
|
||||||
|
host=host,
|
||||||
|
port=port,
|
||||||
|
comm_type=self._determine_comm_type(comm_type),
|
||||||
|
queue=queue.Queue(maxsize=queue_size),
|
||||||
|
queue_size=queue_size,
|
||||||
|
reconnect_delay=reconnect_delay,
|
||||||
|
cv=threading.Condition(),
|
||||||
|
closed=threading.Event(),
|
||||||
|
dropped=list((0,)),
|
||||||
|
)
|
||||||
|
|
||||||
|
self.__thread = None
|
||||||
|
self._initial_pid = os.getpid()
|
||||||
|
self._closing = threading.Event()
|
||||||
|
|
||||||
|
def _alive_thread(self):
|
||||||
|
"""Thread-safe validator confirming background worker availability."""
|
||||||
|
if not (self.__thread and self.__thread.is_alive()):
|
||||||
|
return False
|
||||||
|
|
||||||
|
with self.__state.cv:
|
||||||
|
if self.__thread and self.__thread.is_alive():
|
||||||
|
if hasattr(self.__thread, '_scoreboard') and self.__thread._scoreboard:
|
||||||
|
self.__thread._scoreboard.alive = (time.time(), time.monotonic_ns())
|
||||||
|
return True
|
||||||
|
return False
|
||||||
|
|
||||||
|
def _after_fork(self, current_pid=None):
|
||||||
|
"""
|
||||||
|
Intercepts the Unix process boundary skew. If a fork is identified,
|
||||||
|
it cleans old references and initializes fresh process-isolated primitives.
|
||||||
|
"""
|
||||||
|
|
||||||
|
# Detect if we crossed the Unix fork boundary into Huey's process worker.
|
||||||
|
if current_pid is None:
|
||||||
|
current_pid = os.getpid()
|
||||||
|
|
||||||
|
if current_pid == self._initial_pid:
|
||||||
|
# Return early when we have not forked.
|
||||||
|
return
|
||||||
|
|
||||||
|
self._initial_pid = current_pid
|
||||||
|
|
||||||
|
state = self.__state
|
||||||
|
new_state = _ThreadState(
|
||||||
|
host=state.host,
|
||||||
|
port=state.port,
|
||||||
|
comm_type=state.comm_type,
|
||||||
|
queue=queue.Queue(maxsize=state.queue_size),
|
||||||
|
queue_size=state.queue_size,
|
||||||
|
reconnect_delay=state.reconnect_delay,
|
||||||
|
cv=threading.Condition(),
|
||||||
|
closed=threading.Event(),
|
||||||
|
dropped=list((0,)),
|
||||||
|
)
|
||||||
|
self.__state = new_state
|
||||||
|
self.__thread = None
|
||||||
|
|
||||||
|
def _create_thread(self):
|
||||||
|
"""Aligns process states and builds a clean worker thread context."""
|
||||||
|
|
||||||
|
self._after_fork()
|
||||||
|
|
||||||
|
if self._closing.is_set() or self.__state.closed.is_set() or self._alive_thread():
|
||||||
|
return
|
||||||
|
|
||||||
|
initial = self.__thread is None
|
||||||
|
with self.__state.cv:
|
||||||
|
previous = self.__thread
|
||||||
|
|
||||||
|
self.__thread = threading.Thread(
|
||||||
|
target=_logging_handler_thread,
|
||||||
|
args=(self.__state,),
|
||||||
|
daemon=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
scoreboard = ThreadScoreboard()
|
||||||
|
self.__thread._scoreboard = scoreboard
|
||||||
|
|
||||||
|
if initial:
|
||||||
|
scoreboard.alive = scoreboard.start
|
||||||
|
scoreboard.initialized = scoreboard.start
|
||||||
|
elif previous and hasattr(previous, '_scoreboard') and previous._scoreboard:
|
||||||
|
scoreboard.alive = previous._scoreboard.alive
|
||||||
|
scoreboard.initialized = previous._scoreboard.initialized
|
||||||
|
scoreboard.previous_start = previous._scoreboard.start
|
||||||
|
logger.debug(f'Created a replacement thread: {scoreboard=}')
|
||||||
|
|
||||||
|
self.__thread.start()
|
||||||
|
|
||||||
|
def _determine_comm_type(self, comm_type):
|
||||||
|
"""Maps and cross-checks string connection descriptions to Enumeration definitions."""
|
||||||
|
if isinstance(comm_type, str):
|
||||||
|
needle = comm_type
|
||||||
|
haystack = frozenset(common.CommType.__members__)
|
||||||
|
vary = lambda x: {
|
||||||
|
x, x.upper(),
|
||||||
|
x.casefold(), x.casefold().upper(),
|
||||||
|
x.lower(), x.lower().upper(),
|
||||||
|
}
|
||||||
|
try:
|
||||||
|
matched_elements = tuple(haystack.intersection(vary(needle)))
|
||||||
|
member = matched_elements[0]
|
||||||
|
return common.CommType[member]
|
||||||
|
except (IndexError, KeyError) as e:
|
||||||
|
raise ValueError(f'Specify a valid comm_type from this list: {list(haystack)}') from e
|
||||||
|
|
||||||
|
if not isinstance(comm_type, common.CommType):
|
||||||
|
raise ValueError('Invalid comm_type argument')
|
||||||
|
|
||||||
|
def _parent_class_name(self):
|
||||||
|
return hat_syslog_handler_SyslogHandler.__name__
|
||||||
|
|
||||||
|
def _mangled_name(self, attr_name):
|
||||||
|
return f'_{self._parent_class_name()}{attr_name}'
|
||||||
|
|
||||||
|
def _get_parent_attr(self, attr_name):
|
||||||
|
"""Computes and gets mangled attributes from the super class."""
|
||||||
|
return getattr(self, self._mangled_name(attr_name), None)
|
||||||
|
|
||||||
|
def _set_parent_attr(self, attr_name, value):
|
||||||
|
"""Computes and sets mangled attributes on the super class."""
|
||||||
|
setattr(self, self._mangled_name(attr_name), value)
|
||||||
|
|
||||||
|
def emit(self, record):
|
||||||
|
"""Enqueues new log records and guarantees active connection coverage."""
|
||||||
|
if self._closing.is_set():
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
logger.handle(record)
|
||||||
|
return
|
||||||
|
|
||||||
|
self._create_thread()
|
||||||
|
|
||||||
|
if not self._alive_thread():
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
logger.handle(record)
|
||||||
|
return
|
||||||
|
|
||||||
|
state = self.__state
|
||||||
|
if state.closed.is_set():
|
||||||
|
self._closing.set()
|
||||||
|
logger.warning('Closed in emit')
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
logger.handle(record)
|
||||||
|
return
|
||||||
|
|
||||||
|
msg = _record_to_msg(record)
|
||||||
|
|
||||||
|
try:
|
||||||
|
state.queue.put_nowait(msg)
|
||||||
|
except queue.Full:
|
||||||
|
# ACQUIRE LOCK ON MAIN THREAD BEFORE INCREMENTING COUNTER SLICES
|
||||||
|
with state.cv:
|
||||||
|
dropped_count = state.dropped[-1]
|
||||||
|
if 1_000_000 < dropped_count:
|
||||||
|
state.dropped.append(1)
|
||||||
|
else:
|
||||||
|
state.dropped[-1] = 1 + dropped_count
|
||||||
|
|
||||||
|
logger.warning(f'Dropped a log message in emit due to buffer overflow: {msg.msg!r}')
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
logger.handle(record)
|
||||||
|
|
||||||
|
def flush(self):
|
||||||
|
"""Blocks execution until the internal logging queue is empty."""
|
||||||
|
self._create_thread()
|
||||||
|
|
||||||
|
if not self._alive_thread():
|
||||||
|
return
|
||||||
|
|
||||||
|
state = self.__state
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
state.queue.join()
|
||||||
|
|
||||||
|
def close(self):
|
||||||
|
"""
|
||||||
|
Gracefully flushes the queue and terminates the background logging thread.
|
||||||
|
Cleans up the queue, flags the state as closed, and shuts down
|
||||||
|
the background thread without allowing new ones to be generated.
|
||||||
|
"""
|
||||||
|
|
||||||
|
# Align/verify the background worker thread state immediately
|
||||||
|
self._closing.clear()
|
||||||
|
self._create_thread()
|
||||||
|
state = self.__state
|
||||||
|
if state.closed.is_set():
|
||||||
|
# Only return early when the thread is alive
|
||||||
|
state.closed.clear()
|
||||||
|
self._create_thread()
|
||||||
|
self._closing.set()
|
||||||
|
|
||||||
|
# The native queue.join() blockade is now perfectly synchronized with the internal
|
||||||
|
# tracking flow loop. It will block until retry_queue is 100% empty.
|
||||||
|
if self._alive_thread():
|
||||||
|
logger.debug('Flushing logging queue in close')
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
state.queue.join()
|
||||||
|
|
||||||
|
# Immediately trip the closed flag to end the networking thread
|
||||||
|
self.__state.closed.set()
|
||||||
|
|
||||||
|
# =====================================================================
|
||||||
|
# DYNAMIC GRANDPARENT BYPASS VIA MRO
|
||||||
|
# Instructs Python to search for close() starting *after* our direct
|
||||||
|
# parent class type descriptor. This dynamically resolves grandfather
|
||||||
|
# dependencies while completely avoiding the parent class thread-joins.
|
||||||
|
# =====================================================================
|
||||||
|
super(hat_syslog_handler_SyslogHandler, self).close()
|
||||||
|
|
||||||
|
|
||||||
|
handler = SyslogHandler
|
||||||
|
default_handler = handler(
|
||||||
|
comm_type='UDP',
|
||||||
|
host='127.0.0.1',
|
||||||
|
port=6514,
|
||||||
|
)
|
||||||
7
tubesync/common/logging/syslog/hat/_logger.py
Normal file
7
tubesync/common/logging/syslog/hat/_logger.py
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
import logging
|
||||||
|
|
||||||
|
logger = lambda name=None: logging.getLogger(
|
||||||
|
__name__.rsplit('.', 1)[0] if name is None else name
|
||||||
|
)
|
||||||
|
|
||||||
|
__all__ = ['logger']
|
||||||
3
tubesync/common/logging/syslog/std/__init__.py
Normal file
3
tubesync/common/logging/syslog/std/__init__.py
Normal file
@@ -0,0 +1,3 @@
|
|||||||
|
from ._default import * # noqa: F403
|
||||||
|
from ._logger import logger as logger
|
||||||
|
logger = logger(__name__)
|
||||||
19
tubesync/common/logging/syslog/std/_default.py
Normal file
19
tubesync/common/logging/syslog/std/_default.py
Normal file
@@ -0,0 +1,19 @@
|
|||||||
|
import logging
|
||||||
|
from logging.handlers import SysLogHandler
|
||||||
|
|
||||||
|
handler = SysLogHandler
|
||||||
|
facility = SysLogHandler.LOG_LOCAL0
|
||||||
|
|
||||||
|
default_formatter = logging.Formatter(
|
||||||
|
'%(asctime)s %(name)s: %(message)s',
|
||||||
|
'%b %d %H:%M:%S',
|
||||||
|
)
|
||||||
|
|
||||||
|
default_handler = handler(
|
||||||
|
address='/dev/log',
|
||||||
|
facility=facility,
|
||||||
|
)
|
||||||
|
default_handler.setFormatter(default_formatter)
|
||||||
|
default_handler.setLevel(logging.DEBUG)
|
||||||
|
|
||||||
|
__all__ = ['default_formatter', 'default_handler', 'handler']
|
||||||
7
tubesync/common/logging/syslog/std/_logger.py
Normal file
7
tubesync/common/logging/syslog/std/_logger.py
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
import logging
|
||||||
|
|
||||||
|
logger = lambda name=None: logging.getLogger(
|
||||||
|
__name__.rsplit('.', 1)[0] if name is None else name
|
||||||
|
)
|
||||||
|
|
||||||
|
__all__ = ['logger']
|
||||||
@@ -1,7 +1,7 @@
|
|||||||
from django import VERSION as DJANGO_VERSION
|
from django import VERSION as DJANGO_VERSION
|
||||||
from logging.handlers import SysLogHandler
|
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from common.huey import sqlite_tasks
|
from common.huey import sqlite_tasks
|
||||||
|
from common.logging import syslog
|
||||||
from common.utils import getenv
|
from common.utils import getenv
|
||||||
from sync.choices import TaskQueue
|
from sync.choices import TaskQueue
|
||||||
|
|
||||||
@@ -134,7 +134,7 @@ LOGGING = {
|
|||||||
},
|
},
|
||||||
'handlers': {
|
'handlers': {
|
||||||
'hat_syslog': {
|
'hat_syslog': {
|
||||||
'class': 'common.huey_syslog.SyslogHandler',
|
'class': syslog.hat.handler,
|
||||||
'host': '127.0.0.1',
|
'host': '127.0.0.1',
|
||||||
'port': 6514,
|
'port': 6514,
|
||||||
'comm_type': 'TCP',
|
'comm_type': 'TCP',
|
||||||
@@ -142,7 +142,7 @@ LOGGING = {
|
|||||||
'formatter': 'default',
|
'formatter': 'default',
|
||||||
},
|
},
|
||||||
'hat_syslog_worker_process': {
|
'hat_syslog_worker_process': {
|
||||||
'class': 'common.huey_syslog.SyslogHandler',
|
'class': syslog.hat.handler,
|
||||||
'host': '127.0.0.1',
|
'host': '127.0.0.1',
|
||||||
'port': 6514,
|
'port': 6514,
|
||||||
'comm_type': 'TCP',
|
'comm_type': 'TCP',
|
||||||
@@ -150,7 +150,7 @@ LOGGING = {
|
|||||||
'formatter': 'worker_process',
|
'formatter': 'worker_process',
|
||||||
},
|
},
|
||||||
'hat_syslog_worker_thread': {
|
'hat_syslog_worker_thread': {
|
||||||
'class': 'common.huey_syslog.SyslogHandler',
|
'class': syslog.hat.handler,
|
||||||
'host': '127.0.0.1',
|
'host': '127.0.0.1',
|
||||||
'port': 6514,
|
'port': 6514,
|
||||||
'comm_type': 'TCP',
|
'comm_type': 'TCP',
|
||||||
@@ -173,9 +173,9 @@ LOGGING = {
|
|||||||
'formatter': 'worker_thread',
|
'formatter': 'worker_thread',
|
||||||
},
|
},
|
||||||
'syslog': {
|
'syslog': {
|
||||||
'class': SysLogHandler,
|
'class': syslog.std.handler,
|
||||||
'address': '/dev/log',
|
'address': syslog.std.default_handler.address,
|
||||||
'facility': SysLogHandler.LOG_LOCAL0,
|
'facility': syslog.std.default_handler.facility,
|
||||||
'level': 'DEBUG',
|
'level': 'DEBUG',
|
||||||
'formatter': 'syslog',
|
'formatter': 'syslog',
|
||||||
},
|
},
|
||||||
@@ -185,7 +185,7 @@ LOGGING = {
|
|||||||
'level': 'DEBUG',
|
'level': 'DEBUG',
|
||||||
},
|
},
|
||||||
'loggers': {
|
'loggers': {
|
||||||
'common.huey_syslog': {
|
'common.logging.syslog.hat': {
|
||||||
'handlers': ['syslog'],
|
'handlers': ['syslog'],
|
||||||
'level': 'DEBUG',
|
'level': 'DEBUG',
|
||||||
'propagate': False,
|
'propagate': False,
|
||||||
|
|||||||
Reference in New Issue
Block a user