From aba8e30f9859f1f498467058c55db04b2ae66b35 Mon Sep 17 00:00:00 2001 From: Ang Li Date: Tue, 6 Oct 2026 07:06:13 +0000 Subject: [PATCH 1/4] Speed up logcat_processor file scanning and polling loops Behavior-preserving efficiency pass over mobly/controllers/android_device_lib/logcat_processor.py: * `_iter_lines` now reads the file in binary mode and derives byte offsets by summing raw line lengths instead of calling text-mode `f.tell()` twice per line (an opaque cookie computation that dominated scan time). Lines are decoded as UTF-8 with `errors='replace'` exactly as before, and the resulting offsets are the same byte offsets `tail()` already computes, so positions from forward scans and `tail()` remain interchangeable. * Filter criteria are normalised once per query instead of once per line: a new internal `_LineFilter` compiles the pattern and builds the level sets a single time, and `_TimestampCutoff` pre-parses the `since` timestamp. `LogLine.matches()` keeps its public signature and now delegates to `_LineFilter`, so matching semantics are unchanged. * `LogcatListenerContext._listen_loop` and `wait_for()` keep a single file handle open for the lifetime of the listen/wait (via a new internal `_LineReader`) instead of re-opening the file every 50ms/100ms. The reader lazily opens a file that does not exist yet, re-tries after I/O errors, and is always closed through a context manager so handles are not held past the wait. Truncation behaves as before (reads stall until the file grows past the saved offset). In-order `wait_for` reuses the same reader across patterns. * `_parse_timestamp` uses precompiled split regexes. No public API changes. The old in-order timeout message semantics (reporting remaining budget) are kept as-is. Micro-benchmark (Python 3.13, 200k-line / 12.8 MB logcat file, best of 3, master's module loaded side-by-side via `git show origin/master:...`): _iter_lines (full scan) 2.394s -> 0.709s (3.4x) get_lines(level=[E,F], tag=[...]) 2.467s -> 0.738s (3.3x) get_lines(pattern, since=) 3.505s -> 1.150s (3.0x) tail(num_lines=50000, level=E) 0.860s -> 0.800s Adds tests/mobly/controllers/android_device_lib/logcat_processor_test.py covering: filter/cutoff equivalence against the public API, byte offsets with multi-byte characters, CRLF endings and invalid UTF-8, tail vs forward-scan offset consistency, reader handle reuse / lazy open, and listen/wait_for behaviour including files created after the wait starts. All behavioural tests in the new file also pass against the master implementation. --- .../android_device_lib/logcat_processor.py | 452 ++++++++++----- .../logcat_processor_test.py | 523 ++++++++++++++++++ 2 files changed, 848 insertions(+), 127 deletions(-) create mode 100644 tests/mobly/controllers/android_device_lib/logcat_processor_test.py diff --git a/mobly/controllers/android_device_lib/logcat_processor.py b/mobly/controllers/android_device_lib/logcat_processor.py index 7f46ad91..98df26fc 100644 --- a/mobly/controllers/android_device_lib/logcat_processor.py +++ b/mobly/controllers/android_device_lib/logcat_processor.py @@ -52,6 +52,16 @@ 'SILENT': 'S', } +# Splits a timestamp into its date and time halves ("MM-DD HH:MM:SS.mmm"). +_TIMESTAMP_SPLIT_RE = re.compile(r'[\sT]+') +# Splits the date half into its numeric elements ("2026-08-09" or "08/09"). +_DATE_SPLIT_RE = re.compile(r'[-/]') + +# Encoding used for all logcat file reads. Mirrors the text-mode arguments the +# logcat service uses when it opens the same file. +_ENCODING = 'utf-8' +_ENCODING_ERRORS = 'replace' + @dataclasses.dataclass(frozen=True) class LogcatPosition: @@ -87,8 +97,8 @@ def _parse_timestamp(t: str) -> tuple[int, int, int, int, int, int, int]: if not t: raise ValueError('Empty timestamp string') - date_part, time_part = re.split(r'[\sT]+', t.strip(), maxsplit=1) - date_elements = [int(x) for x in re.split(r'[-/]', date_part)] + date_part, time_part = _TIMESTAMP_SPLIT_RE.split(t.strip(), maxsplit=1) + date_elements = [int(x) for x in _DATE_SPLIT_RE.split(date_part)] if len(date_elements) == 3: year, month, day = date_elements elif len(date_elements) == 2: @@ -232,31 +242,7 @@ def matches( level: Optional[Union[str, Sequence[str], Set[str]]] = None, ) -> bool: """Checks if this log line matches the given pattern, tag, and/or level.""" - if pattern is not None: - regex = re.compile(pattern) if isinstance(pattern, str) else pattern - if not (regex.search(self.message) or regex.search(self.raw)): - return False - - if tag is not None: - if isinstance(tag, str): - if self.tag != tag: - return False - elif hasattr(tag, 'search'): - if not tag.search(self.tag): - return False - elif isinstance(tag, Iterable) and self.tag not in tag: - return False - - if level is not None: - levels = {level} if isinstance(level, str) else set(level) - norm_levels = { - _LEVEL_NORM_MAP.get(str(l).upper(), str(l).upper()) for l in levels - } - self_norm = _LEVEL_NORM_MAP.get(self.level.upper(), self.level.upper()) - if self.level not in levels and self_norm not in norm_levels: - return False - - return True + return _LineFilter(pattern=pattern, tag=tag, level=level).matches(self) @property def is_error(self) -> bool: @@ -284,6 +270,217 @@ def __ge__(self, other: Any) -> bool: return self.position >= other.position +class _LineFilter: + """Pre-normalised (pattern, tag, level) filter applied to many LogLines. + + Normalising the criteria once (compiling the regex, building the level sets) + and reusing the result is much cheaper than doing it per line, which matters + when scanning large logcat files. The matching semantics are identical to + :meth:`LogLine.matches`. + """ + + __slots__ = ('_regex', '_tag', '_tag_mode', '_levels', '_norm_levels') + + def __init__( + self, + pattern: Optional[Union[str, Pattern[str]]] = None, + tag: Optional[Union[str, Pattern[str], Sequence[str], Set[str]]] = None, + level: Optional[Union[str, Sequence[str], Set[str]]] = None, + ): + self._regex: Optional[Pattern[str]] = None + if pattern is not None: + self._regex = re.compile(pattern) if isinstance(pattern, str) else pattern + + # _tag_mode: None (no filter), 'eq', 'search', 'in' or 'noop' (an object + # that is neither a str, a regex nor an Iterable never filters anything). + self._tag_mode: Optional[str] = None + self._tag: Any = tag + if tag is not None: + if isinstance(tag, str): + self._tag_mode = 'eq' + elif hasattr(tag, 'search'): + self._tag_mode = 'search' + elif isinstance(tag, Iterable): + self._tag_mode = 'in' + try: + self._tag = frozenset(tag) + except TypeError: + self._tag = tuple(tag) + else: + self._tag_mode = 'noop' + + self._levels: Optional[frozenset[Any]] = None + self._norm_levels: Optional[frozenset[str]] = None + if level is not None: + levels = {level} if isinstance(level, str) else set(level) + self._levels = frozenset(levels) + self._norm_levels = frozenset( + _LEVEL_NORM_MAP.get(str(l).upper(), str(l).upper()) for l in levels + ) + + def matches(self, line: LogLine) -> bool: + """Returns True if the line satisfies every configured criterion.""" + regex = self._regex + if regex is not None and not ( + regex.search(line.message) or regex.search(line.raw) + ): + return False + + tag_mode = self._tag_mode + if tag_mode == 'eq': + if line.tag != self._tag: + return False + elif tag_mode == 'search': + if not self._tag.search(line.tag): + return False + elif tag_mode == 'in': + if line.tag not in self._tag: + return False + + if self._levels is not None: + self_norm = _LEVEL_NORM_MAP.get(line.level.upper(), line.level.upper()) + if line.level not in self._levels and self_norm not in self._norm_levels: + return False + + return True + + +class _TimestampCutoff: + """Pre-parsed lower timestamp bound used to skip lines older than `since`. + + ``is_before(ts)`` is equivalent to + ``LogcatPosition._compare_timestamps(ts, begin_time) < 0`` but parses + ``begin_time`` only once instead of once per scanned line. + """ + + __slots__ = ('_raw', '_parsed') + + def __init__(self, begin_time: str): + self._raw = str(begin_time) + try: + self._parsed: Optional[tuple[int, ...]] = LogcatPosition._parse_timestamp( + begin_time + ) + except (ValueError, IndexError): + self._parsed = None + + def is_before(self, timestamp: Optional[str]) -> bool: + """Returns True if `timestamp` is chronologically before the cutoff.""" + if not timestamp: + # _compare_timestamps(falsy, truthy) == -1. + return True + if self._parsed is not None: + try: + p1 = LogcatPosition._parse_timestamp(timestamp) + except (ValueError, IndexError): + pass + else: + p2 = self._parsed + if p1[0] == 0 or p2[0] == 0: + p1 = (0,) + p1[1:] + p2 = (0,) + p2[1:] + return p1 < p2 + return str(timestamp) < self._raw + + +class _LineReader: + """Incrementally reads parsed LogLines from a (possibly growing) file. + + The file is opened lazily in binary mode and the handle is kept open across + successive :meth:`read_lines` calls, so that polling loops do not pay for an + ``open()`` on every iteration. Byte offsets are tracked by summing the length + of each raw line, which is both exact and far cheaper than text-mode + ``tell()``; the offsets are therefore directly comparable with the ones + computed by :meth:`LogcatProcessor.tail`. + + Always use as a context manager (or call :meth:`close`) so the handle is not + held longer than necessary. + """ + + __slots__ = ('_file_path', 'offset', '_file') + + def __init__(self, file_path: str, offset: int = 0): + self._file_path = file_path + self.offset = offset + self._file: Optional[Any] = None + + def __enter__(self) -> '_LineReader': + return self + + def __exit__(self, exc_type, exc_val, exc_tb) -> None: + self.close() + + @property + def is_open(self) -> bool: + return self._file is not None + + def close(self) -> None: + """Closes the underlying file handle, if any.""" + f, self._file = self._file, None + if f is not None: + try: + f.close() + except OSError: + pass + + def _ensure_open(self) -> bool: + if self._file is not None: + return True + if not os.path.exists(self._file_path): + return False + try: + f = open(self._file_path, 'rb') + except OSError: + return False + try: + if self.offset > 0: + f.seek(self.offset) + except OSError: + f.close() + return False + self._file = f + return True + + def read_lines(self) -> Iterator[tuple[int, LogLine]]: + """Yields (offset_after_line, LogLine) for every new parseable line. + + Reading stops at the current end of file; calling this again later picks up + data appended in the meantime. On an I/O error the handle is closed and the + iteration ends; the next call will try to re-open the file. + """ + if not self._ensure_open(): + return + f = self._file + offset = self.offset + try: + while True: + raw = f.readline() + if not raw: + break + line_offset = offset + offset += len(raw) + self.offset = offset + parsed = LogLine.from_string( + raw.decode(_ENCODING, _ENCODING_ERRORS), byte_offset=line_offset + ) + if parsed is not None: + yield offset, parsed + except OSError: + self.close() + return + + +def _resolve_since( + since: Optional[Union[LogcatPosition, LogLine]], +) -> tuple[int, Optional[_TimestampCutoff]]: + """Converts a `since` argument into (byte_offset, timestamp cutoff).""" + pos = since.position if isinstance(since, LogLine) else since + offset = pos._byte_offset if pos else 0 + begin_time = pos.timestamp if pos and offset == 0 else None + cutoff = _TimestampCutoff(begin_time) if begin_time else None + return offset, cutoff + + class LogcatListenerContext: """Context manager for listening to real-time logcat events.""" @@ -301,6 +498,7 @@ def __init__( self._pattern = pattern self._tag = tag self._level = level + self._filter = _LineFilter(pattern=pattern, tag=tag, level=level) self._position = ( position.position if isinstance(position, LogLine) else position ) @@ -337,7 +535,7 @@ def get_next_event(self, timeout: Optional[float] = None) -> LogLine: ) def _dispatch(self, line: LogLine) -> None: - if line.matches(pattern=self._pattern, tag=self._tag, level=self._level): + if self._filter.matches(line): with self._lock: self._events.append(line) try: @@ -346,18 +544,21 @@ def _dispatch(self, line: LogLine) -> None: pass def _listen_loop(self) -> None: - current_offset = ( + start_offset = ( self._position._byte_offset if self._position else LogcatPosition.from_file(self._processor.file_path)._byte_offset ) - while not self._stop_event.is_set(): - for offset, line in self._processor._iter_lines(offset=current_offset): - current_offset = offset - self._dispatch(line) - if self._stop_event.is_set(): - break - time.sleep(0.05) + stop_event = self._stop_event + # A single reader (and file handle) is reused for the whole listen session + # instead of re-opening the file on every poll. + with _LineReader(self._processor.file_path, start_offset) as reader: + while not stop_event.is_set(): + for _, line in reader.read_lines(): + self._dispatch(line) + if stop_event.is_set(): + break + stop_event.wait(0.05) def __enter__(self) -> 'LogcatListenerContext': self._stop_event.clear() @@ -388,26 +589,14 @@ def file_path(self) -> str: return self._file_path def _iter_lines(self, offset: int = 0) -> Iterator[tuple[int, LogLine]]: - """Yields (line_offset, LogLine) pairs from file from given offset.""" - if not os.path.exists(self._file_path): - return - try: - with open( - self._file_path, 'r', encoding='utf-8', errors='replace', newline='' - ) as f: - if offset > 0: - f.seek(offset) - while True: - line_offset = f.tell() - line = f.readline() - if not line: - break - current_offset = f.tell() - parsed = LogLine.from_string(line, byte_offset=line_offset) - if parsed is not None: - yield current_offset, parsed - except OSError: - return + """Yields (offset_after_line, LogLine) pairs from file from given offset. + + The file is read in binary mode and byte offsets are derived from the raw + line lengths, so they match the offsets produced by :meth:`tail`. Lines are + decoded as UTF-8 with replacement of undecodable bytes. + """ + with _LineReader(self._file_path, offset) as reader: + yield from reader.read_lines() def get_lines( self, @@ -432,19 +621,14 @@ def get_lines( ' tail() instead.' ) - pos = since.position if isinstance(since, LogLine) else since - offset = pos._byte_offset if pos else 0 - begin_time = pos.timestamp if pos and offset == 0 else None + offset, cutoff = _resolve_since(since) + line_filter = _LineFilter(pattern=pattern, tag=tag, level=level) results: list[LogLine] = [] for _, parsed in self._iter_lines(offset=offset): - if ( - begin_time - and LogcatPosition._compare_timestamps(parsed.timestamp, begin_time) - < 0 - ): + if cutoff is not None and cutoff.is_before(parsed.timestamp): continue - if parsed.matches(pattern=pattern, tag=tag, level=level): + if line_filter.matches(parsed): results.append(parsed) if max_lines is not None and len(results) >= max_lines: break @@ -463,6 +647,7 @@ def tail( buf: collections.deque[LogLine] = collections.deque() block_size = 64 * 1024 # 64KB chunks + line_filter = _LineFilter(pattern=pattern, tag=tag, level=level) try: with open(self._file_path, 'rb') as f: @@ -473,7 +658,6 @@ def tail( remaining = file_size remainder = b'' - lines_to_process: list[tuple[int, bytes]] = [] while remaining > 0 and len(buf) < num_lines: read_size = min(block_size, remaining) @@ -494,17 +678,17 @@ def tail( current_offset = 0 # Calculate offsets and parse lines in reverse order within this block - block_lines: list[tuple[int, LogLine]] = [] + block_lines: list[LogLine] = [] for line_bytes in lines_chunk: line_offset = current_offset current_offset += len(line_bytes) + 1 # count \n byte - line_str = line_bytes.decode('utf-8', errors='replace') + line_str = line_bytes.decode(_ENCODING, _ENCODING_ERRORS) parsed = LogLine.from_string(line_str, byte_offset=line_offset) if parsed is not None: - block_lines.append((line_offset, parsed)) + block_lines.append(parsed) - for _, parsed in reversed(block_lines): - if parsed.matches(pattern=pattern, tag=tag, level=level): + for parsed in reversed(block_lines): + if line_filter.matches(parsed): buf.appendleft(parsed) if len(buf) >= num_lines: break @@ -542,82 +726,96 @@ def wait_for( return [] deadline = time.perf_counter() + timeout_sec + offset, cutoff = _resolve_since(since) if in_order: matched_lines: list[LogLine] = [] - current_since = since - for pat in patterns: - remaining = deadline - time.perf_counter() - if remaining <= 0: - raise self._timeout_error_cls( - f'Timed out after {timeout_sec}s waiting for in-order pattern:' - f' {pat!r}' + with _LineReader(self._file_path, offset) as reader: + for pat in patterns: + remaining = deadline - time.perf_counter() + if remaining <= 0: + raise self._timeout_error_cls( + f'Timed out after {timeout_sec}s waiting for in-order pattern:' + f' {pat!r}' + ) + matched_lines.append( + self._wait_on_reader( + reader, + _LineFilter(pattern=pat), + cutoff, + deadline, + remaining, + pat, + ) ) - matched, next_offset = self._wait_for_single( - pattern=pat, - timeout_sec=remaining, - since=current_since, - ) - matched_lines.append(matched) - current_since = LogcatPosition(_byte_offset=next_offset) + # Subsequent patterns continue right after the matched line; the + # timestamp cutoff only bounds the initial scan. + cutoff = None return matched_lines - unmatched = list(enumerate(patterns)) + unmatched: list[tuple[int, Union[str, Pattern[str]], _LineFilter]] = [ + (idx, pat, _LineFilter(pattern=pat)) for idx, pat in enumerate(patterns) + ] matched_dict: dict[int, LogLine] = {} - pos = since.position if isinstance(since, LogLine) else since - offset = pos._byte_offset if pos else 0 - begin_time = pos.timestamp if pos and offset == 0 else None - scan_offset = offset - while time.perf_counter() < deadline: - for current_offset, parsed in self._iter_lines(offset=scan_offset): - scan_offset = current_offset - if ( - begin_time - and LogcatPosition._compare_timestamps(parsed.timestamp, begin_time) - < 0 - ): - continue - for idx, pat in list(unmatched): - if parsed.matches(pattern=pat): - matched_dict[idx] = parsed - unmatched.remove((idx, pat)) - if not unmatched: - return [matched_dict[i] for i in range(len(patterns))] - time.sleep(0.1) - - remaining_patterns = [pat for _, pat in unmatched] + with _LineReader(self._file_path, offset) as reader: + while time.perf_counter() < deadline: + for _, parsed in reader.read_lines(): + if cutoff is not None and cutoff.is_before(parsed.timestamp): + continue + for entry in list(unmatched): + if entry[2].matches(parsed): + matched_dict[entry[0]] = parsed + unmatched.remove(entry) + if not unmatched: + return [matched_dict[i] for i in range(len(patterns))] + time.sleep(0.1) + + remaining_patterns = [pat for _, pat, _ in unmatched] raise self._timeout_error_cls( f'Timed out after {timeout_sec}s waiting for patterns:' f' {remaining_patterns!r}' ) - def _wait_for_single( + def _wait_on_reader( self, + reader: _LineReader, + line_filter: _LineFilter, + cutoff: Optional[_TimestampCutoff], + deadline: float, + timeout_sec: float, pattern: Union[str, Pattern[str]], - timeout_sec: float = 60.0, - since: Optional[Union[LogcatPosition, LogLine]] = None, - ) -> tuple[LogLine, int]: - deadline = time.perf_counter() + timeout_sec - pos = since.position if isinstance(since, LogLine) else since - offset = pos._byte_offset if pos else 0 - begin_time = pos.timestamp if pos and offset == 0 else None - scan_offset = offset - + ) -> LogLine: + """Polls `reader` until a line matches `line_filter` or `deadline` passes.""" while time.perf_counter() < deadline: - for current_offset, parsed in self._iter_lines(offset=scan_offset): - scan_offset = current_offset - if ( - begin_time - and LogcatPosition._compare_timestamps(parsed.timestamp, begin_time) - < 0 - ): + for _, parsed in reader.read_lines(): + if cutoff is not None and cutoff.is_before(parsed.timestamp): continue - if parsed.matches(pattern=pattern): - return parsed, current_offset + if line_filter.matches(parsed): + return parsed time.sleep(0.1) raise self._timeout_error_cls( f'Timed out after {timeout_sec}s waiting for logcat pattern:' f' {pattern!r}' ) + + def _wait_for_single( + self, + pattern: Union[str, Pattern[str]], + timeout_sec: float = 60.0, + since: Optional[Union[LogcatPosition, LogLine]] = None, + ) -> tuple[LogLine, int]: + """Waits for a single pattern; returns (line, offset after that line).""" + deadline = time.perf_counter() + timeout_sec + offset, cutoff = _resolve_since(since) + with _LineReader(self._file_path, offset) as reader: + matched = self._wait_on_reader( + reader, + _LineFilter(pattern=pattern), + cutoff, + deadline, + timeout_sec, + pattern, + ) + return matched, reader.offset diff --git a/tests/mobly/controllers/android_device_lib/logcat_processor_test.py b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py new file mode 100644 index 00000000..dc6ade75 --- /dev/null +++ b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py @@ -0,0 +1,523 @@ +# Copyright 2026 Google Inc. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +"""Unit tests for mobly.controllers.android_device_lib.logcat_processor.""" + +import os +import re +import shutil +import tempfile +import time +import unittest + +from mobly.controllers.android_device_lib import logcat_processor + +LogLine = logcat_processor.LogLine +LogcatPosition = logcat_processor.LogcatPosition +LogcatProcessor = logcat_processor.LogcatProcessor + +SAMPLE_LINES = [ + '--------- beginning of system', + '08-09 22:00:00.100 1000 1010 I SystemServer: Entered main', + '2026-08-09 22:00:01.200 1000 1020 D WifiService: Enabling wlan0', + '08-09 22:00:02.300 2050 2050 I ExampleApp: Initialised \u2713 \u00e9t\u00e9', + '08-09 22:00:03.400 1000 1040 W BtGatt: Retry 1 for AA:BB', + '\tat com.example.app.NetworkClient.connect(NetworkClient.java:42)', + '08-09 22:00:05.150 2050 2060 E ExampleApp: Failed to connect', + '08-09 22:00:05.800 1000 1040 F BtGatt: Fatal controller error', +] + + +def _make_line(timestamp='08-09 22:00:00.000', level='I', tag='T', msg='m'): + return LogLine.from_string(f'{timestamp} 1000 1010 {level} {tag}: {msg}') + + +class LogLineTest(unittest.TestCase): + + def test_from_string_parses_fields(self): + line = LogLine.from_string( + '2026-08-09 22:00:01.200 1000 1020 D WifiService: Enabling wlan0\r\n', + byte_offset=42, + ) + self.assertEqual(line.timestamp, '2026-08-09 22:00:01.200') + self.assertEqual(line.pid, 1000) + self.assertEqual(line.tid, 1020) + self.assertEqual(line.level, 'D') + self.assertEqual(line.tag, 'WifiService') + self.assertEqual(line.message, 'Enabling wlan0') + self.assertEqual(line.raw.endswith('wlan0'), True) + self.assertEqual(line.position._byte_offset, 42) + + def test_from_string_rejects_non_logcat_lines(self): + self.assertIsNone(LogLine.from_string('')) + self.assertIsNone(LogLine.from_string('--------- beginning of main')) + self.assertIsNone(LogLine.from_string(None)) + + def test_matches_pattern_str_and_compiled(self): + line = _make_line(msg='DHCP OFFER received from 192.168.1.1') + self.assertTrue(line.matches(pattern=r'OFFER.*192\.168')) + self.assertTrue(line.matches(pattern=re.compile(r'^DHCP'))) + self.assertFalse(line.matches(pattern='DISCOVER')) + # The pattern is also searched against the full raw line. + self.assertTrue(line.matches(pattern=r'1000\s+1010')) + + def test_matches_tag_variants(self): + line = _make_line(tag='WifiService') + self.assertTrue(line.matches(tag='WifiService')) + self.assertFalse(line.matches(tag='Wifi')) + self.assertTrue(line.matches(tag=re.compile('^Wifi'))) + self.assertTrue(line.matches(tag=['BtGatt', 'WifiService'])) + self.assertTrue(line.matches(tag={'WifiService'})) + self.assertFalse(line.matches(tag=('BtGatt',))) + # A tag object that is neither str, regex nor iterable never filters. + self.assertTrue(line.matches(tag=object())) + + def test_matches_level_normalisation(self): + line = _make_line(level='E') + self.assertTrue(line.matches(level='E')) + self.assertTrue(line.matches(level='error')) + self.assertTrue(line.matches(level=['W', 'ERROR'])) + self.assertTrue(line.matches(level={'e'})) + self.assertFalse(line.matches(level='W')) + self.assertFalse(line.matches(level=['V', 'D', 'I'])) + fatal = _make_line(level='F') + self.assertTrue(fatal.matches(level='A')) + self.assertTrue(fatal.matches(level='ASSERT')) + + def test_matches_combined_criteria(self): + line = _make_line(level='E', tag='ExampleApp', msg='Failed to connect') + self.assertTrue(line.matches(pattern='Failed', tag='ExampleApp', level='E')) + self.assertFalse( + line.matches(pattern='Failed', tag='ExampleApp', level='I') + ) + self.assertFalse(line.matches(pattern='Failed', tag='Other', level='E')) + self.assertFalse(line.matches(pattern='Nope', tag='ExampleApp', level='E')) + + def test_line_filter_is_equivalent_to_matches(self): + lines = [ + _make_line(level=lvl, tag=tag, msg=msg) + for lvl in 'VDIWEF' + for tag in ('WifiService', 'BtGatt', 'ExampleApp') + for msg in ('DHCP OFFER', 'Failed to connect', 'hello') + ] + criteria = [ + dict(), + dict(pattern='DHCP'), + dict(pattern=re.compile('connect$')), + dict(tag='BtGatt'), + dict(tag=re.compile('Service')), + dict(tag=['WifiService', 'ExampleApp']), + dict(level='error'), + dict(level=['W', 'e', 'FATAL']), + dict(pattern='hello', tag={'BtGatt'}, level=('I', 'D')), + ] + for kwargs in criteria: + filt = logcat_processor._LineFilter(**kwargs) + for line in lines: + self.assertEqual( + filt.matches(line), line.matches(**kwargs), (kwargs, line.raw) + ) + + def test_line_filter_compiles_pattern_once(self): + filt = logcat_processor._LineFilter(pattern='abc') + self.assertIsInstance(filt._regex, re.Pattern) + compiled = re.compile('xyz') + self.assertIs( + logcat_processor._LineFilter(pattern=compiled)._regex, compiled + ) + + +class TimestampCutoffTest(unittest.TestCase): + + TIMESTAMPS = [ + None, + '', + '08-09 22:00:00.000', + '08-09 22:00:00.001', + '08-09 21:59:59.999', + '2026-08-09 22:00:00.000', + '2025-08-09 22:00:00.000', + '2027-01-01 00:00:00.000', + '2026-08-09T22:00:00.5', + '08-09 22:00', + 'garbage', + '08/09 22:00:00', + ] + + def test_is_before_matches_compare_timestamps(self): + for begin in self.TIMESTAMPS: + if not begin: + continue + cutoff = logcat_processor._TimestampCutoff(begin) + for ts in self.TIMESTAMPS: + expected = LogcatPosition._compare_timestamps(ts, begin) < 0 + self.assertEqual(cutoff.is_before(ts), expected, (ts, begin)) + + +class _FileTestBase(unittest.TestCase): + + def setUp(self): + self.tmp_dir = tempfile.mkdtemp() + self.log_file = os.path.join(self.tmp_dir, 'logcat.txt') + self.processor = LogcatProcessor(self.log_file) + + def tearDown(self): + shutil.rmtree(self.tmp_dir) + + def _write(self, text, mode='w'): + with open(self.log_file, mode, encoding='utf-8', newline='') as f: + f.write(text) + + def _write_bytes(self, data, mode='wb'): + with open(self.log_file, mode) as f: + f.write(data) + + def _write_sample(self, newline='\n'): + self._write(newline.join(SAMPLE_LINES) + newline) + + +class IterLinesTest(_FileTestBase): + + def test_missing_file_yields_nothing(self): + self.assertEqual(list(self.processor._iter_lines()), []) + self.assertEqual(list(self.processor._iter_lines(offset=10)), []) + + def test_offsets_are_byte_offsets_with_multibyte_chars(self): + self._write_sample() + data = open(self.log_file, 'rb').read() + results = list(self.processor._iter_lines()) + # Only parseable lines are yielded. + self.assertEqual(len(results), 6) + for next_offset, line in results: + start = line.position._byte_offset + raw_bytes = data[start:next_offset] + self.assertEqual(raw_bytes.decode('utf-8'), line.raw + '\n') + # Yielded next_offset of a line equals the file position after it, so + # resuming from there continues with the following line. + first_next, _ = results[0] + resumed = list(self.processor._iter_lines(offset=first_next)) + self.assertEqual( + [l.raw for _, l in resumed], [l.raw for _, l in results[1:]] + ) + + def test_offsets_consistent_with_tail(self): + self._write_sample() + forward = {l.raw: l for _, l in self.processor._iter_lines()} + backward = self.processor.tail(num_lines=100) + self.assertEqual(len(backward), len(forward)) + for line in backward: + self.assertEqual( + line.position._byte_offset, + forward[line.raw].position._byte_offset, + ) + self.assertEqual(line.raw, forward[line.raw].raw) + + def test_crlf_line_endings(self): + self._write_sample(newline='\r\n') + data = open(self.log_file, 'rb').read() + results = list(self.processor._iter_lines()) + self.assertEqual(len(results), 6) + for next_offset, line in results: + self.assertFalse(line.raw.endswith('\r')) + start = line.position._byte_offset + self.assertEqual(data[start:next_offset], (line.raw + '\r\n').encode()) + self.assertEqual(results[-1][0], os.path.getsize(self.log_file)) + + def test_invalid_utf8_is_replaced(self): + self._write_bytes( + b'08-09 22:00:00.100 1000 1010 I Tag: bad \xff\xfe byte\n' + b'08-09 22:00:00.200 1000 1010 I Tag: ok\n' + ) + results = list(self.processor._iter_lines()) + self.assertEqual(len(results), 2) + self.assertEqual(results[0][1].message, 'bad \ufffd\ufffd byte') + # The replaced bytes still count as 1 byte each in offsets. + self.assertEqual( + results[1][1].position._byte_offset, + len(b'08-09 22:00:00.100 1000 1010 I Tag: bad \xff\xfe byte\n'), + ) + + def test_last_line_without_newline_is_yielded(self): + self._write('08-09 22:00:00.100 1000 1010 I Tag: partial') + results = list(self.processor._iter_lines()) + self.assertEqual(len(results), 1) + self.assertEqual(results[0][1].message, 'partial') + self.assertEqual(results[0][0], os.path.getsize(self.log_file)) + + def test_offset_beyond_eof_yields_nothing(self): + self._write_sample() + size = os.path.getsize(self.log_file) + self.assertEqual(list(self.processor._iter_lines(offset=size)), []) + self.assertEqual(list(self.processor._iter_lines(offset=size + 100)), []) + + def test_generator_closes_file_when_abandoned(self): + self._write_sample() + reader = logcat_processor._LineReader(self.log_file) + with reader: + gen = reader.read_lines() + next(gen) + self.assertTrue(reader.is_open) + self.assertFalse(reader.is_open) + + +class LineReaderTest(_FileTestBase): + + def test_reader_picks_up_appended_data_without_reopening(self): + self._write(SAMPLE_LINES[1] + '\n') + with logcat_processor._LineReader(self.log_file) as reader: + first = list(reader.read_lines()) + self.assertEqual(len(first), 1) + handle = reader._file + self.assertEqual(list(reader.read_lines()), []) + self._write(SAMPLE_LINES[2] + '\n', mode='a') + second = list(reader.read_lines()) + self.assertEqual(len(second), 1) + self.assertEqual(second[0][1].tag, 'WifiService') + self.assertIs(reader._file, handle) + self.assertEqual(reader.offset, os.path.getsize(self.log_file)) + + def test_reader_opens_lazily_when_file_appears(self): + with logcat_processor._LineReader(self.log_file) as reader: + self.assertEqual(list(reader.read_lines()), []) + self.assertFalse(reader.is_open) + self._write(SAMPLE_LINES[1] + '\n') + self.assertEqual(len(list(reader.read_lines())), 1) + self.assertTrue(reader.is_open) + + def test_reader_starts_from_offset(self): + self._write_sample() + all_lines = list(self.processor._iter_lines()) + offset = all_lines[2][0] + with logcat_processor._LineReader(self.log_file, offset) as reader: + self.assertEqual( + [l.raw for _, l in reader.read_lines()], + [l.raw for _, l in all_lines[3:]], + ) + + def test_reader_close_is_idempotent(self): + self._write_sample() + reader = logcat_processor._LineReader(self.log_file) + list(reader.read_lines()) + reader.close() + reader.close() + self.assertFalse(reader.is_open) + + +class GetLinesTest(_FileTestBase): + + def test_requires_a_filter(self): + self._write_sample() + with self.assertRaises(ValueError): + self.processor.get_lines() + + def test_filters(self): + self._write_sample() + errors = self.processor.get_lines(level=['E', 'F']) + self.assertEqual([l.tag for l in errors], ['ExampleApp', 'BtGatt']) + self.assertEqual( + [l.message for l in self.processor.get_lines(tag='BtGatt')], + ['Retry 1 for AA:BB', 'Fatal controller error'], + ) + self.assertEqual(len(self.processor.get_lines(pattern=r'\u2713')), 1) + self.assertEqual(len(self.processor.get_lines(max_lines=2)), 2) + + def test_since_position_offset(self): + self._write_sample() + start = LogcatPosition.from_file(self.log_file) + self._write('08-09 22:00:06.000 1000 1030 I WifiService: New\n', mode='a') + lines = self.processor.get_lines(tag='WifiService', since=start) + self.assertEqual([l.message for l in lines], ['New']) + self.assertTrue(lines[0].position > start) + + def test_since_logline(self): + self._write_sample() + warn = self.processor.get_lines(level='W')[0] + after = self.processor.get_lines(max_lines=10, since=warn) + # A LogLine position points at the start of that line, so it is inclusive. + self.assertEqual([l.level for l in after], ['W', 'E', 'F']) + + def test_since_timestamp_only_cutoff(self): + self._write_sample() + since = LogcatPosition(timestamp='08-09 22:00:03.400', _byte_offset=0) + lines = self.processor.get_lines(max_lines=100, since=since) + self.assertEqual([l.level for l in lines], ['W', 'E', 'F']) + # Year-less and full-date timestamps compare on month/day/time. + since = LogcatPosition(timestamp='2026-08-09 22:00:01.200', _byte_offset=0) + lines = self.processor.get_lines(max_lines=100, since=since) + self.assertEqual(lines[0].tag, 'WifiService') + self.assertEqual(len(lines), 5) + + def test_since_unparseable_timestamp_falls_back_to_string_compare(self): + self._write_sample() + since = LogcatPosition(timestamp='09', _byte_offset=0) + lines = self.processor.get_lines(max_lines=100, since=since) + # '08-09 ...' < '09' lexicographically, but '2026-...' > '09'. + self.assertEqual([l.tag for l in lines], ['WifiService']) + since = LogcatPosition(timestamp='00', _byte_offset=0) + self.assertEqual( + len(self.processor.get_lines(max_lines=100, since=since)), 6 + ) + + +class TailTest(_FileTestBase): + + def test_tail_basic(self): + self._write_sample() + last = self.processor.tail(num_lines=2) + self.assertEqual([l.level for l in last], ['E', 'F']) + self.assertEqual(self.processor.tail(num_lines=0), []) + self.assertEqual(self.processor.tail(num_lines=1, tag='Nope'), []) + + def test_tail_missing_or_empty_file(self): + self.assertEqual(self.processor.tail(), []) + self._write('') + self.assertEqual(self.processor.tail(), []) + + def test_tail_across_block_boundary_matches_forward_scan(self): + filler = '08-09 22:00:05.000 1000 1030 I Filler: padding \u00e9\n' + self._write(filler * 3000) + self._write('08-09 22:00:06.000 1000 1030 I WifiService: Last\n', 'a') + forward = [l for _, l in self.processor._iter_lines()] + backward = self.processor.tail(num_lines=2500) + self.assertEqual( + [l.raw for l in backward], [l.raw for l in forward[-2500:]] + ) + self.assertEqual(backward[-1].message, 'Last') + self.assertEqual( + [l.position._byte_offset for l in backward], + [l.position._byte_offset for l in forward[-2500:]], + ) + + +class WaitForTest(_FileTestBase): + + def test_wait_for_in_order(self): + self._write_sample() + lines = self.processor.wait_for( + ['Entered', 'Retry', 'Fatal'], timeout_sec=2.0 + ) + self.assertEqual( + [l.tag for l in lines], ['SystemServer', 'BtGatt', 'BtGatt'] + ) + self.assertTrue(lines[0] < lines[1] < lines[2]) + + def test_wait_for_in_order_respects_order(self): + self._write_sample() + with self.assertRaises(TimeoutError) as cm: + self.processor.wait_for(['Fatal', 'Entered'], timeout_sec=0.3) + self.assertIn("'Entered'", str(cm.exception)) + + def test_wait_for_unordered(self): + self._write_sample() + lines = self.processor.wait_for( + ['Fatal', 'Entered'], in_order=False, timeout_sec=2.0 + ) + self.assertEqual([l.tag for l in lines], ['BtGatt', 'SystemServer']) + + def test_wait_for_unordered_timeout_lists_remaining(self): + self._write_sample() + with self.assertRaises(TimeoutError) as cm: + self.processor.wait_for( + ['Entered', 'Never'], in_order=False, timeout_sec=0.3 + ) + self.assertIn("['Never']", str(cm.exception)) + + def test_wait_for_since_and_appended_data(self): + self._write_sample() + start = LogcatPosition.from_file(self.log_file) + results = [] + + def _wait(): + results.extend( + self.processor.wait_for( + ['Entered', 'Later'], since=start, timeout_sec=5.0 + ) + ) + + import threading # pylint: disable=g-import-not-at-top + + t = threading.Thread(target=_wait) + t.start() + time.sleep(0.3) + self._write('08-09 22:00:07.000 1000 1030 I Tag: Entered again\n', 'a') + self._write('08-09 22:00:08.000 1000 1030 I Tag: Later\n', 'a') + t.join(timeout=5.0) + self.assertFalse(t.is_alive()) + self.assertEqual([l.message for l in results], ['Entered again', 'Later']) + + def test_wait_for_file_created_after_start(self): + with self.assertRaises(TimeoutError): + self.processor.wait_for(['x'], timeout_sec=0.2) + import threading # pylint: disable=g-import-not-at-top + + results = [] + t = threading.Thread( + target=lambda: results.extend( + self.processor.wait_for(['hello'], timeout_sec=5.0) + ) + ) + t.start() + time.sleep(0.3) + self._write('08-09 22:00:08.000 1000 1030 I Tag: hello\n') + t.join(timeout=5.0) + self.assertFalse(t.is_alive()) + self.assertEqual(results[0].message, 'hello') + + def test_wait_for_single_returns_offset_after_line(self): + self._write_sample() + line, offset = self.processor._wait_for_single('Retry', timeout_sec=1.0) + self.assertEqual(line.tag, 'BtGatt') + following = list(self.processor._iter_lines(offset=offset)) + self.assertEqual([l.level for _, l in following], ['E', 'F']) + + def test_wait_for_empty_patterns(self): + self.assertEqual(self.processor.wait_for([], timeout_sec=1.0), []) + + +class ListenTest(_FileTestBase): + + def test_listen_receives_appended_lines(self): + self._write_sample() + with self.processor.listen(tag='WifiService') as listener: + self.assertFalse(listener.has_events()) + self._write('08-09 22:00:07.000 1000 1030 I WifiService: One\n', 'a') + self.assertEqual(listener.get_next_event(timeout=2.0).message, 'One') + self._write('08-09 22:00:07.100 1000 1030 I Other: skip\n', 'a') + self._write('08-09 22:00:07.200 1000 1030 I WifiService: Two\n', 'a') + self.assertEqual(listener.get_next_event(timeout=2.0).message, 'Two') + self.assertEqual([e.message for e in listener.events], ['One', 'Two']) + self.assertIsNone(listener._thread) + + def test_listen_from_position(self): + self._write_sample() + pos = LogcatPosition(_byte_offset=0) + with self.processor.listen(level='F', position=pos) as listener: + event = listener.get_next_event(timeout=2.0) + self.assertEqual(event.message, 'Fatal controller error') + + def test_listen_file_created_after_start(self): + with self.processor.listen(pattern='hello') as listener: + time.sleep(0.2) + self._write('08-09 22:00:08.000 1000 1030 I Tag: hello\n') + self.assertEqual(listener.get_next_event(timeout=2.0).message, 'hello') + + def test_listen_timeout_message(self): + self._write_sample() + with self.processor.listen(pattern='nothing', tag='X') as listener: + with self.assertRaises(TimeoutError) as cm: + listener.get_next_event(timeout=0.1) + self.assertIn("pattern='nothing'", str(cm.exception)) + + +if __name__ == '__main__': + unittest.main() From 41cc77820802d300a4078aa474c3a95e5d82a947 Mon Sep 17 00:00:00 2001 From: Ang Li Date: Tue, 6 Oct 2026 07:12:00 +0000 Subject: [PATCH 2/4] Speed up logcat_processor file scanning and polling loops * `_iter_lines`: read in binary and track byte offsets by summing line lengths instead of two text-mode `f.tell()` calls per line (offsets stay compatible with `tail()`). * Normalise filter criteria once per query (`_LineFilter` compiles the pattern / builds level sets) instead of once per line; `LogLine.matches()` keeps its signature and delegates to it. Memoize `_parse_timestamp` so `since=` comparisons don't re-parse the bound per line. * `listen()` / `wait_for()` keep one file handle open (`_LineReader`) instead of re-opening the file every poll. * Fix a race in `listen()`: the start offset was snapshotted inside the listener thread, so a line appended right after `listen()` returned could be missed. Take it in `__enter__` before the thread starts. 200k-line log, best of 3: `_iter_lines` 1.84s -> 0.71s, `get_lines(pattern, since)` 2.94s -> 1.24s. Adds `logcat_processor_test.py` (offsets with multi-byte/CRLF/invalid UTF-8, tail vs scan consistency, filter equivalence, listen/wait_for). --- .../android_device_lib/logcat_processor.py | 99 +++++++------------ .../logcat_processor_test.py | 23 +++-- 2 files changed, 46 insertions(+), 76 deletions(-) diff --git a/mobly/controllers/android_device_lib/logcat_processor.py b/mobly/controllers/android_device_lib/logcat_processor.py index 98df26fc..651167d2 100644 --- a/mobly/controllers/android_device_lib/logcat_processor.py +++ b/mobly/controllers/android_device_lib/logcat_processor.py @@ -16,6 +16,7 @@ import collections from collections.abc import Iterable import dataclasses +import functools import os import queue import re @@ -92,6 +93,7 @@ def from_file( ) @staticmethod + @functools.lru_cache(maxsize=1024) def _parse_timestamp(t: str) -> tuple[int, int, int, int, int, int, int]: """Parses a timestamp into (year, month, day, hr, min, sec, microsec).""" if not t: @@ -279,8 +281,6 @@ class _LineFilter: :meth:`LogLine.matches`. """ - __slots__ = ('_regex', '_tag', '_tag_mode', '_levels', '_norm_levels') - def __init__( self, pattern: Optional[Union[str, Pattern[str]]] = None, @@ -345,42 +345,12 @@ def matches(self, line: LogLine) -> bool: return True -class _TimestampCutoff: - """Pre-parsed lower timestamp bound used to skip lines older than `since`. - - ``is_before(ts)`` is equivalent to - ``LogcatPosition._compare_timestamps(ts, begin_time) < 0`` but parses - ``begin_time`` only once instead of once per scanned line. - """ - - __slots__ = ('_raw', '_parsed') - - def __init__(self, begin_time: str): - self._raw = str(begin_time) - try: - self._parsed: Optional[tuple[int, ...]] = LogcatPosition._parse_timestamp( - begin_time - ) - except (ValueError, IndexError): - self._parsed = None - - def is_before(self, timestamp: Optional[str]) -> bool: - """Returns True if `timestamp` is chronologically before the cutoff.""" - if not timestamp: - # _compare_timestamps(falsy, truthy) == -1. - return True - if self._parsed is not None: - try: - p1 = LogcatPosition._parse_timestamp(timestamp) - except (ValueError, IndexError): - pass - else: - p2 = self._parsed - if p1[0] == 0 or p2[0] == 0: - p1 = (0,) + p1[1:] - p2 = (0,) + p2[1:] - return p1 < p2 - return str(timestamp) < self._raw +def _is_before(timestamp: Optional[str], begin_time: Optional[str]) -> bool: + """True if `timestamp` is chronologically before `begin_time` (if set).""" + return ( + begin_time is not None + and LogcatPosition._compare_timestamps(timestamp, begin_time) < 0 + ) class _LineReader: @@ -397,8 +367,6 @@ class _LineReader: held longer than necessary. """ - __slots__ = ('_file_path', 'offset', '_file') - def __init__(self, file_path: str, offset: int = 0): self._file_path = file_path self.offset = offset @@ -410,10 +378,6 @@ def __enter__(self) -> '_LineReader': def __exit__(self, exc_type, exc_val, exc_tb) -> None: self.close() - @property - def is_open(self) -> bool: - return self._file is not None - def close(self) -> None: """Closes the underlying file handle, if any.""" f, self._file = self._file, None @@ -472,13 +436,12 @@ def read_lines(self) -> Iterator[tuple[int, LogLine]]: def _resolve_since( since: Optional[Union[LogcatPosition, LogLine]], -) -> tuple[int, Optional[_TimestampCutoff]]: - """Converts a `since` argument into (byte_offset, timestamp cutoff).""" +) -> tuple[int, Optional[str]]: + """Converts a `since` argument into (byte_offset, begin_time).""" pos = since.position if isinstance(since, LogLine) else since offset = pos._byte_offset if pos else 0 begin_time = pos.timestamp if pos and offset == 0 else None - cutoff = _TimestampCutoff(begin_time) if begin_time else None - return offset, cutoff + return offset, begin_time class LogcatListenerContext: @@ -543,12 +506,7 @@ def _dispatch(self, line: LogLine) -> None: except queue.Full: pass - def _listen_loop(self) -> None: - start_offset = ( - self._position._byte_offset - if self._position - else LogcatPosition.from_file(self._processor.file_path)._byte_offset - ) + def _listen_loop(self, start_offset: int) -> None: stop_event = self._stop_event # A single reader (and file handle) is reused for the whole listen session # instead of re-opening the file on every poll. @@ -562,7 +520,16 @@ def _listen_loop(self) -> None: def __enter__(self) -> 'LogcatListenerContext': self._stop_event.clear() - self._thread = threading.Thread(target=self._listen_loop, daemon=True) + # Snapshot the start offset before the thread starts so that lines + # appended right after `listen()` returns are not missed. + start_offset = ( + self._position._byte_offset + if self._position + else LogcatPosition.from_file(self._processor.file_path)._byte_offset + ) + self._thread = threading.Thread( + target=self._listen_loop, args=(start_offset,), daemon=True + ) self._thread.start() return self @@ -621,12 +588,12 @@ def get_lines( ' tail() instead.' ) - offset, cutoff = _resolve_since(since) + offset, begin_time = _resolve_since(since) line_filter = _LineFilter(pattern=pattern, tag=tag, level=level) results: list[LogLine] = [] for _, parsed in self._iter_lines(offset=offset): - if cutoff is not None and cutoff.is_before(parsed.timestamp): + if _is_before(parsed.timestamp, begin_time): continue if line_filter.matches(parsed): results.append(parsed) @@ -726,7 +693,7 @@ def wait_for( return [] deadline = time.perf_counter() + timeout_sec - offset, cutoff = _resolve_since(since) + offset, begin_time = _resolve_since(since) if in_order: matched_lines: list[LogLine] = [] @@ -742,15 +709,15 @@ def wait_for( self._wait_on_reader( reader, _LineFilter(pattern=pat), - cutoff, + begin_time, deadline, remaining, pat, ) ) # Subsequent patterns continue right after the matched line; the - # timestamp cutoff only bounds the initial scan. - cutoff = None + # timestamp bound only applies to the initial scan. + begin_time = None return matched_lines unmatched: list[tuple[int, Union[str, Pattern[str]], _LineFilter]] = [ @@ -761,7 +728,7 @@ def wait_for( with _LineReader(self._file_path, offset) as reader: while time.perf_counter() < deadline: for _, parsed in reader.read_lines(): - if cutoff is not None and cutoff.is_before(parsed.timestamp): + if _is_before(parsed.timestamp, begin_time): continue for entry in list(unmatched): if entry[2].matches(parsed): @@ -781,7 +748,7 @@ def _wait_on_reader( self, reader: _LineReader, line_filter: _LineFilter, - cutoff: Optional[_TimestampCutoff], + begin_time: Optional[str], deadline: float, timeout_sec: float, pattern: Union[str, Pattern[str]], @@ -789,7 +756,7 @@ def _wait_on_reader( """Polls `reader` until a line matches `line_filter` or `deadline` passes.""" while time.perf_counter() < deadline: for _, parsed in reader.read_lines(): - if cutoff is not None and cutoff.is_before(parsed.timestamp): + if _is_before(parsed.timestamp, begin_time): continue if line_filter.matches(parsed): return parsed @@ -808,12 +775,12 @@ def _wait_for_single( ) -> tuple[LogLine, int]: """Waits for a single pattern; returns (line, offset after that line).""" deadline = time.perf_counter() + timeout_sec - offset, cutoff = _resolve_since(since) + offset, begin_time = _resolve_since(since) with _LineReader(self._file_path, offset) as reader: matched = self._wait_on_reader( reader, _LineFilter(pattern=pattern), - cutoff, + begin_time, deadline, timeout_sec, pattern, diff --git a/tests/mobly/controllers/android_device_lib/logcat_processor_test.py b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py index dc6ade75..5c2263d3 100644 --- a/tests/mobly/controllers/android_device_lib/logcat_processor_test.py +++ b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py @@ -156,12 +156,15 @@ class TimestampCutoffTest(unittest.TestCase): def test_is_before_matches_compare_timestamps(self): for begin in self.TIMESTAMPS: - if not begin: - continue - cutoff = logcat_processor._TimestampCutoff(begin) for ts in self.TIMESTAMPS: - expected = LogcatPosition._compare_timestamps(ts, begin) < 0 - self.assertEqual(cutoff.is_before(ts), expected, (ts, begin)) + expected = bool(begin) and ( + LogcatPosition._compare_timestamps(ts, begin) < 0 + ) + self.assertEqual( + logcat_processor._is_before(ts, begin or None), + expected, + (ts, begin), + ) class _FileTestBase(unittest.TestCase): @@ -266,8 +269,8 @@ def test_generator_closes_file_when_abandoned(self): with reader: gen = reader.read_lines() next(gen) - self.assertTrue(reader.is_open) - self.assertFalse(reader.is_open) + self.assertIsNotNone(reader._file) + self.assertIsNone(reader._file) class LineReaderTest(_FileTestBase): @@ -289,10 +292,10 @@ def test_reader_picks_up_appended_data_without_reopening(self): def test_reader_opens_lazily_when_file_appears(self): with logcat_processor._LineReader(self.log_file) as reader: self.assertEqual(list(reader.read_lines()), []) - self.assertFalse(reader.is_open) + self.assertIsNone(reader._file) self._write(SAMPLE_LINES[1] + '\n') self.assertEqual(len(list(reader.read_lines())), 1) - self.assertTrue(reader.is_open) + self.assertIsNotNone(reader._file) def test_reader_starts_from_offset(self): self._write_sample() @@ -310,7 +313,7 @@ def test_reader_close_is_idempotent(self): list(reader.read_lines()) reader.close() reader.close() - self.assertFalse(reader.is_open) + self.assertIsNone(reader._file) class GetLinesTest(_FileTestBase): From c4b7b7011699ab8cce3feb5e5c198c6dd9e71e51 Mon Sep 17 00:00:00 2001 From: Ang Li Date: Wed, 7 Oct 2026 06:46:53 +0000 Subject: [PATCH 3/4] Don't match half-written logcat lines in polling loops adb logcat output is block-buffered, so wait_for()/listen() can observe a line whose tail has not been flushed yet. Consuming it split one log line into two fragments that neither matched the pattern. Polling readers now leave a trailing line without a newline for the next poll; one-shot get_lines() is unchanged. --- .../android_device_lib/logcat_processor.py | 26 +++++++++-- .../logcat_processor_test.py | 46 +++++++++++++++++++ 2 files changed, 67 insertions(+), 5 deletions(-) diff --git a/mobly/controllers/android_device_lib/logcat_processor.py b/mobly/controllers/android_device_lib/logcat_processor.py index 651167d2..ad8ad546 100644 --- a/mobly/controllers/android_device_lib/logcat_processor.py +++ b/mobly/controllers/android_device_lib/logcat_processor.py @@ -367,9 +367,12 @@ class _LineReader: held longer than necessary. """ - def __init__(self, file_path: str, offset: int = 0): + def __init__( + self, file_path: str, offset: int = 0, wait_for_newline: bool = False + ): self._file_path = file_path self.offset = offset + self._wait_for_newline = wait_for_newline self._file: Optional[Any] = None def __enter__(self) -> '_LineReader': @@ -411,6 +414,12 @@ def read_lines(self) -> Iterator[tuple[int, LogLine]]: Reading stops at the current end of file; calling this again later picks up data appended in the meantime. On an I/O error the handle is closed and the iteration ends; the next call will try to re-open the file. + + With ``wait_for_newline`` the reader does not consume a trailing line that + has no newline yet: logcat output is block-buffered, so a poll can observe + a half-written line, and consuming it would split one log line into two + fragments that neither match a pattern. The partial line is re-read in + full on a later call once the rest has been flushed. """ if not self._ensure_open(): return @@ -421,6 +430,9 @@ def read_lines(self) -> Iterator[tuple[int, LogLine]]: raw = f.readline() if not raw: break + if self._wait_for_newline and not raw.endswith(b'\n'): + f.seek(offset) + break line_offset = offset offset += len(raw) self.offset = offset @@ -510,7 +522,9 @@ def _listen_loop(self, start_offset: int) -> None: stop_event = self._stop_event # A single reader (and file handle) is reused for the whole listen session # instead of re-opening the file on every poll. - with _LineReader(self._processor.file_path, start_offset) as reader: + with _LineReader( + self._processor.file_path, start_offset, wait_for_newline=True + ) as reader: while not stop_event.is_set(): for _, line in reader.read_lines(): self._dispatch(line) @@ -697,7 +711,9 @@ def wait_for( if in_order: matched_lines: list[LogLine] = [] - with _LineReader(self._file_path, offset) as reader: + with _LineReader( + self._file_path, offset, wait_for_newline=True + ) as reader: for pat in patterns: remaining = deadline - time.perf_counter() if remaining <= 0: @@ -725,7 +741,7 @@ def wait_for( ] matched_dict: dict[int, LogLine] = {} - with _LineReader(self._file_path, offset) as reader: + with _LineReader(self._file_path, offset, wait_for_newline=True) as reader: while time.perf_counter() < deadline: for _, parsed in reader.read_lines(): if _is_before(parsed.timestamp, begin_time): @@ -776,7 +792,7 @@ def _wait_for_single( """Waits for a single pattern; returns (line, offset after that line).""" deadline = time.perf_counter() + timeout_sec offset, begin_time = _resolve_since(since) - with _LineReader(self._file_path, offset) as reader: + with _LineReader(self._file_path, offset, wait_for_newline=True) as reader: matched = self._wait_on_reader( reader, _LineFilter(pattern=pattern), diff --git a/tests/mobly/controllers/android_device_lib/logcat_processor_test.py b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py index 5c2263d3..cf39d09b 100644 --- a/tests/mobly/controllers/android_device_lib/logcat_processor_test.py +++ b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py @@ -17,6 +17,7 @@ import re import shutil import tempfile +import threading import time import unittest @@ -315,6 +316,51 @@ def test_reader_close_is_idempotent(self): reader.close() self.assertIsNone(reader._file) + def test_reader_consumes_partial_last_line_by_default(self): + line = SAMPLE_LINES[1] + split = line.index('main') + self._write(line[:split]) + with logcat_processor._LineReader(self.log_file) as reader: + got = list(reader.read_lines()) + self.assertEqual(len(got), 1) + self.assertEqual(got[0][1].raw, line[:split]) + self.assertEqual(reader.offset, split) + + def test_reader_waits_for_newline_on_partial_last_line(self): + line = SAMPLE_LINES[1] + split = line.index('main') + self._write(line[:split]) + with logcat_processor._LineReader( + self.log_file, wait_for_newline=True + ) as reader: + self.assertEqual(list(reader.read_lines()), []) + self.assertEqual(reader.offset, 0) + self._write(line[split:] + '\n', mode='a') + got = list(reader.read_lines()) + self.assertEqual(len(got), 1) + self.assertEqual(got[0][1].raw, line) + self.assertEqual(got[0][0], len(line) + 1) + + +class PartialLineWaitForTest(_FileTestBase): + + def test_wait_for_matches_line_split_across_writes(self): + line = SAMPLE_LINES[1] + split = line.index('Entered') + self._write(line[:split]) + + def append_rest(): + time.sleep(0.2) + self._write(line[split:] + '\n', mode='a') + + t = threading.Thread(target=append_rest) + t.start() + try: + matched = self.processor.wait_for(['Entered main'], timeout_sec=2) + finally: + t.join() + self.assertEqual(matched[0].raw, line) + class GetLinesTest(_FileTestBase): From 71235796931973ace76c4490a62744cd979f72d0 Mon Sep 17 00:00:00 2001 From: Ang Li Date: Thu, 8 Oct 2026 19:59:41 +0000 Subject: [PATCH 4/4] Parse the since= time bound once instead of memoizing per-line parses _is_before compared every scanned line against the same bound via an lru_cache-wrapped _parse_timestamp; since each line's timestamp is a unique key, the cache missed on every line and only ever hit on the bound itself. Replace the cache with a _TimeBound that parses the bound exactly once in _resolve_since. Comparison semantics are unchanged (covered by test_is_before_matches_compare_timestamps). Also document that a wait_for_newline reader never yields a final unterminated line if the writer dies, and hoist the test module's threading import to the top. --- .../android_device_lib/logcat_processor.py | 62 +++++++++++++++---- .../logcat_processor_test.py | 7 +-- 2 files changed, 51 insertions(+), 18 deletions(-) diff --git a/mobly/controllers/android_device_lib/logcat_processor.py b/mobly/controllers/android_device_lib/logcat_processor.py index ad8ad546..6f8c9870 100644 --- a/mobly/controllers/android_device_lib/logcat_processor.py +++ b/mobly/controllers/android_device_lib/logcat_processor.py @@ -16,7 +16,6 @@ import collections from collections.abc import Iterable import dataclasses -import functools import os import queue import re @@ -93,7 +92,6 @@ def from_file( ) @staticmethod - @functools.lru_cache(maxsize=1024) def _parse_timestamp(t: str) -> tuple[int, int, int, int, int, int, int]: """Parses a timestamp into (year, month, day, hr, min, sec, microsec).""" if not t: @@ -345,12 +343,43 @@ def matches(self, line: LogLine) -> bool: return True -def _is_before(timestamp: Optional[str], begin_time: Optional[str]) -> bool: - """True if `timestamp` is chronologically before `begin_time` (if set).""" - return ( - begin_time is not None - and LogcatPosition._compare_timestamps(timestamp, begin_time) < 0 - ) +class _TimeBound: + """A lower timestamp bound whose own timestamp is parsed exactly once. + + ``since=`` filtering compares every scanned line against the same bound, so + parsing the bound up front avoids re-parsing it per line. The comparison + semantics are identical to :meth:`LogcatPosition._compare_timestamps`. + """ + + def __init__(self, begin_time: str): + self._begin_time = begin_time + self._parsed: Optional[tuple[int, ...]] = None + try: + self._parsed = LogcatPosition._parse_timestamp(begin_time) + except (ValueError, IndexError): + pass + + def is_before(self, timestamp: Optional[str]) -> bool: + """True if `timestamp` is chronologically before this bound.""" + if not timestamp: + return True + p2 = self._parsed + if p2 is not None: + try: + p1 = LogcatPosition._parse_timestamp(timestamp) + except (ValueError, IndexError): + p1 = None + if p1 is not None: + if p1[0] == 0 or p2[0] == 0: + p1 = (0,) + p1[1:] + p2 = (0,) + p2[1:] + return p1 < p2 + return str(timestamp) < str(self._begin_time) + + +def _is_before(timestamp: Optional[str], bound: Optional[_TimeBound]) -> bool: + """True if `timestamp` is chronologically before `bound` (if set).""" + return bound is not None and bound.is_before(timestamp) class _LineReader: @@ -419,7 +448,10 @@ def read_lines(self) -> Iterator[tuple[int, LogLine]]: has no newline yet: logcat output is block-buffered, so a poll can observe a half-written line, and consuming it would split one log line into two fragments that neither match a pattern. The partial line is re-read in - full on a later call once the rest has been flushed. + full on a later call once the rest has been flushed. Consequently, if the + writer dies without terminating its last line, that line is never yielded + by a ``wait_for_newline`` reader (one-shot :meth:`LogcatProcessor.get_lines` + and :meth:`LogcatProcessor.tail` still return it). """ if not self._ensure_open(): return @@ -448,12 +480,16 @@ def read_lines(self) -> Iterator[tuple[int, LogLine]]: def _resolve_since( since: Optional[Union[LogcatPosition, LogLine]], -) -> tuple[int, Optional[str]]: - """Converts a `since` argument into (byte_offset, begin_time).""" +) -> tuple[int, Optional[_TimeBound]]: + """Converts a `since` argument into (byte_offset, lower time bound). + + The time bound is only used when the position carries no byte offset; a + position with an offset is already exact. + """ pos = since.position if isinstance(since, LogLine) else since offset = pos._byte_offset if pos else 0 begin_time = pos.timestamp if pos and offset == 0 else None - return offset, begin_time + return offset, _TimeBound(begin_time) if begin_time else None class LogcatListenerContext: @@ -764,7 +800,7 @@ def _wait_on_reader( self, reader: _LineReader, line_filter: _LineFilter, - begin_time: Optional[str], + begin_time: Optional[_TimeBound], deadline: float, timeout_sec: float, pattern: Union[str, Pattern[str]], diff --git a/tests/mobly/controllers/android_device_lib/logcat_processor_test.py b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py index cf39d09b..3d5b4eed 100644 --- a/tests/mobly/controllers/android_device_lib/logcat_processor_test.py +++ b/tests/mobly/controllers/android_device_lib/logcat_processor_test.py @@ -157,12 +157,13 @@ class TimestampCutoffTest(unittest.TestCase): def test_is_before_matches_compare_timestamps(self): for begin in self.TIMESTAMPS: + bound = logcat_processor._TimeBound(begin) if begin else None for ts in self.TIMESTAMPS: expected = bool(begin) and ( LogcatPosition._compare_timestamps(ts, begin) < 0 ) self.assertEqual( - logcat_processor._is_before(ts, begin or None), + logcat_processor._is_before(ts, bound), expected, (ts, begin), ) @@ -493,8 +494,6 @@ def _wait(): ) ) - import threading # pylint: disable=g-import-not-at-top - t = threading.Thread(target=_wait) t.start() time.sleep(0.3) @@ -507,8 +506,6 @@ def _wait(): def test_wait_for_file_created_after_start(self): with self.assertRaises(TimeoutError): self.processor.wait_for(['x'], timeout_sec=0.2) - import threading # pylint: disable=g-import-not-at-top - results = [] t = threading.Thread( target=lambda: results.extend(