Skip to content

Commit ca5cc31

Browse files
committed
fix(log): redirect the run log when the first log stream comes back empty
1 parent ac60559 commit ca5cc31

2 files changed

Lines changed: 190 additions & 39 deletions

File tree

‎src/apify_client/_streamed_log.py‎

Lines changed: 91 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -32,15 +32,24 @@ class StreamedLogBase:
3232
duration of the run (Impit currently maps it to an effective 24-hour cap) and mirrors the JS client.
3333
"""
3434

35+
_empty_stream_retry_s: ClassVar[float] = 0.5
36+
"""Pause before reopening a log stream that ended before the run logged anything.
37+
38+
The API serves the log of a run that has not logged anything yet as an empty stream that ends at once.
39+
"""
40+
3541
def __init__(self, to_logger: logging.Logger, *, from_start: bool = True) -> None:
3642
if self._force_propagate:
3743
to_logger.propagate = True
3844
self._to_logger = to_logger
3945
self._stream_buffer = list[bytes]()
4046
self._split_marker = re.compile(rb'(?:\n|^)(\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z)')
4147
self._relevancy_time_limit: datetime | None = None if from_start else datetime.now(tz=UTC)
48+
self._received_data = False
4249

4350
def _process_new_data(self, data: bytes) -> None:
51+
if data:
52+
self._received_data = True
4453
new_chunk = data
4554
self._stream_buffer.append(new_chunk)
4655
if re.findall(self._split_marker, new_chunk):
@@ -75,6 +84,12 @@ def _log_buffer_content(self, *, include_last_part: bool = False) -> None:
7584
message = decoded_marker + decoded_content
7685
self._to_logger.log(level=self._guess_log_level_from_message(message), msg=message.strip())
7786

87+
def _process_whole_log(self, log: bytes | None) -> None:
88+
"""Redirect a log read in one request, for a stop that came before the stream delivered anything."""
89+
if log:
90+
self._process_new_data(log)
91+
self._log_buffer_content(include_last_part=True)
92+
7893
@staticmethod
7994
def _guess_log_level_from_message(message: str) -> int:
8095
"""Guess the log level from the message."""
@@ -121,6 +136,7 @@ def __init__(self, log_client: LogClient, *, to_logger: logging.Logger, from_sta
121136
self._streaming_thread: Thread | None = None
122137
self._log_stream: HttpResponse | None = None
123138
self._stop_logging = False
139+
self._stop_event = threading.Event()
124140

125141
def start(self) -> Thread:
126142
"""Start the streaming thread.
@@ -130,6 +146,7 @@ def start(self) -> Thread:
130146
if self._streaming_thread and self._streaming_thread.is_alive():
131147
raise RuntimeError('Streaming thread already active')
132148
self._stop_logging = False
149+
self._stop_event.clear()
133150
# A daemon thread so a stream still blocked on a read can never hold up interpreter shutdown.
134151
self._streaming_thread = threading.Thread(target=self._stream_log, daemon=True)
135152
self._streaming_thread.start()
@@ -140,11 +157,13 @@ def stop(self) -> None:
140157
141158
A thread that outlives the wait is a daemon with `_stop_logging` set, so it exits after at most one more chunk,
142159
and only then does its buffered tail reach the logger. Its handle is kept while it is alive, so `start` cannot
143-
revive it beside a second thread on the same buffer.
160+
revive it beside a second thread on the same buffer. If no stream has delivered anything yet, the thread reads
161+
the whole log in one request before it ends.
144162
"""
145163
if not self._streaming_thread:
146164
raise RuntimeError('Streaming thread is not active')
147165
self._stop_logging = True
166+
self._stop_event.set()
148167
# Read once; the streaming thread clears the attribute as soon as the stream ends.
149168
log_stream = self._log_stream
150169
if log_stream is not None:
@@ -173,39 +192,57 @@ def __exit__(
173192

174193
def _stream_log(self) -> None:
175194
try:
176-
with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
177-
if not log_stream:
195+
# An empty stream means the run has not logged anything yet, so reopen it until the first bytes arrive.
196+
while not self._stop_logging:
197+
if not self._stream_log_once() or self._received_data:
178198
return
179-
# Published so `stop` can close the response.
180-
self._log_stream = log_stream
181-
try:
182-
# `stop` may have run before the response existed for it to close.
183-
if self._stop_logging:
184-
return
185-
for data in log_stream.iter_bytes():
186-
self._process_new_data(data)
187-
if self._stop_logging:
188-
break
189-
finally:
190-
self._log_stream = None
191-
try:
192-
# Flush the last buffered part even if the read timed out or was stopped.
193-
self._log_buffer_content(include_last_part=True)
194-
except Exception:
195-
# A truncated stream leaves an undecodable tail, which is worth a traceback even while a stop
196-
# is in progress.
197-
self._to_logger.exception('Log redirection stopped due to unexpected error:')
199+
self._stop_event.wait(self._empty_stream_retry_s)
198200
except Exception as exc:
199201
if self._stop_logging:
200202
# `stop` closed the stream out from under the read, so the failure is expected.
201203
self._to_logger.debug('Log streaming stopped while `stop` was in progress: %r', exc)
202-
return
203-
if self._log_client._http_client.is_timeout_error(exc): # noqa: SLF001
204+
elif self._log_client._http_client.is_timeout_error(exc): # noqa: SLF001
204205
# The stream cannot continue, so warn and let the thread end instead of leaking a traceback.
205206
self._to_logger.warning('Log streaming stopped: the log stream request timed out.')
207+
return
206208
else:
207209
# Any other failure in log redirection must not escape the background thread; log it instead.
208210
self._to_logger.exception('Log redirection stopped due to unexpected error:')
211+
return
212+
if self._received_data:
213+
return
214+
# Stopped before any stream delivered a byte, which a run that finishes quickly can cause.
215+
try:
216+
self._process_whole_log(self._log_client.get_as_bytes(raw=True))
217+
except Exception:
218+
self._to_logger.exception('Log redirection stopped due to unexpected error:')
219+
220+
def _stream_log_once(self) -> bool:
221+
"""Redirect one log stream until it ends or `stop` is called. Return `False` when the log does not exist."""
222+
with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
223+
if not log_stream:
224+
return False
225+
# Published so `stop` can close the response.
226+
self._log_stream = log_stream
227+
try:
228+
# `stop` may have run before the response existed for it to close. A stream opened this late would
229+
# end after its first chunk, so the whole log is read in one request instead.
230+
if self._stop_logging:
231+
return True
232+
for data in log_stream.iter_bytes():
233+
self._process_new_data(data)
234+
if self._stop_logging:
235+
break
236+
finally:
237+
self._log_stream = None
238+
try:
239+
# Flush the last buffered part even if the read timed out or was stopped.
240+
self._log_buffer_content(include_last_part=True)
241+
except Exception:
242+
# A truncated stream leaves an undecodable tail, which is worth a traceback even while a stop is
243+
# in progress.
244+
self._to_logger.exception('Log redirection stopped due to unexpected error:')
245+
return True
209246

210247

211248
@docs_group('Other')
@@ -244,17 +281,28 @@ def start(self) -> Task:
244281
return self._streaming_task
245282

246283
async def stop(self) -> None:
247-
"""Stop the streaming task."""
284+
"""Stop the streaming task.
285+
286+
If no stream has delivered anything yet, read the whole log in one request instead.
287+
"""
248288
if not self._streaming_task:
249289
raise RuntimeError('Streaming task is not active')
250290

291+
was_streaming = not self._streaming_task.done()
251292
self._streaming_task.cancel()
252293
try:
253294
await self._streaming_task
254295
except asyncio.CancelledError:
255296
pass
256297
finally:
257298
self._streaming_task = None
299+
if not was_streaming or self._received_data:
300+
return
301+
# Stopped before any stream delivered a byte, which a run that finishes quickly can cause.
302+
try:
303+
self._process_whole_log(await self._log_client.get_as_bytes(raw=True))
304+
except Exception:
305+
self._to_logger.exception('Log redirection stopped due to unexpected error:')
258306

259307
async def __aenter__(self) -> Self:
260308
"""Start the streaming task within the context. Exiting the context will cancel the streaming task."""
@@ -269,20 +317,25 @@ async def __aexit__(
269317

270318
async def _stream_log(self) -> None:
271319
try:
272-
async with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
273-
if not log_stream:
274-
return
275-
try:
276-
async for data in log_stream.aiter_bytes():
277-
self._process_new_data(data)
278-
finally:
320+
# An empty stream means the run has not logged anything yet, so reopen it until the first bytes arrive.
321+
while True:
322+
async with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
323+
if not log_stream:
324+
return
279325
try:
280-
# Flush the last buffered part even if the task is cancelled by `stop()`.
281-
self._log_buffer_content(include_last_part=True)
282-
except Exception:
283-
# A truncated stream leaves an undecodable tail. Keeping the failure here also keeps the
284-
# cancellation `stop` raised propagating, so the task ends up cancelled as asyncio expects.
285-
self._to_logger.exception('Log redirection stopped due to unexpected error:')
326+
async for data in log_stream.aiter_bytes():
327+
self._process_new_data(data)
328+
finally:
329+
try:
330+
# Flush the last buffered part even if the task is cancelled by `stop()`.
331+
self._log_buffer_content(include_last_part=True)
332+
except Exception:
333+
# A truncated stream leaves an undecodable tail. Keeping the failure here also keeps the
334+
# cancellation `stop` raised propagating, so the task ends up cancelled as asyncio expects.
335+
self._to_logger.exception('Log redirection stopped due to unexpected error:')
336+
if self._received_data:
337+
return
338+
await asyncio.sleep(self._empty_stream_retry_s)
286339
except Exception as exc:
287340
if self._log_client._http_client.is_timeout_error(exc): # noqa: SLF001
288341
# A timeout on the long-lived stream is an expected terminal condition, not an error.

‎tests/unit/test_logging.py‎

Lines changed: 99 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
import itertools
55
import json
66
import logging
7+
import math
78
import threading
89
import time
910
from datetime import datetime, timedelta
@@ -20,7 +21,7 @@
2021
from apify_client._streamed_log import StreamedLog, StreamedLogAsync, StreamedLogBase
2122

2223
if TYPE_CHECKING:
23-
from collections.abc import Iterator
24+
from collections.abc import Callable, Iterator
2425

2526
from _pytest.logging import LogCaptureFixture
2627
from pytest_httpserver import HTTPServer
@@ -1393,6 +1394,103 @@ def test_streamed_log_sync_stop_reports_failing_stream_close(
13931394
streaming_thread.join(timeout=5)
13941395

13951396

1397+
def serve_log_after_empty_streams(httpserver: HTTPServer, *, empty_streams: float) -> list[Request]:
1398+
"""Serve the mocked log, but answer the first `empty_streams` stream requests with an empty body.
1399+
1400+
That is how the API answers a log stream request before the run has logged anything. Return the list the stream
1401+
requests are recorded in.
1402+
"""
1403+
stream_requests: list[Request] = []
1404+
1405+
def handler(request: Request) -> Response:
1406+
if 'stream' in request.args:
1407+
stream_requests.append(request)
1408+
if len(stream_requests) <= empty_streams:
1409+
return Response(b'', status=200, mimetype='application/octet-stream')
1410+
return Response(b''.join(_MOCKED_ACTOR_LOGS), status=200, mimetype='application/octet-stream')
1411+
1412+
httpserver.expect_request(f'/v2/actor-runs/{_MOCKED_RUN_ID}/log', method='GET').respond_with_handler(handler)
1413+
return stream_requests
1414+
1415+
1416+
def wait_until(condition: Callable[[], bool], *, timeout: float = 5) -> None:
1417+
"""Poll `condition` until it holds. Async tests call it through `asyncio.to_thread` to keep the event loop free."""
1418+
deadline = time.monotonic() + timeout
1419+
while not condition() and time.monotonic() < deadline:
1420+
time.sleep(0.01)
1421+
assert condition(), 'condition not met in time'
1422+
1423+
1424+
def test_streamed_log_sync_reopens_empty_stream(caplog: LogCaptureFixture, httpserver: HTTPServer) -> None:
1425+
"""A log stream that ends empty is reopened, and the reopened stream delivers the log."""
1426+
serve_log_after_empty_streams(httpserver, empty_streams=1)
1427+
logger = logging.getLogger('apify_client.tests.reopen_empty_stream_sync')
1428+
api_url = httpserver.url_for('/').removesuffix('/')
1429+
log_client = ApifyClient(token='mocked_token', api_url=api_url).run(run_id=_MOCKED_RUN_ID).log()
1430+
streamed_log = StreamedLog(log_client=log_client, to_logger=logger)
1431+
1432+
with caplog.at_level(logging.DEBUG, logger=logger.name):
1433+
streamed_log.start()
1434+
# Only the reopened stream can deliver the log before `stop` would read it in one request.
1435+
wait_until(lambda: len(caplog.records) == len(_EXPECTED_MESSAGES_AND_LEVELS))
1436+
streamed_log.stop()
1437+
1438+
assert [(record.message, record.levelno) for record in caplog.records] == list(_EXPECTED_MESSAGES_AND_LEVELS)
1439+
1440+
1441+
async def test_streamed_log_async_reopens_empty_stream(caplog: LogCaptureFixture, httpserver: HTTPServer) -> None:
1442+
"""A log stream that ends empty is reopened, and the reopened stream delivers the log."""
1443+
serve_log_after_empty_streams(httpserver, empty_streams=1)
1444+
logger = logging.getLogger('apify_client.tests.reopen_empty_stream_async')
1445+
api_url = httpserver.url_for('/').removesuffix('/')
1446+
log_client = ApifyClientAsync(token='mocked_token', api_url=api_url).run(run_id=_MOCKED_RUN_ID).log()
1447+
streamed_log = StreamedLogAsync(log_client=log_client, to_logger=logger)
1448+
1449+
with caplog.at_level(logging.DEBUG, logger=logger.name):
1450+
streamed_log.start()
1451+
# Only the reopened stream can deliver the log before `stop` would read it in one request.
1452+
await asyncio.to_thread(wait_until, lambda: len(caplog.records) == len(_EXPECTED_MESSAGES_AND_LEVELS))
1453+
await streamed_log.stop()
1454+
1455+
assert [(record.message, record.levelno) for record in caplog.records] == list(_EXPECTED_MESSAGES_AND_LEVELS)
1456+
1457+
1458+
def test_streamed_log_sync_stop_reads_log_when_streams_stay_empty(
1459+
caplog: LogCaptureFixture, httpserver: HTTPServer
1460+
) -> None:
1461+
"""When every log stream is empty until `stop`, the whole log is read in one request."""
1462+
stream_requests = serve_log_after_empty_streams(httpserver, empty_streams=math.inf)
1463+
logger = logging.getLogger('apify_client.tests.empty_streams_sync')
1464+
api_url = httpserver.url_for('/').removesuffix('/')
1465+
log_client = ApifyClient(token='mocked_token', api_url=api_url).run(run_id=_MOCKED_RUN_ID).log()
1466+
streamed_log = StreamedLog(log_client=log_client, to_logger=logger)
1467+
1468+
with caplog.at_level(logging.DEBUG, logger=logger.name):
1469+
streamed_log.start()
1470+
wait_until(lambda: len(stream_requests) >= 2)
1471+
streamed_log.stop()
1472+
1473+
assert [(record.message, record.levelno) for record in caplog.records] == list(_EXPECTED_MESSAGES_AND_LEVELS)
1474+
1475+
1476+
async def test_streamed_log_async_stop_reads_log_when_streams_stay_empty(
1477+
caplog: LogCaptureFixture, httpserver: HTTPServer
1478+
) -> None:
1479+
"""When every log stream is empty until `stop`, the whole log is read in one request."""
1480+
stream_requests = serve_log_after_empty_streams(httpserver, empty_streams=math.inf)
1481+
logger = logging.getLogger('apify_client.tests.empty_streams_async')
1482+
api_url = httpserver.url_for('/').removesuffix('/')
1483+
log_client = ApifyClientAsync(token='mocked_token', api_url=api_url).run(run_id=_MOCKED_RUN_ID).log()
1484+
streamed_log = StreamedLogAsync(log_client=log_client, to_logger=logger)
1485+
1486+
with caplog.at_level(logging.DEBUG, logger=logger.name):
1487+
streamed_log.start()
1488+
await asyncio.to_thread(wait_until, lambda: len(stream_requests) >= 2)
1489+
await streamed_log.stop()
1490+
1491+
assert [(record.message, record.levelno) for record in caplog.records] == list(_EXPECTED_MESSAGES_AND_LEVELS)
1492+
1493+
13961494
def test_logger_once_logs_the_first_call(caplog: LogCaptureFixture) -> None:
13971495
"""Test the first call with a given key is logged."""
13981496
logger = logging.getLogger('apify_client.tests.log_once_first')

0 commit comments

Comments
 (0)