Skip to content

Commit da831de

Browse files
authored
[Feature] Improve CLI Realtime Logs (#100)
* Adding Execution Tree to the Realtime log with step jump feature * Refining the refresh feature of the CLI's Realtime log * Using now a per-client broadcast pattern for multiple simultaneous connections Also, dropped events during reconnection are now buffered and replayed instead of silently lost. * Fixing related unit tests for CLI
1 parent 6c23d87 commit da831de

7 files changed

Lines changed: 1363 additions & 820 deletions

File tree

tests/test_run/test_log_stream_handler.py

Lines changed: 29 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -46,10 +46,11 @@ def test_is_running_initially_false(self):
4646
h = LogStreamHandler()
4747
assert h.is_running is False
4848

49-
def test_log_queue_is_queue(self):
49+
def test_clients_is_empty_set(self):
5050
with patch("th_cli.test_run.log_stream_handler.LogsHTTPServer"):
5151
h = LogStreamHandler()
52-
assert isinstance(h.log_queue, queue.Queue)
52+
assert isinstance(h._clients, set)
53+
assert len(h._clients) == 0
5354

5455
def test_log_file_path_initially_none(self):
5556
with patch("th_cli.test_run.log_stream_handler.LogsHTTPServer"):
@@ -163,24 +164,23 @@ def test_sets_is_running_false(self):
163164

164165
assert h.is_running is False
165166

166-
def test_puts_none_sentinel_into_queue(self):
167+
def test_broadcasts_none_sentinel_on_stop(self):
167168
h, mock_srv = _make_handler()
168169
h.is_running = True
170+
client_q = queue.Queue()
171+
h._clients.add(client_q)
169172

170173
with patch("th_cli.test_run.log_stream_handler.logger"):
171174
h.stop()
172175

173-
assert h.log_queue.get_nowait() is None
176+
assert client_q.get_nowait() is None
174177

175-
def test_skips_sentinel_when_queue_full(self):
178+
def test_stop_does_not_raise_when_client_queue_full(self):
176179
h, mock_srv = _make_handler()
177180
h.is_running = True
178-
# Fill the queue to capacity
179-
for _ in range(h.log_queue.maxsize):
180-
try:
181-
h.log_queue.put_nowait("entry")
182-
except queue.Full:
183-
break
181+
client_q = queue.Queue(maxsize=1)
182+
client_q.put_nowait("existing") # fill to capacity
183+
h._clients.add(client_q)
184184

185185
with patch("th_cli.test_run.log_stream_handler.logger"):
186186
h.stop() # must not raise
@@ -196,45 +196,54 @@ class TestAddLogEntry:
196196
def test_noop_when_not_running(self):
197197
h, _ = _make_handler()
198198
h.is_running = False
199-
h.add_log_entry(message="ignored")
200-
assert h.log_queue.empty()
199+
with patch.object(h, "_broadcast") as mock_bcast:
200+
h.add_log_entry(message="ignored")
201+
mock_bcast.assert_not_called()
201202

202203
def test_adds_entry_to_queue_when_running(self):
203204
h, _ = _make_handler()
204205
h.is_running = True
206+
client_q = queue.Queue()
207+
h._clients.add(client_q)
205208
h.add_log_entry(message="hello", level="INFO")
206-
entry = h.log_queue.get_nowait()
209+
entry = client_q.get_nowait()
207210
assert entry["message"] == "hello"
208211
assert entry["level"] == "INFO"
209212

210213
def test_auto_generates_timestamp_when_not_provided(self):
211214
h, _ = _make_handler()
212215
h.is_running = True
216+
client_q = queue.Queue()
217+
h._clients.add(client_q)
213218
h.add_log_entry(message="msg")
214-
entry = h.log_queue.get_nowait()
219+
entry = client_q.get_nowait()
215220
assert "timestamp" in entry
216221
assert entry["timestamp"] is not None
217222

218223
def test_uses_provided_timestamp(self):
219224
h, _ = _make_handler()
220225
h.is_running = True
226+
client_q = queue.Queue()
227+
h._clients.add(client_q)
221228
h.add_log_entry(message="msg", timestamp="2025-01-01T00:00:00")
222-
entry = h.log_queue.get_nowait()
229+
entry = client_q.get_nowait()
223230
assert entry["timestamp"] == "2025-01-01T00:00:00"
224231

225232
def test_level_uppercased(self):
226233
h, _ = _make_handler()
227234
h.is_running = True
235+
client_q = queue.Queue()
236+
h._clients.add(client_q)
228237
h.add_log_entry(message="msg", level="warning")
229-
entry = h.log_queue.get_nowait()
238+
entry = client_q.get_nowait()
230239
assert entry["level"] == "WARNING"
231240

232241
def test_silently_drops_when_queue_full(self):
233242
h, _ = _make_handler()
234243
h.is_running = True
235-
# Fill queue to max
236-
for _ in range(h.log_queue.maxsize):
237-
h.log_queue.put_nowait({"message": "x"})
244+
client_q = queue.Queue(maxsize=1)
245+
client_q.put_nowait({"message": "x"}) # fill to capacity
246+
h._clients.add(client_q)
238247
h.add_log_entry(message="overflow") # must not raise
239248

240249

tests/test_run/test_logs_http_server.py

Lines changed: 63 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,12 @@ def _make_handler(path="/", server_attrs=None):
4545
handler.path = path
4646

4747
mock_server = MagicMock()
48+
# Set sensible defaults for the broadcast-pattern attributes so stream_logs()
49+
# doesn't hit MagicMock internals (e.g. deepcopy of threading.Lock).
50+
mock_server.active_clients = set()
51+
mock_server.clients_lock = threading.Lock()
52+
mock_server.tree_state = {} # empty dict avoids deepcopy branch
53+
mock_server.tree_lock = None
4854
for attr, value in (server_attrs or {}).items():
4955
setattr(mock_server, attr, value)
5056
handler.server = mock_server
@@ -64,6 +70,14 @@ def _make_handler(path="/", server_attrs=None):
6470
return handler
6571

6672

73+
def _pre_filled_queue(*items):
74+
"""Return a queue pre-loaded with items for use when patching queue.Queue."""
75+
q = queue.Queue()
76+
for item in items:
77+
q.put(item)
78+
return q
79+
80+
6781
# ---------------------------------------------------------------------------
6882
# do_GET routing
6983
# ---------------------------------------------------------------------------
@@ -202,49 +216,41 @@ def test_data_is_valid_json(self):
202216
@pytest.mark.unit
203217
class TestStreamLogs:
204218
def test_sends_sse_headers(self):
205-
q = queue.Queue()
206-
q.put(None) # immediate end-of-stream
207-
h = _make_handler(server_attrs={"log_queue": q})
208-
209-
with patch("th_cli.test_run.logs_http_server.logger"):
210-
h.stream_logs()
219+
h = _make_handler()
220+
with patch("th_cli.test_run.logs_http_server.queue.Queue", return_value=_pre_filled_queue(None)):
221+
with patch("th_cli.test_run.logs_http_server.logger"):
222+
h.stream_logs()
211223

212224
assert h._response_code == 200
213225
assert h._headers_sent.get("Content-Type") == "text/event-stream"
214226

215227
def test_sends_end_event_on_none_sentinel(self):
216-
q = queue.Queue()
217-
q.put(None)
218-
h = _make_handler(server_attrs={"log_queue": q})
228+
h = _make_handler()
219229
sent_events = []
220-
original_send = h._send_sse_event
221230
h._send_sse_event = lambda ev, data: sent_events.append(ev) or True
222231

223-
with patch("th_cli.test_run.logs_http_server.logger"):
224-
h.stream_logs()
232+
with patch("th_cli.test_run.logs_http_server.queue.Queue", return_value=_pre_filled_queue(None)):
233+
with patch("th_cli.test_run.logs_http_server.logger"):
234+
h.stream_logs()
225235

226236
assert "end" in sent_events
227237

228238
def test_streams_log_entries_from_queue(self):
229-
q = queue.Queue()
230-
q.put({"message": "hello", "level": "INFO", "timestamp": "2025-01-01"})
231-
q.put(None)
232-
h = _make_handler(server_attrs={"log_queue": q})
239+
h = _make_handler()
233240
sent_events = []
234241
h._send_sse_event = lambda ev, data: sent_events.append((ev, data)) or True
235242

236-
with patch("th_cli.test_run.logs_http_server.logger"):
237-
h.stream_logs()
243+
entry = {"message": "hello", "level": "INFO", "timestamp": "2025-01-01"}
244+
with patch("th_cli.test_run.logs_http_server.queue.Queue", return_value=_pre_filled_queue(entry, None)):
245+
with patch("th_cli.test_run.logs_http_server.logger"):
246+
h.stream_logs()
238247

239248
log_events = [d for ev, d in sent_events if ev == "log"]
240249
assert len(log_events) == 1
241250
assert log_events[0]["message"] == "hello"
242251

243252
def test_stops_when_send_sse_returns_false(self):
244-
q = queue.Queue()
245-
q.put({"message": "msg", "level": "INFO", "timestamp": "t"})
246-
q.put({"message": "msg2", "level": "INFO", "timestamp": "t"})
247-
h = _make_handler(server_attrs={"log_queue": q})
253+
h = _make_handler()
248254
call_count = [0]
249255

250256
def _send(ev, data):
@@ -254,33 +260,33 @@ def _send(ev, data):
254260
return False # disconnect immediately
255261

256262
h._send_sse_event = _send
263+
entry1 = {"message": "msg", "level": "INFO", "timestamp": "t"}
264+
entry2 = {"message": "msg2", "level": "INFO", "timestamp": "t"}
265+
with patch("th_cli.test_run.logs_http_server.queue.Queue", return_value=_pre_filled_queue(entry1, entry2)):
266+
with patch("th_cli.test_run.logs_http_server.logger"):
267+
h.stream_logs()
257268

258-
with patch("th_cli.test_run.logs_http_server.logger"):
259-
h.stream_logs()
260-
261-
# Should not have kept reading after disconnect
262269
assert call_count[0] <= 3
263270

264-
def test_no_log_queue_on_server_returns_early(self):
265-
h = _make_handler(server_attrs={"log_queue": None})
271+
def test_no_active_clients_on_server_returns_early(self):
272+
h = _make_handler(server_attrs={"active_clients": None, "clients_lock": None})
266273
with patch("th_cli.test_run.logs_http_server.logger"):
267274
h.stream_logs()
268-
# Should have sent 200 headers but not crashed
269275
assert h._response_code == 200
270276

271277
def test_debug_log_at_100_entries(self):
272-
"""Covers line 134: `if sent_count % 100 == 0` debug message."""
273-
q = queue.Queue()
274-
# Put 100 log entries then sentinel
275-
for i in range(100):
276-
q.put({"message": f"msg{i}", "level": "INFO", "timestamp": "t"})
277-
q.put(None)
278-
279-
h = _make_handler(server_attrs={"log_queue": q})
280-
with patch("th_cli.test_run.logs_http_server.logger") as mock_logger:
281-
h.stream_logs()
278+
"""Covers line `if sent_count % 100 == 0` debug message."""
279+
entries = [{"message": f"msg{i}", "level": "INFO", "timestamp": "t"} for i in range(100)]
280+
entries.append(None) # sentinel
281+
pre_q = queue.Queue()
282+
for e in entries:
283+
pre_q.put(e)
284+
285+
h = _make_handler()
286+
with patch("th_cli.test_run.logs_http_server.queue.Queue", return_value=pre_q):
287+
with patch("th_cli.test_run.logs_http_server.logger") as mock_logger:
288+
h.stream_logs()
282289

283-
# 100 entries should have triggered the debug log
284290
mock_logger.debug.assert_called()
285291

286292

@@ -381,7 +387,8 @@ def test_server_thread_initially_none(self):
381387
class TestLogsHTTPServerStart:
382388
def test_creates_threading_http_server(self):
383389
srv = LogsHTTPServer(port=0)
384-
q = queue.Queue()
390+
clients = set()
391+
lock = threading.Lock()
385392

386393
with patch("th_cli.test_run.logs_http_server.ThreadingHTTPServer") as mock_cls:
387394
mock_ths = MagicMock()
@@ -390,23 +397,32 @@ def test_creates_threading_http_server(self):
390397
mock_thread = MagicMock()
391398
mock_thread_cls.return_value = mock_thread
392399
with patch("th_cli.test_run.logs_http_server.logger"):
393-
srv.start(log_queue=q, test_run_title="Run")
400+
srv.start(active_clients=clients, clients_lock=lock, tree_state={}, test_run_title="Run")
394401

395402
mock_cls.assert_called_once()
396403
mock_thread.start.assert_called_once()
397404

398405
def test_sets_server_attributes(self):
399406
srv = LogsHTTPServer(port=0)
400-
q = queue.Queue()
407+
clients = set()
408+
lock = threading.Lock()
401409

402410
with patch("th_cli.test_run.logs_http_server.ThreadingHTTPServer") as mock_cls:
403411
mock_ths = MagicMock()
404412
mock_cls.return_value = mock_ths
405413
with patch("th_cli.test_run.logs_http_server.threading.Thread", return_value=MagicMock()):
406414
with patch("th_cli.test_run.logs_http_server.logger"):
407-
srv.start(log_queue=q, test_run_title="MyTitle", local_ip="1.2.3.4", log_file_path="/tmp/f.log")
408-
409-
assert mock_ths.log_queue is q
415+
srv.start(
416+
active_clients=clients,
417+
clients_lock=lock,
418+
tree_state={},
419+
test_run_title="MyTitle",
420+
local_ip="1.2.3.4",
421+
log_file_path="/tmp/f.log",
422+
)
423+
424+
assert mock_ths.active_clients is clients
425+
assert mock_ths.clients_lock is lock
410426
assert mock_ths.test_run_title == "MyTitle"
411427
assert mock_ths.local_ip == "1.2.3.4"
412428
assert mock_ths.log_file_path == "/tmp/f.log"
@@ -416,7 +432,7 @@ def test_propagates_oserror(self):
416432
with patch("th_cli.test_run.logs_http_server.ThreadingHTTPServer", side_effect=OSError("port in use")):
417433
with patch("th_cli.test_run.logs_http_server.logger"):
418434
with pytest.raises(OSError):
419-
srv.start(log_queue=queue.Queue(), test_run_title="run")
435+
srv.start(active_clients=set(), clients_lock=threading.Lock(), tree_state={}, test_run_title="run")
420436

421437

422438
# ---------------------------------------------------------------------------

0 commit comments

Comments
 (0)