Asynchronously read process output
(self)
| 239 | return cmd |
| 240 | |
| 241 | async def _read_output(self): |
| 242 | """Asynchronously read process output""" |
| 243 | loop = asyncio.get_event_loop() |
| 244 | |
| 245 | try: |
| 246 | while self.process and self.process.poll() is None: |
| 247 | # Read a line in thread pool |
| 248 | line = await loop.run_in_executor( |
| 249 | None, self.process.stdout.readline |
| 250 | ) |
| 251 | if line: |
| 252 | line = line.strip() |
| 253 | if line: |
| 254 | level = self._parse_log_level(line) |
| 255 | entry = self._create_log_entry(line, level) |
| 256 | await self._push_log(entry) |
| 257 | |
| 258 | # Read remaining output |
| 259 | if self.process and self.process.stdout: |
| 260 | remaining = await loop.run_in_executor( |
| 261 | None, self.process.stdout.read |
| 262 | ) |
| 263 | if remaining: |
| 264 | for line in remaining.strip().split('\n'): |
| 265 | if line.strip(): |
| 266 | level = self._parse_log_level(line) |
| 267 | entry = self._create_log_entry(line.strip(), level) |
| 268 | await self._push_log(entry) |
| 269 | |
| 270 | # Process ended |
| 271 | if self.status == "running": |
| 272 | exit_code = self.process.returncode if self.process else -1 |
| 273 | if exit_code == 0: |
| 274 | entry = self._create_log_entry("Crawler completed successfully", "success") |
| 275 | else: |
| 276 | entry = self._create_log_entry(f"Crawler exited with code: {exit_code}", "warning") |
| 277 | await self._push_log(entry) |
| 278 | self.status = "idle" |
| 279 | |
| 280 | except asyncio.CancelledError: |
| 281 | pass |
| 282 | except Exception as e: |
| 283 | entry = self._create_log_entry(f"Error reading output: {str(e)}", "error") |
| 284 | await self._push_log(entry) |
| 285 | |
| 286 | |
| 287 | # Global singleton |
no test coverage detected