Fix ZMQ socket monitor leak and guard monitor recv against races
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <[email protected]>
This commit is contained in:
+8
-3
@@ -9,6 +9,8 @@ import struct
|
|||||||
import logging
|
import logging
|
||||||
import time
|
import time
|
||||||
|
|
||||||
|
from .constants import STATS_CONNECTION_DELAY
|
||||||
|
|
||||||
logger = logging.getLogger('network')
|
logger = logging.getLogger('network')
|
||||||
|
|
||||||
|
|
||||||
@@ -37,11 +39,11 @@ def check_monitor(monitor):
|
|||||||
"""Check monitor socket for events"""
|
"""Check monitor socket for events"""
|
||||||
try:
|
try:
|
||||||
event_monitor = monitor.recv(zmq.NOBLOCK)
|
event_monitor = monitor.recv(zmq.NOBLOCK)
|
||||||
|
event_id, event_name, event_value = read_socket_event(event_monitor)
|
||||||
|
event_endpoint = monitor.recv(zmq.NOBLOCK)
|
||||||
except zmq.Again:
|
except zmq.Again:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
event_id, event_name, event_value = read_socket_event(event_monitor)
|
|
||||||
event_endpoint = monitor.recv(zmq.NOBLOCK)
|
|
||||||
logger.debug(f'Monitor: {event_name} {event_value} endpoint {event_endpoint}')
|
logger.debug(f'Monitor: {event_name} {event_value} endpoint {event_endpoint}')
|
||||||
return (event_id, event_value)
|
return (event_id, event_value)
|
||||||
|
|
||||||
@@ -104,6 +106,9 @@ class RconConnection:
|
|||||||
|
|
||||||
def close(self):
|
def close(self):
|
||||||
"""Close connection"""
|
"""Close connection"""
|
||||||
|
if self.monitor:
|
||||||
|
self.monitor.setsockopt(zmq.LINGER, 0) # Don't wait for unsent messages
|
||||||
|
self.monitor.close()
|
||||||
if self.socket:
|
if self.socket:
|
||||||
self.socket.setsockopt(zmq.LINGER, 0) # Don't wait for unsent messages
|
self.socket.setsockopt(zmq.LINGER, 0) # Don't wait for unsent messages
|
||||||
self.socket.close()
|
self.socket.close()
|
||||||
@@ -143,7 +148,7 @@ class StatsConnection:
|
|||||||
logger.debug('Setting ZMQ_SUBSCRIBE to empty (all messages)')
|
logger.debug('Setting ZMQ_SUBSCRIBE to empty (all messages)')
|
||||||
self.socket.setsockopt(zmq.SUBSCRIBE, b'')
|
self.socket.setsockopt(zmq.SUBSCRIBE, b'')
|
||||||
|
|
||||||
time.sleep(0.5)
|
time.sleep(STATS_CONNECTION_DELAY)
|
||||||
self.connected = True
|
self.connected = True
|
||||||
logger.info('Stats stream connected')
|
logger.info('Stats stream connected')
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user