| 153 | |
| 154 | @dataclass |
| 155 | class LiveClient: |
| 156 | name: str |
| 157 | proc: subprocess.Popen |
| 158 | master_fd: int |
| 159 | stop: threading.Event = field(default_factory=threading.Event) |
| 160 | thread: threading.Thread | None = None |
| 161 | total_bytes: int = 0 |
| 162 | |
| 163 | def _pump(self) -> None: |
| 164 | buffer = b"" |
| 165 | while not self.stop.is_set(): |
| 166 | try: |
| 167 | rlist, _, _ = select.select([self.master_fd], [], [], 0.1) |
| 168 | except (OSError, ValueError): |
| 169 | break |
| 170 | if not rlist: |
| 171 | if self.proc.poll() is not None: |
| 172 | break |
| 173 | continue |
| 174 | try: |
| 175 | chunk = os.read(self.master_fd, 65536) |
| 176 | except (BlockingIOError, OSError): |
| 177 | continue |
| 178 | if not chunk: |
| 179 | break |
| 180 | self.total_bytes += len(chunk) |
| 181 | buffer = (buffer + chunk)[-8192:] |
| 182 | changed = True |
| 183 | while changed: |
| 184 | changed = False |
| 185 | for query, response in _TERM_REPLIES: |
| 186 | if query in buffer: |
| 187 | try: |
| 188 | os.write(self.master_fd, response) |
| 189 | except OSError: |
| 190 | pass |
| 191 | buffer = buffer.replace(query, b"") |
| 192 | changed = True |
| 193 | |
| 194 | def start_pump(self) -> None: |
| 195 | self.thread = threading.Thread(target=self._pump, daemon=True) |
| 196 | self.thread.start() |
| 197 | |
| 198 | def alive(self) -> bool: |
| 199 | return self.proc.poll() is None |
| 200 | |
| 201 | def send_keys(self, data: bytes) -> None: |
| 202 | os.write(self.master_fd, data) |
| 203 | |
| 204 | def shutdown(self) -> None: |
| 205 | self.stop.set() |
| 206 | if self.thread: |
| 207 | self.thread.join(timeout=1.0) |
| 208 | try: |
| 209 | os.killpg(self.proc.pid, signal.SIGTERM) |
| 210 | except (ProcessLookupError, PermissionError): |
| 211 | pass |
| 212 | try: |