11import os
2+ import socket
23import time as ttime
34
45import pytest
56import requests
67from bluesky_queueserver .manager .comms import zmq_single_request
7- from bluesky_queueserver .manager .tests .common import set_qserver_zmq_encoding # noqa: F401
8+ from bluesky_queueserver .manager .tests .common import ( # noqa: F401
9+ ReManager ,
10+ condition_manager_idle ,
11+ set_qserver_zmq_address ,
12+ set_qserver_zmq_encoding ,
13+ wait_for_condition ,
14+ )
815from xprocess import ProcessStarter
916
1017import bluesky_httpserver .server as bqss
1118
12- SERVER_ADDRESS = "localhost"
13- SERVER_PORT = "60610"
19+
20+ def _get_free_tcp_port () -> int :
21+ with socket .socket (socket .AF_INET , socket .SOCK_STREAM ) as sock :
22+ sock .bind (("" , 0 ))
23+ sock .listen (1 )
24+ return sock .getsockname ()[1 ]
25+
26+
27+ SERVER_BIND_HOST = os .getenv ("QSERVER_HTTP_TEST_BIND_HOST" , "127.0.0.1" )
28+ SERVER_ADDRESS = os .getenv ("QSERVER_HTTP_TEST_HOST" , "localhost" )
29+ SERVER_PORT = os .getenv ("QSERVER_HTTP_TEST_PORT" , str (_get_free_tcp_port ()))
30+
31+ QSERVER_TEST_REDIS_ADDR = os .getenv ("QSERVER_TEST_REDIS_ADDR" , "localhost" )
32+ QSERVER_TEST_ZMQ_CLIENT_HOST = os .getenv ("QSERVER_TEST_ZMQ_CLIENT_HOST" , "localhost" )
1433
1534# Single-user API key used for most of the tests
1635API_KEY_FOR_TESTS = "APIKEYFORTESTS"
@@ -41,7 +60,7 @@ class Starter(ProcessStarter):
4160 env ["QSERVER_HTTP_SERVER_SINGLE_USER_API_KEY" ] = API_KEY_FOR_TESTS
4261
4362 pattern = "Bluesky HTTP Server started successfully"
44- args = f"uvicorn --host={ SERVER_ADDRESS } --port { SERVER_PORT } { bqss .__name__ } :app" .split ()
63+ args = f"uvicorn --host={ SERVER_BIND_HOST } --port { SERVER_PORT } { bqss .__name__ } :app" .split ()
4564 # args = f"start-bluesky-httpserver --host={SERVER_ADDRESS} --port {SERVER_PORT}".split()
4665
4766 xprocess .ensure ("fastapi_server" , Starter )
@@ -60,7 +79,7 @@ def fastapi_server_fs(xprocess):
6079 to perform additional steps (such as setting environmental variables) before the server is started.
6180 """
6281
63- def start (http_server_host = SERVER_ADDRESS , http_server_port = SERVER_PORT , api_key = API_KEY_FOR_TESTS ):
82+ def start (http_server_host = SERVER_BIND_HOST , http_server_port = SERVER_PORT , api_key = API_KEY_FOR_TESTS ):
6483 class Starter (ProcessStarter ):
6584 max_read_lines = 53
6685
@@ -79,6 +98,100 @@ class Starter(ProcessStarter):
7998 xprocess .getinfo ("fastapi_server" ).terminate ()
8099
81100
101+ def _reset_queue_mode ():
102+ """Reset queue mode to default mode."""
103+ response , message = zmq_single_request ("queue_mode_set" , {"mode" : "default" })
104+ if response ["success" ] is not True :
105+ raise RuntimeError (message )
106+
107+
108+ def _build_re_manager_params (params = None ):
109+ params = list (params or [])
110+
111+ if not any (_ .startswith ("--zmq-control-addr" ) for _ in params ):
112+ control_port = _get_free_tcp_port ()
113+ params .append (f"--zmq-control-addr=tcp://*:{ control_port } " )
114+ else :
115+ control_port = int (next (_ .split (":" )[- 1 ] for _ in params if _ .startswith ("--zmq-control-addr" )).strip ())
116+
117+ if not any (_ .startswith ("--zmq-info-addr" ) for _ in params ):
118+ info_port = _get_free_tcp_port ()
119+ params .append (f"--zmq-info-addr=tcp://*:{ info_port } " )
120+
121+ if not any (_ .startswith ("--redis-addr" ) for _ in params ):
122+ params .append (f"--redis-addr={ QSERVER_TEST_REDIS_ADDR } " )
123+
124+ if not any (_ .startswith ("--redis-name-prefix" ) for _ in params ):
125+ params .append (f"--redis-name-prefix=qs_unit_tests_{ control_port } " )
126+
127+ client_zmq_address = f"tcp://{ QSERVER_TEST_ZMQ_CLIENT_HOST } :{ control_port } "
128+ return params , client_zmq_address
129+
130+
131+ @pytest .fixture
132+ def re_manager (monkeypatch ):
133+ """Start RE Manager on dynamically assigned ports for each test."""
134+ params , client_zmq_address = _build_re_manager_params ()
135+ set_qserver_zmq_address (monkeypatch , zmq_server_address = client_zmq_address )
136+
137+ manager = ReManager (params )
138+ failed_to_start = False
139+
140+ if not wait_for_condition (time = 10 , condition = condition_manager_idle ):
141+ failed_to_start = True
142+ manager .kill_manager ()
143+ raise TimeoutError ("Timeout: RE Manager failed to start." )
144+
145+ _reset_queue_mode ()
146+
147+ yield manager
148+
149+ if not failed_to_start :
150+ manager .stop_manager ()
151+ else :
152+ manager .kill_manager ()
153+
154+
155+ @pytest .fixture
156+ def re_manager_cmd (monkeypatch ):
157+ """Create RE Manager with dynamic ports and return handle factory."""
158+ state = {"manager" : None }
159+ failed_to_start = False
160+
161+ def _close ():
162+ manager = state ["manager" ]
163+ if manager is not None :
164+ if not failed_to_start :
165+ manager .stop_manager ()
166+ else :
167+ manager .kill_manager ()
168+ state ["manager" ] = None
169+
170+ def create_re_manager (params = None , * , stdout = None , stderr = None , set_redis_name_prefix = True ):
171+ nonlocal failed_to_start
172+ _close ()
173+ params_built , client_zmq_address = _build_re_manager_params (params = params )
174+ set_qserver_zmq_address (monkeypatch , zmq_server_address = client_zmq_address )
175+ state ["manager" ] = ReManager (
176+ params_built ,
177+ stdout = stdout ,
178+ stderr = stderr ,
179+ set_redis_name_prefix = set_redis_name_prefix ,
180+ )
181+
182+ if not wait_for_condition (time = 10 , condition = condition_manager_idle ):
183+ failed_to_start = True
184+ state ["manager" ].kill_manager ()
185+ raise TimeoutError ("Timeout: RE Manager failed to start." )
186+
187+ failed_to_start = False
188+ _reset_queue_mode ()
189+ return state ["manager" ]
190+
191+ yield create_re_manager
192+ _close ()
193+
194+
82195def setup_server_with_config_file (* , config_file_str , tmpdir , monkeypatch ):
83196 """
84197 Creates config file for the server in ``tmpdir/config/`` directory and
0 commit comments