-
Notifications
You must be signed in to change notification settings - Fork 35
Fix decoding error (issue 168) #169
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
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -27,6 +27,7 @@ | |
| import subprocess | ||
| import threading | ||
| import time | ||
| from codecs import getincrementaldecoder | ||
|
|
||
| from aexpect.exceptions import ( | ||
| ExpectError, | ||
|
|
@@ -56,6 +57,9 @@ | |
|
|
||
| LOG = logging.getLogger(__name__) | ||
|
|
||
| # Buffer size in byte for pipe reads | ||
| READ_BUFFER_SIZE = 1024 | ||
|
|
||
|
|
||
| def kill_tail_threads(): | ||
| """ | ||
|
|
@@ -731,6 +735,8 @@ def _print_line(text): | |
| poller = select.poll() | ||
| poller.register(tail_pipe, select.POLLIN) | ||
| bfr = "" | ||
| decoder_class = getincrementaldecoder(self.encoding) | ||
| decoder = decoder_class(errors="ignore") | ||
| while True: | ||
| if _THREAD_KILL_REQUESTED.is_set(): | ||
| try: | ||
|
|
@@ -745,10 +751,10 @@ def _print_line(text): | |
| break | ||
| if poll_status: | ||
| # Some data is available; read it | ||
| new_data = os.read(tail_pipe, 1024) | ||
| if not new_data: | ||
| new_bytes = os.read(tail_pipe, READ_BUFFER_SIZE) | ||
| if not new_bytes: | ||
| break | ||
| new_data = new_data.decode(self.encoding, "ignore") | ||
| new_data = decoder.decode(input=new_bytes) | ||
| if not new_data: # all chars were ignored, skip round | ||
| continue | ||
| bfr += new_data | ||
|
|
@@ -901,23 +907,23 @@ def _read_nonblocking(self, internal_timeout=None, timeout=None): | |
| expect_pipe = self._get_fd("expect") | ||
| poller = select.poll() | ||
| poller.register(expect_pipe, select.POLLIN) | ||
| data = "" | ||
| data = b"" | ||
|
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. 🎯 Functional Correctness | 🟠 Major | ⚡ Quick win 🔎 Supported by static analysis🤖 get_repo_knowledge executed:
Length of output: 690 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- changed hunk ---'
git diff -- aexpect/client.py | sed -n '1,240p'
printf '%s\n' '--- relevant symbols ---'
rg -n -C 8 'class (Expect|Tail)|def (_read_nonblocking|_tail)|read_nonblocking|IncrementalDecoder|errors=' aexpect/client.pyRepository: avocado-framework/aexpect Length of output: 19692 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- Tail._tail ---'
sed -n '715,802p' aexpect/client.py
printf '%s\n' '--- Expect initialization and _read_nonblocking ---'
sed -n '810,930p' aexpect/client.py
printf '%s\n' '--- decoder-related imports and tests ---'
sed -n '1,90p' aexpect/client.py
rg -n -C 6 'read_nonblocking|嗨|incremental|decode|encoding' tests aexpect 2>/dev/null | sed -n '1,260p'Repository: avocado-framework/aexpect Length of output: 27215 Preserve decoder state between
🤖 Prompt for AI Agents |
||
| read = 0 | ||
| while True: | ||
| try: | ||
| poll_status = poller.poll(internal_timeout) | ||
| except select.error: | ||
| return read, data | ||
| return read, data.decode(self.encoding, "ignore") | ||
| if poll_status: | ||
| raw_data = os.read(expect_pipe, 1024) | ||
| raw_data = os.read(expect_pipe, READ_BUFFER_SIZE) | ||
| if not raw_data: | ||
| return read, data | ||
| return read, data.decode(self.encoding, "ignore") | ||
|
Contributor
Author
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. I assume from your comment in #168 (comment) this could be instead be turned into
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. I'd suggest another commit to do mass ignore->replace which could be reverted in case someone depended on the ignore. |
||
| read += len(raw_data) | ||
| data += raw_data.decode(self.encoding, "ignore") | ||
| data += raw_data | ||
|
Contributor
Author
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. Definitely better not to decode raw data until the very end, I think this change improves the clarity and related better to the choice of naming. |
||
| else: | ||
| return read, data | ||
| return read, data.decode(self.encoding, "ignore") | ||
| if end_time and time.monotonic() > end_time: | ||
| return read, data | ||
| return read, data.decode(self.encoding, "ignore") | ||
|
|
||
| def read_nonblocking(self, internal_timeout=None, timeout=None): | ||
| """ | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -18,6 +18,7 @@ | |
| import random | ||
| import string | ||
| import sys | ||
| import time | ||
| import unittest | ||
|
|
||
| from aexpect import client | ||
|
|
@@ -194,5 +195,88 @@ def get_proc_fds(): | |
| ) | ||
|
|
||
|
|
||
| class EncodingTest(unittest.TestCase): | ||
|
|
||
| TEXT = "嗨😀" | ||
| MAX_OFFSET = 10 | ||
|
|
||
| def _multibyte_write_cmd(self, offset, count=1): | ||
| """Build a Python command that writes multibyte text to stdout.""" | ||
| encoded = self.TEXT.encode("utf-8") | ||
| reps = 1024 // len(encoded) + 1 | ||
| writes = "; ".join(["f.write(t); f.flush()"] * count) | ||
| return ( | ||
| f"import os,sys; t=b' '*{offset}+{encoded!r}*{reps}+b'\\n'; " | ||
| f"f=os.fdopen(sys.stdout.fileno(),'wb',closefd=False); {writes}" | ||
| ) | ||
|
|
||
| @unittest.skipUnless(os.name == "posix", "Unix/Linux/macOS only") | ||
| def test_shell(self): | ||
| """Test multibyte decoding in ShellSession across buffer boundaries.""" | ||
| sess = client.ShellSession("/bin/sh") | ||
| sess.cmd_output("echo init") | ||
| lengths = [] | ||
| for offset in range(self.MAX_OFFSET): | ||
| cmd = self._multibyte_write_cmd(offset) | ||
| result = sess.cmd_output(f'{sys.executable} -c "{cmd}"').lstrip() | ||
| self.assertTrue( | ||
| result.startswith(self.TEXT), | ||
| f"offset {offset}: unexpected start: {result[:20]!r}", | ||
| ) | ||
| lengths.append(len(result)) | ||
| sess.close() | ||
| self.assertTrue(lengths, "No output collected") | ||
| self.assertEqual( | ||
| len(set(lengths)), | ||
| 1, | ||
| f"Output lengths vary across offsets: {lengths}", | ||
| ) | ||
|
|
||
| @unittest.skipUnless(os.name == "posix", "Unix/Linux/macOS only") | ||
| def test_tail(self): | ||
| """Test multibyte decoding in Tail across buffer boundaries.""" | ||
| tail_lines = 3 | ||
| lengths = [] | ||
| output_buffer = [] | ||
| for offset in range(self.MAX_OFFSET): | ||
| output_buffer = [] | ||
| terminated = False | ||
|
|
||
| def on_output(text): | ||
| nonlocal output_buffer | ||
| output_buffer.append(text) | ||
|
|
||
| def on_terminate(_status): | ||
| nonlocal terminated | ||
| terminated = True | ||
|
|
||
| cmd = self._multibyte_write_cmd(offset, count=tail_lines) | ||
| tail = client.Tail( | ||
| f'{sys.executable} -c "{cmd}"', | ||
| output_func=on_output, | ||
| termination_func=on_terminate, | ||
| ) | ||
| for _ in range(1000): | ||
| if terminated: | ||
| break | ||
| time.sleep(0.01) | ||
| tail.close() | ||
| for line in output_buffer: | ||
| if line.startswith("(Process terminated "): | ||
| continue | ||
| stripped = line.lstrip() | ||
| self.assertTrue( | ||
| stripped.startswith(self.TEXT), | ||
| f"offset {offset}: unexpected start: {stripped[:20]!r}", | ||
| ) | ||
| lengths.append(len(stripped)) | ||
|
Comment on lines
+264
to
+272
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. 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win Require data output for every offset. This loop only validates lines that exist. If one offset produces no non-status callback, other offsets can still populate 🤖 Prompt for AI Agents |
||
| self.assertTrue(lengths, "No output collected") | ||
| self.assertEqual( | ||
| len(set(lengths)), | ||
| 1, | ||
| f"Output lengths vary across offsets: {lengths}", | ||
| ) | ||
|
|
||
|
|
||
| if __name__ == "__main__": | ||
| unittest.main() | ||
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.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
🔎 Supported by static analysis
🤖 get_repo_knowledge executed:
get_repo_knowledge avocado-framework/aexpect /tmp/coderabbit-repo-knowledge/avocado-framework-aexpect-ffa119cfLength of output: 662
🏁 Script executed:
Repository: avocado-framework/aexpect
Length of output: 5847
🏁 Script executed:
Repository: avocado-framework/aexpect
Length of output: 18886
🏁 Script executed:
Repository: avocado-framework/aexpect
Length of output: 191
Finalize the decoder and preserve replacement characters at EOF.
Tail._tailuseserrors="ignore", so an incomplete sequence is discarded even when the decoder is finalized. Useerrors="replace"and calldecoder.decode(b"", final=True)before the final buffer dispatch.🤖 Prompt for AI Agents