Skip to content

Commit 000341a

Browse files
committed
the issue is dameon rep socket is now blocking, so we wait
1 parent 930549c commit 000341a

2 files changed

Lines changed: 16 additions & 16 deletions

File tree

src/cortex/discovery/daemon.py

Lines changed: 15 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@
1010

1111
import contextlib
1212
import os
13-
import threading
1413

1514
import zmq
1615

@@ -57,9 +56,8 @@ def __init__(
5756
self._context: zmq.Context | None = None
5857
self._socket: zmq.Socket | None = None
5958

60-
# Control flags
59+
# Control flag
6160
self._running = False
62-
self._shutdown_event = threading.Event()
6361

6462
def _ensure_ipc_path(self) -> None:
6563
"""Ensure the IPC socket directory exists."""
@@ -90,7 +88,6 @@ def start(self) -> None:
9088
self._socket.setsockopt(zmq.LINGER, 0) # Immediate shutdown
9189

9290
self._running = True
93-
self._shutdown_event.clear()
9491

9592
logger.info("=" * 50)
9693
logger.info("DISCOVERY DAEMON STARTED")
@@ -102,24 +99,27 @@ def start(self) -> None:
10299
except KeyboardInterrupt:
103100
logger.info("Received interrupt signal")
104101
finally:
105-
self.stop()
102+
self._cleanup()
106103

107104
def _run_loop(self) -> None:
108105
"""Main event loop."""
109-
while self._running and not self._shutdown_event.is_set():
106+
while self._running:
110107
try:
111-
# Try to receive a request
108+
# Try to receive a request (blocks up to RCVTIMEO)
112109
try:
113110
request_bytes = self._socket.recv(copy=False)
114111

115112
# Process the request
116113
response = self._handle_request(request_bytes)
117114
self._socket.send(response.to_bytes())
118115
except zmq.Again:
119-
# No message available, continue
116+
# Timeout, check _running and continue
120117
continue
121118

122119
except Exception as e:
120+
if not self._running:
121+
# We're shutting down, exit cleanly
122+
break
123123
logger.error(f"Error in discovery loop: {e}")
124124
# Send error response if we received a request
125125
try:
@@ -253,12 +253,8 @@ def _handle_shutdown(self) -> DiscoveryResponse:
253253
logger.info("SHUTDOWN command received")
254254
return DiscoveryResponse(status=DiscoveryStatus.OK, message="Shutting down")
255255

256-
def stop(self) -> None:
257-
"""Stop the discovery daemon."""
258-
logger.info("Stopping discovery daemon")
259-
self._shutdown_event.set()
260-
self._running = False
261-
256+
def _cleanup(self) -> None:
257+
"""Clean up resources."""
262258
if self._socket:
263259
try:
264260
self._socket.close()
@@ -284,6 +280,11 @@ def stop(self) -> None:
284280
logger.info("DISCOVERY DAEMON STOPPED")
285281
logger.info("=" * 50)
286282

283+
def stop(self) -> None:
284+
"""Stop the discovery daemon."""
285+
logger.info("Stopping discovery daemon")
286+
self._running = False
287+
287288

288289
def main():
289290
"""Entry point for the discovery daemon."""

tests/conftest.py

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -77,8 +77,7 @@ def run_daemon():
7777
def stop(self) -> None:
7878
"""Stop the discovery daemon."""
7979
if self._daemon:
80-
self._daemon._running = False
81-
self._daemon._shutdown_event.set()
80+
self._daemon.stop()
8281

8382
# Wait for thread to finish (daemon cleans up its own context)
8483
if self._thread:

0 commit comments

Comments
 (0)