From 29a9d417d6dbbf3ed70cd7a170eb9709579e161e Mon Sep 17 00:00:00 2001 From: Jarrod Johnson Date: Thu, 20 Aug 2026 07:23:40 -0400 Subject: [PATCH] Do not aggressively respawn buffer daemon. Start buffer daemon only when needed. Limit restarts to once every 30 seconds. --- confluent_server/confluent/consoleserver.py | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/confluent_server/confluent/consoleserver.py b/confluent_server/confluent/consoleserver.py index 92e3762f..d1db8ff3 100644 --- a/confluent_server/confluent/consoleserver.py +++ b/confluent_server/confluent/consoleserver.py @@ -61,6 +61,8 @@ def chunk_output(output, n): yield output[i:i + n] def get_buffer_output(nodename): + if _bufferdaemon is None: + eventlet.spawn(run_buffer_daemon) out = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) out.setsockopt(socket.SOL_SOCKET, socket.SO_PASSCRED, 1) out.connect("\x00confluent-vtbuffer") @@ -85,6 +87,8 @@ def get_buffer_output(nodename): def send_output(nodename, output): if not isinstance(nodename, bytes): nodename = nodename.encode('utf8') + if _bufferdaemon is None: + eventlet.spawn(run_buffer_daemon) out = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) out.setsockopt(socket.SOL_SOCKET, socket.SO_PASSCRED, 1) out.connect("\x00confluent-vtbuffer") @@ -602,16 +606,20 @@ running = True def run_buffer_daemon(): global _bufferdaemon while running: + minrestartdeadline = time.time() + 30 # Do not restart more than once every 30 seconds _bufferdaemon = subprocess.Popen( ['/opt/confluent/bin/vtbufferd', 'confluent-vtbuffer'], bufsize=0, stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL) _bufferdaemon.wait() + # Ensure we do not restart more than once every 30 seconds + sleep_time = minrestartdeadline - time.time() + if sleep_time > 0: + eventlet.sleep(sleep_time) def initialize(): global _tracelog global _bufferdaemon _tracelog = log.Logger('trace') - eventlet.spawn(run_buffer_daemon) def start_console_sessions(): configmodule.hook_new_configmanagers(_start_tenant_sessions)