Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
148 changes: 105 additions & 43 deletions src/apify_client/_streamed_log.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,15 +32,29 @@ class StreamedLogBase:
duration of the run (Impit currently maps it to an effective 24-hour cap) and mirrors the JS client.
"""

_empty_stream_retry_s: ClassVar[float] = 0.5
"""Pause before the first reopen of a log stream that ended before the run logged anything.

The API serves the log of a run that has not logged anything yet as an empty stream that ends at once. Each further
reopen doubles the pause up to `_empty_stream_max_retry_s`, which bounds the request rate while a run waits long to
start, for example for free memory.
"""

_empty_stream_max_retry_s: ClassVar[float] = 5
"""Upper bound on the pause between reopens of an empty log stream."""

def __init__(self, to_logger: logging.Logger, *, from_start: bool = True) -> None:
if self._force_propagate:
to_logger.propagate = True
self._to_logger = to_logger
self._stream_buffer = list[bytes]()
self._split_marker = re.compile(rb'(?:\n|^)(\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z)')
self._relevancy_time_limit: datetime | None = None if from_start else datetime.now(tz=UTC)
self._received_data = False

def _process_new_data(self, data: bytes) -> None:
if data:
self._received_data = True
Comment thread
vdusek marked this conversation as resolved.
new_chunk = data
self._stream_buffer.append(new_chunk)
if re.findall(self._split_marker, new_chunk):
Expand Down Expand Up @@ -75,6 +89,12 @@ def _log_buffer_content(self, *, include_last_part: bool = False) -> None:
message = decoded_marker + decoded_content
self._to_logger.log(level=self._guess_log_level_from_message(message), msg=message.strip())

def _process_whole_log(self, log: bytes | None) -> None:
"""Redirect a whole log read in one request, including its last part."""
if log:
self._process_new_data(log)
self._log_buffer_content(include_last_part=True)

@staticmethod
def _guess_log_level_from_message(message: str) -> int:
"""Guess the log level from the message."""
Expand Down Expand Up @@ -120,7 +140,7 @@ def __init__(self, log_client: LogClient, *, to_logger: logging.Logger, from_sta
self._log_client = log_client
self._streaming_thread: Thread | None = None
self._log_stream: HttpResponse | None = None
self._stop_logging = False
self._stop_event = threading.Event()

def start(self) -> Thread:
"""Start the streaming thread.
Expand All @@ -129,7 +149,8 @@ def start(self) -> Thread:
"""
if self._streaming_thread and self._streaming_thread.is_alive():
raise RuntimeError('Streaming thread already active')
self._stop_logging = False
self._stop_event.clear()
self._received_data = False
# A daemon thread so a stream still blocked on a read can never hold up interpreter shutdown.
self._streaming_thread = threading.Thread(target=self._stream_log, daemon=True)
self._streaming_thread.start()
Expand All @@ -138,13 +159,14 @@ def start(self) -> Thread:
def stop(self) -> None:
"""Signal the streaming thread to stop logging and wait up to `_stop_timeout_s` for it to finish.

A thread that outlives the wait is a daemon with `_stop_logging` set, so it exits after at most one more chunk,
A thread that outlives the wait is a daemon with `_stop_event` set, so it exits after at most one more chunk,
and only then does its buffered tail reach the logger. Its handle is kept while it is alive, so `start` cannot
revive it beside a second thread on the same buffer.
revive it beside a second thread on the same buffer. If no stream has delivered anything yet, the thread reads
the whole log in one request before it ends.
"""
if not self._streaming_thread:
raise RuntimeError('Streaming thread is not active')
self._stop_logging = True
self._stop_event.set()
# Read once; the streaming thread clears the attribute as soon as the stream ends.
log_stream = self._log_stream
if log_stream is not None:
Expand Down Expand Up @@ -173,39 +195,60 @@ def __exit__(

def _stream_log(self) -> None:
try:
with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
if not log_stream:
# An empty stream means the run has not logged anything yet, so reopen it until the first bytes arrive.
retry_s = self._empty_stream_retry_s
while not self._stop_event.is_set():
if not self._stream_log_once() or self._received_data:
return
# Published so `stop` can close the response.
self._log_stream = log_stream
try:
# `stop` may have run before the response existed for it to close.
if self._stop_logging:
return
for data in log_stream.iter_bytes():
self._process_new_data(data)
if self._stop_logging:
break
finally:
self._log_stream = None
try:
# Flush the last buffered part even if the read timed out or was stopped.
self._log_buffer_content(include_last_part=True)
except Exception:
# A truncated stream leaves an undecodable tail, which is worth a traceback even while a stop
# is in progress.
self._to_logger.exception('Log redirection stopped due to unexpected error:')
self._stop_event.wait(retry_s)
retry_s = min(retry_s * 2, self._empty_stream_max_retry_s)
except Exception as exc:
if self._stop_logging:
if self._stop_event.is_set():
# `stop` closed the stream out from under the read, so the failure is expected.
self._to_logger.debug('Log streaming stopped while `stop` was in progress: %r', exc)
return
if self._log_client._http_client.is_timeout_error(exc): # noqa: SLF001
elif self._log_client._http_client.is_timeout_error(exc): # noqa: SLF001
# The stream cannot continue, so warn and let the thread end instead of leaking a traceback.
self._to_logger.warning('Log streaming stopped: the log stream request timed out.')
return
else:
# Any other failure in log redirection must not escape the background thread; log it instead.
self._to_logger.exception('Log redirection stopped due to unexpected error:')
return
if self._received_data:
return
# Stopped before any stream delivered a byte, which a run that finishes quickly can cause.
try:
self._process_whole_log(self._log_client.get_as_bytes(raw=True))
except Exception:
self._to_logger.exception('Log redirection stopped due to unexpected error:')

def _stream_log_once(self) -> bool:
"""Redirect one log stream until it ends or `stop` is called. Return `False` when the log does not exist."""
with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
if not log_stream:
return False
# Published so `stop` can close the response.
self._log_stream = log_stream
try:
# `stop` may have run before the response existed for it to close. No earlier stream delivered
# anything, and one opened this late would end after its first chunk, so read the whole log in one
# request.
if self._stop_event.is_set():
return True
for data in log_stream.iter_bytes():
self._process_new_data(data)
if self._stop_event.is_set():
break
finally:
self._log_stream = None
try:
# Flush the last buffered part even if the read timed out or was stopped.
self._log_buffer_content(include_last_part=True)
except Exception:
# A truncated stream leaves an undecodable tail, which is worth a traceback even while a stop is
# in progress.
self._to_logger.exception('Log redirection stopped due to unexpected error:')
return True


@docs_group('Other')
Expand Down Expand Up @@ -240,21 +283,33 @@ def start(self) -> Task:
"""
if self._streaming_task and not self._streaming_task.done():
raise RuntimeError('Streaming task already active')
self._received_data = False
self._streaming_task = asyncio.create_task(self._stream_log())
return self._streaming_task

async def stop(self) -> None:
"""Stop the streaming task."""
"""Stop the streaming task.

If no stream has delivered anything yet, read the whole log in one request instead.
"""
if not self._streaming_task:
raise RuntimeError('Streaming task is not active')

was_streaming = not self._streaming_task.done()
self._streaming_task.cancel()
try:
await self._streaming_task
except asyncio.CancelledError:
pass
finally:
self._streaming_task = None
if not was_streaming or self._received_data:
return
# Stopped before any stream delivered a byte, which a run that finishes quickly can cause.
try:
self._process_whole_log(await self._log_client.get_as_bytes(raw=True))
except Exception:
self._to_logger.exception('Log redirection stopped due to unexpected error:')

async def __aenter__(self) -> Self:
"""Start the streaming task within the context. Exiting the context will cancel the streaming task."""
Expand All @@ -269,20 +324,27 @@ async def __aexit__(

async def _stream_log(self) -> None:
try:
async with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
if not log_stream:
return
try:
async for data in log_stream.aiter_bytes():
self._process_new_data(data)
finally:
# An empty stream means the run has not logged anything yet, so reopen it until the first bytes arrive.
retry_s = self._empty_stream_retry_s
while True:
async with self._log_client.stream(raw=True, timeout=self._stream_timeout) as log_stream:
if not log_stream:
return
try:
# Flush the last buffered part even if the task is cancelled by `stop()`.
self._log_buffer_content(include_last_part=True)
except Exception:
# A truncated stream leaves an undecodable tail. Keeping the failure here also keeps the
# cancellation `stop` raised propagating, so the task ends up cancelled as asyncio expects.
self._to_logger.exception('Log redirection stopped due to unexpected error:')
async for data in log_stream.aiter_bytes():
self._process_new_data(data)
finally:
try:
# Flush the last buffered part even if the task is cancelled by `stop()`.
self._log_buffer_content(include_last_part=True)
except Exception:
# A truncated stream leaves an undecodable tail. Keeping the failure here also keeps the
# cancellation `stop` raised propagating, so the task ends up cancelled as asyncio expects.
self._to_logger.exception('Log redirection stopped due to unexpected error:')
if self._received_data:
return
await asyncio.sleep(retry_s)
retry_s = min(retry_s * 2, self._empty_stream_max_retry_s)
except Exception as exc:
if self._log_client._http_client.is_timeout_error(exc): # noqa: SLF001
# A timeout on the long-lived stream is an expected terminal condition, not an error.
Expand Down
Loading
Loading