Skip to content
Open
136 changes: 93 additions & 43 deletions src/apify_client/_streamed_log.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,15 +32,24 @@ 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 reopening 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.
"""

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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

_received_data is never reset to False. StreamedLog and StreamedLogAsync fall back to the old behavior after stop() and a second start().

new_chunk = data
self._stream_buffer.append(new_chunk)
if re.findall(self._split_marker, new_chunk):
Expand Down Expand Up @@ -75,6 +84,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 +135,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 +144,7 @@ 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()
# 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 +153,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 +189,57 @@ 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.
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(self._empty_stream_retry_s)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is a fixed _empty_stream_retry_s enough, or should we use a backoff here?

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. A stream opened this late would
# end after its first chunk, so the whole log is read in one request instead.
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 @@ -244,17 +278,28 @@ def start(self) -> Task:
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 +314,25 @@ 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.
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(self._empty_stream_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