Repository navigation
fix(log): redirect the run log when the first log stream comes back empty #1084
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
vdusek
wants to merge
7
commits into
master
Choose a base branch
from
fix/streamed-log-empty-start
base: master
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
+342
−45
Open
Changes from all commits
Commits
Show all changes
7 commits
Select commit
Hold shift + click to select a range
ca5cc31
fix(log): redirect the run log when the first log stream comes back e…
vdusek efb25ce
refactor(log): track the sync stop request with a single event
vdusek f0953be
test(log): check only the redirect logger's records in the empty stre…
vdusek 96aa922
test(log): cover a sync stop that lands while a stream reopens
vdusek 5ac8831
test(log): cover a missing log being neither reopened nor read
vdusek 644c8ad
test(log): cover an async stop that lands while a stream reopens
vdusek 88b51da
test(log): cover a failing one-shot log read on stop
vdusek File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
| new_chunk = data | ||
| self._stream_buffer.append(new_chunk) | ||
| if re.findall(self._split_marker, new_chunk): | ||
|
|
@@ -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.""" | ||
|
|
@@ -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. | ||
|
|
@@ -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() | ||
|
|
@@ -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: | ||
|
|
@@ -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) | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Is a fixed |
||
| 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') | ||
|
|
@@ -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.""" | ||
|
|
@@ -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. | ||
|
|
||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
_received_datais never reset to False.StreamedLogandStreamedLogAsyncfall back to the old behavior afterstop()and a secondstart().