Coverage for documents/management/commands/document_consumer.py: 0%

276 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 09:07 +0000

1""" 

2Document consumer management command. 

3 

4Watches a consumption directory for new documents and queues them for processing. 

5Uses watchfiles for efficient file system monitoring with support for both 

6native OS notifications and polling fallback. 

7""" 

8 

9from __future__ import annotations 

10 

11import logging 

12from dataclasses import dataclass 

13from pathlib import Path 

14from threading import Event 

15from time import monotonic 

16from typing import TYPE_CHECKING 

17from typing import Final 

18 

19from django import db 

20from django.conf import settings 

21from django.core.management.base import BaseCommand 

22from django.core.management.base import CommandError 

23from watchfiles import Change 

24from watchfiles import DefaultFilter 

25from watchfiles import watch 

26 

27from documents.data_models import ConsumableDocument 

28from documents.data_models import DocumentMetadataOverrides 

29from documents.data_models import DocumentSource 

30from documents.models import PaperlessTask 

31from documents.models import Tag 

32from documents.parsers import get_supported_file_extensions 

33from documents.tasks import consume_file 

34 

35if TYPE_CHECKING: 

36 from collections.abc import Iterator 

37 

38 

39logger = logging.getLogger("paperless.management.consumer") 

40 

41 

42@dataclass 

43class TrackedFile: 

44 """Represents a file being tracked for stability.""" 

45 

46 path: Path 

47 last_event_time: float 

48 last_mtime: float | None = None 

49 last_size: int | None = None 

50 

51 def update_stats(self) -> bool: 

52 """ 

53 Update file stats. Returns True if file exists and stats were updated. 

54 """ 

55 try: 

56 stat = self.path.stat() 

57 self.last_mtime = stat.st_mtime 

58 self.last_size = stat.st_size 

59 return True 

60 except OSError: 

61 return False 

62 

63 def is_unchanged(self) -> bool: 

64 """ 

65 Check if file stats match the previously recorded values. 

66 Returns False if file doesn't exist or stats changed. 

67 """ 

68 try: 

69 stat = self.path.stat() 

70 return stat.st_mtime == self.last_mtime and stat.st_size == self.last_size 

71 except OSError: 

72 return False 

73 

74 

75class FileStabilityTracker: 

76 """ 

77 Tracks file events and determines when files are stable for consumption. 

78 

79 A file is considered stable when: 

80 1. No new events have been received for it within the stability delay 

81 2. Its size and modification time haven't changed 

82 3. It still exists as a regular file 

83 

84 This handles various edge cases: 

85 - Network copies that write in chunks 

86 - Scanners that open/close files multiple times 

87 - Temporary files that get renamed 

88 - Files that are deleted before becoming stable 

89 """ 

90 

91 def __init__(self, stability_delay: float = 1.0) -> None: 

92 """ 

93 Initialize the tracker. 

94 

95 Args: 

96 stability_delay: Time in seconds a file must remain unchanged 

97 before being considered stable. 

98 """ 

99 self.stability_delay = stability_delay 

100 self._tracked: dict[Path, TrackedFile] = {} 

101 

102 def track(self, path: Path, change: Change) -> None: 

103 """ 

104 Register a file event. 

105 

106 Args: 

107 path: The file path that changed. 

108 change: The type of change (added, modified, deleted). 

109 """ 

110 path = path.resolve() 

111 

112 match change: 

113 case Change.deleted: 

114 self._tracked.pop(path, None) 

115 logger.debug(f"Stopped tracking deleted file: {path}") 

116 case Change.added | Change.modified: 

117 current_time = monotonic() 

118 if path in self._tracked: 

119 tracked = self._tracked[path] 

120 tracked.last_event_time = current_time 

121 tracked.update_stats() 

122 logger.debug(f"Updated tracking for: {path}") 

123 else: 

124 tracked = TrackedFile(path=path, last_event_time=current_time) 

125 if tracked.update_stats(): 

126 self._tracked[path] = tracked 

127 logger.debug(f"Started tracking: {path}") 

128 else: 

129 logger.debug(f"Could not stat file, not tracking: {path}") 

130 

131 def get_stable_files(self) -> Iterator[Path]: 

132 """ 

133 Yield files that have been stable for the configured delay. 

134 

135 Files are removed from tracking once yielded or determined to be invalid. 

136 """ 

137 current_time = monotonic() 

138 to_remove: list[Path] = [] 

139 to_yield: list[Path] = [] 

140 

141 for path, tracked in self._tracked.items(): 

142 time_since_event = current_time - tracked.last_event_time 

143 

144 if time_since_event < self.stability_delay: 

145 continue 

146 

147 # File has waited long enough, verify it's unchanged 

148 if not tracked.is_unchanged(): 

149 # Stats changed or file gone - update and wait again 

150 if tracked.update_stats(): 

151 tracked.last_event_time = current_time 

152 logger.debug(f"File changed during stability check: {path}") 

153 else: 

154 # File no longer exists, remove from tracking 

155 to_remove.append(path) 

156 logger.debug(f"File disappeared during stability check: {path}") 

157 continue 

158 

159 # Stable, but empty: some scanners create a zero byte placeholder 

160 # and only write the page some time later. Consuming it now can 

161 # only fail so drop it and let the writer's next event 

162 # (or the periodic rescan) bring it back once it has content 

163 if not tracked.last_size: 

164 to_remove.append(path) 

165 logger.debug("Ignoring stable but empty file: %s", path) 

166 continue 

167 

168 # File is stable, we can return it 

169 to_yield.append(path) 

170 logger.info(f"File is stable: {path}") 

171 

172 # Remove files that are no longer valid 

173 for path in to_remove: 

174 self._tracked.pop(path, None) 

175 

176 # Remove and yield stable files 

177 for path in to_yield: 

178 self._tracked.pop(path, None) 

179 yield path 

180 

181 def is_tracking(self, path: Path) -> bool: 

182 """Check whether a path is currently being tracked for stability.""" 

183 return path.resolve() in self._tracked 

184 

185 def has_pending_files(self) -> bool: 

186 """Check if there are files waiting for stability check.""" 

187 return len(self._tracked) > 0 

188 

189 @property 

190 def pending_count(self) -> int: 

191 """Number of files being tracked.""" 

192 return len(self._tracked) 

193 

194 

195class ConsumerFilter(DefaultFilter): 

196 """ 

197 Filter for watchfiles that accepts only supported document types 

198 and ignores system files/directories. 

199 

200 Extends DefaultFilter leveraging its built-in filtering: 

201 - `ignore_dirs`: Directory names to ignore (and all their contents) 

202 - `ignore_entity_patterns`: Regex patterns matched against filename/dirname only 

203 

204 We add custom logic for file extension filtering (only accept supported 

205 document types), which the library doesn't provide. 

206 """ 

207 

208 # Regex patterns for files to always ignore (matched against filename only) 

209 # These are passed to DefaultFilter.ignore_entity_patterns 

210 DEFAULT_IGNORE_PATTERNS: Final[tuple[str, ...]] = ( 

211 r"^\.DS_Store$", 

212 r"^\.DS_STORE$", 

213 r"^\._.*", 

214 r"^desktop\.ini$", 

215 r"^Thumbs\.db$", 

216 ) 

217 

218 # Directories to always ignore (passed to DefaultFilter.ignore_dirs) 

219 # These are matched by directory name, not full path 

220 DEFAULT_IGNORE_DIRS: Final[tuple[str, ...]] = ( 

221 ".stfolder", # Syncthing 

222 ".stversions", # Syncthing 

223 ".localized", # macOS 

224 "@eaDir", # Synology NAS 

225 ".Spotlight-V100", # macOS 

226 ".Trashes", # macOS 

227 "__MACOSX", # macOS archive artifacts 

228 ) 

229 

230 def __init__( 

231 self, 

232 *, 

233 supported_extensions: frozenset[str] | None = None, 

234 ignore_patterns: list[str] | None = None, 

235 ignore_dirs: list[str] | None = None, 

236 ) -> None: 

237 """ 

238 Initialize the consumer filter. 

239 

240 Args: 

241 supported_extensions: Set of file extensions to accept (e.g., {".pdf", ".png"}). 

242 If None, uses get_supported_file_extensions(). 

243 ignore_patterns: Additional regex patterns to ignore (matched against filename). 

244 ignore_dirs: Additional directory names to ignore (merged with defaults). 

245 """ 

246 # Get supported extensions 

247 if supported_extensions is None: 

248 supported_extensions = frozenset(get_supported_file_extensions()) 

249 self._supported_extensions = supported_extensions 

250 

251 # Combine default and user patterns 

252 all_patterns: list[str] = list(self.DEFAULT_IGNORE_PATTERNS) 

253 if ignore_patterns: 

254 all_patterns.extend(ignore_patterns) 

255 

256 # Combine default and user ignore_dirs 

257 all_ignore_dirs: list[str] = list(self.DEFAULT_IGNORE_DIRS) 

258 if ignore_dirs: 

259 all_ignore_dirs.extend(ignore_dirs) 

260 

261 # Let DefaultFilter handle all the pattern and directory filtering 

262 super().__init__( 

263 ignore_dirs=tuple(all_ignore_dirs), 

264 ignore_entity_patterns=tuple(all_patterns), 

265 ignore_paths=(), 

266 ) 

267 

268 def __call__(self, change: Change, path: str) -> bool: 

269 """ 

270 Filter function for watchfiles. 

271 

272 Returns True if the path should be watched, False to ignore. 

273 

274 The parent DefaultFilter handles: 

275 - Hidden files/directories (starting with .) 

276 - Directories in ignore_dirs 

277 - Files/directories matching ignore_entity_patterns 

278 

279 We additionally filter files by extension. 

280 """ 

281 # Let parent filter handle directory ignoring and pattern matching 

282 if not super().__call__(change, path): 

283 return False 

284 

285 path_obj = Path(path) 

286 

287 # For directories, parent filter already handled everything 

288 if path_obj.is_dir(): 

289 return True 

290 

291 # For files, check extension 

292 return self._has_supported_extension(path_obj) 

293 

294 def _has_supported_extension(self, path: Path) -> bool: 

295 """Check if the file has a supported extension.""" 

296 suffix = path.suffix.lower() 

297 return suffix in self._supported_extensions 

298 

299 

300def _tags_from_path(filepath: Path, consumption_dir: Path) -> list[int]: 

301 """ 

302 Walk up the directory tree from filepath to consumption_dir 

303 and get or create Tag IDs for every directory. 

304 

305 Returns list of Tag primary keys. 

306 """ 

307 db.close_old_connections() 

308 tag_ids: set[int] = set() 

309 path_parts = filepath.relative_to(consumption_dir).parent.parts 

310 

311 for part in path_parts: 

312 tag, _ = Tag.objects.get_or_create( 

313 name__iexact=part, 

314 defaults={"name": part}, 

315 ) 

316 tag_ids.add(tag.pk) 

317 

318 return list(tag_ids) 

319 

320 

321def _consume_file( 

322 filepath: Path, 

323 consumption_dir: Path, 

324 *, 

325 subdirs_as_tags: bool, 

326) -> bool: 

327 """ 

328 Queue a file for consumption. 

329 

330 Args: 

331 filepath: Path to the file to consume. 

332 consumption_dir: Base consumption directory. 

333 subdirs_as_tags: Whether to create tags from subdirectory names. 

334 

335 Returns: 

336 True if the file was successfully handed to Celery, False otherwise. 

337 Callers must not record the file as queued on failure, or the rescan 

338 will never retry it. 

339 """ 

340 # Verify file still exists and is accessible 

341 try: 

342 if not filepath.is_file(): 

343 logger.debug(f"Not consuming {filepath}: not a file or doesn't exist") 

344 return False 

345 except OSError as e: 

346 logger.warning(f"Not consuming {filepath}: {e}") 

347 return False 

348 

349 # Get tags from path if configured 

350 tag_ids: list[int] | None = None 

351 if subdirs_as_tags: 

352 try: 

353 tag_ids = _tags_from_path(filepath, consumption_dir) 

354 except Exception: 

355 logger.exception(f"Error creating tags from path for {filepath}") 

356 

357 # Queue for consumption 

358 try: 

359 logger.info(f"Adding {filepath} to the task queue") 

360 consume_file.apply_async( 

361 kwargs={ 

362 "input_doc": ConsumableDocument( 

363 source=DocumentSource.ConsumeFolder, 

364 original_file=filepath, 

365 ), 

366 "overrides": DocumentMetadataOverrides(tag_ids=tag_ids), 

367 }, 

368 headers={"trigger_source": PaperlessTask.TriggerSource.FOLDER_CONSUME}, 

369 ) 

370 except Exception: 

371 logger.exception(f"Error while queuing document {filepath}") 

372 return False 

373 

374 return True 

375 

376 

377class Command(BaseCommand): 

378 """ 

379 Watch a consumption directory and queue new documents for processing. 

380 

381 Uses watchfiles for efficient file system monitoring. Supports both 

382 native OS notifications (inotify on Linux, FSEvents on macOS) and 

383 polling for network filesystems. 

384 """ 

385 

386 help = "Watch the consumption directory for new documents" 

387 

388 # For testing - allows tests to stop the consumer 

389 stop_flag: Event = Event() 

390 

391 # Testing timeout in seconds 

392 testing_timeout_s: Final[float] = 0.5 

393 

394 # How often to perform a full-glob rescan of the consume directory as a 

395 # safety net. Each watchfiles watcher is torn down and recreated on every 

396 # batch to reconfigure its timeout, and a fresh watcher silently adopts the 

397 # current directory contents as its baseline. A file that appears between 

398 # one batch and the next watcher's baseline is therefore never reported and 

399 # would sit in the consume directory forever. This periodic rescan re-injects 

400 # such files into the stability tracker (see GH issue #13011). Not currently 

401 # user-configurable; instances may override for testing. 

402 rescan_interval_s: float = 300.0 

403 

404 def add_arguments(self, parser) -> None: 

405 parser.add_argument( 

406 "directory", 

407 default=None, 

408 nargs="?", 

409 help="The consumption directory (defaults to CONSUMPTION_DIR setting)", 

410 ) 

411 parser.add_argument( 

412 "--oneshot", 

413 action="store_true", 

414 help="Process existing files and exit without watching", 

415 ) 

416 parser.add_argument( 

417 "--testing", 

418 action="store_true", 

419 help="Enable testing mode with shorter timeouts", 

420 default=False, 

421 ) 

422 

423 def handle(self, *args, **options) -> None: 

424 # Resolve consumption directory 

425 directory = options.get("directory") 

426 if not directory: 

427 directory = getattr(settings, "CONSUMPTION_DIR", None) 

428 if not directory: 

429 raise CommandError("CONSUMPTION_DIR is not configured") 

430 

431 directory = Path(directory).resolve() 

432 

433 if not directory.exists(): 

434 raise CommandError(f"Consumption directory does not exist: {directory}") 

435 

436 if not directory.is_dir(): 

437 raise CommandError(f"Consumption path is not a directory: {directory}") 

438 

439 # Ensure scratch directory exists 

440 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True) 

441 

442 # Get settings 

443 recursive: bool = settings.CONSUMER_RECURSIVE 

444 subdirs_as_tags: bool = settings.CONSUMER_SUBDIRS_AS_TAGS 

445 polling_interval: float = settings.CONSUMER_POLLING_INTERVAL 

446 stability_delay: float = settings.CONSUMER_STABILITY_DELAY 

447 ignore_patterns: list[str] = settings.CONSUMER_IGNORE_PATTERNS 

448 ignore_dirs: list[str] = settings.CONSUMER_IGNORE_DIRS 

449 is_testing: bool = options.get("testing", False) 

450 is_oneshot: bool = options.get("oneshot", False) 

451 

452 # Create filter 

453 consumer_filter = ConsumerFilter( 

454 ignore_patterns=ignore_patterns, 

455 ignore_dirs=ignore_dirs, 

456 ) 

457 

458 # Process existing files 

459 queued = self._process_existing_files( 

460 directory=directory, 

461 recursive=recursive, 

462 subdirs_as_tags=subdirs_as_tags, 

463 consumer_filter=consumer_filter, 

464 ) 

465 

466 if is_oneshot: 

467 logger.info("Oneshot mode: processed existing files, exiting") 

468 return 

469 

470 # Start watching 

471 self._watch_directory( 

472 directory=directory, 

473 recursive=recursive, 

474 subdirs_as_tags=subdirs_as_tags, 

475 consumer_filter=consumer_filter, 

476 polling_interval=polling_interval, 

477 stability_delay=stability_delay, 

478 is_testing=is_testing, 

479 queued=queued, 

480 ) 

481 

482 logger.debug("Consumer exiting") 

483 

484 def _process_existing_files( 

485 self, 

486 *, 

487 directory: Path, 

488 recursive: bool, 

489 subdirs_as_tags: bool, 

490 consumer_filter: ConsumerFilter, 

491 ) -> set[Path]: 

492 """ 

493 Process any existing files in the consumption directory. 

494 

495 Returns the set of resolved paths that were queued, so the watch loop 

496 can seed its in-flight set and avoid re-queuing them on the first 

497 rescan before the consume tasks have removed them from disk. 

498 """ 

499 logger.info(f"Processing existing files in {directory}") 

500 

501 glob_pattern = "**/*" if recursive else "*" 

502 queued: set[Path] = set() 

503 

504 for filepath in directory.glob(glob_pattern): 

505 # Use filter to check if file should be processed 

506 if not filepath.is_file(): 

507 continue 

508 

509 if not consumer_filter(Change.added, str(filepath)): 

510 continue 

511 

512 if _consume_file( 

513 filepath=filepath, 

514 consumption_dir=directory, 

515 subdirs_as_tags=subdirs_as_tags, 

516 ): 

517 queued.add(filepath.resolve()) 

518 

519 return queued 

520 

521 def _rescan_existing_files( 

522 self, 

523 *, 

524 directory: Path, 

525 recursive: bool, 

526 consumer_filter: ConsumerFilter, 

527 tracker: FileStabilityTracker, 

528 queued: set[Path], 

529 ) -> None: 

530 """ 

531 Re-inject on-disk files the watcher never reported into the tracker. 

532 

533 Acts as a safety net for files stranded by the watcher-recreation gap 

534 (see ``rescan_interval_s``). Files already being tracked or already 

535 queued and awaiting consumption are skipped, so a file is never queued 

536 twice. Queued paths that have since left the directory are pruned so a 

537 later file reusing the same name is not skipped forever. 

538 """ 

539 # Prune in-flight paths that have left the directory 

540 for path in list(queued): 

541 if not path.exists(): 

542 queued.discard(path) 

543 

544 glob_pattern = "**/*" if recursive else "*" 

545 

546 for filepath in directory.glob(glob_pattern): 

547 if not filepath.is_file(): 

548 continue 

549 

550 if not consumer_filter(Change.added, str(filepath)): 

551 continue 

552 

553 resolved = filepath.resolve() 

554 if tracker.is_tracking(resolved) or resolved in queued: 

555 continue 

556 

557 logger.debug(f"Rescan found untracked file: {resolved}") 

558 tracker.track(resolved, Change.added) 

559 

560 def _watch_directory( 

561 self, 

562 *, 

563 directory: Path, 

564 recursive: bool, 

565 subdirs_as_tags: bool, 

566 consumer_filter: ConsumerFilter, 

567 polling_interval: float, 

568 stability_delay: float, 

569 is_testing: bool, 

570 queued: set[Path] | None = None, 

571 ) -> None: 

572 """Watch directory for changes and process stable files.""" 

573 use_polling = polling_interval > 0 

574 poll_delay_ms = int(polling_interval * 1000) if use_polling else 0 

575 

576 # Resolved paths that have been queued and are awaiting consumption. 

577 # Seeded from the startup scan so the first rescan does not re-queue 

578 # files whose consume tasks have not yet removed them from disk. 

579 queued = set() if queued is None else queued 

580 

581 # Full-glob safety net cadence (0 disables) 

582 rescan_interval_s = self.rescan_interval_s 

583 rescan_timeout_ms = ( 

584 int(rescan_interval_s * 1000) if rescan_interval_s > 0 else 0 

585 ) 

586 last_rescan = monotonic() 

587 

588 if use_polling: 

589 logger.info( 

590 f"Watching {directory} using polling (interval: {polling_interval}s)", 

591 ) 

592 else: 

593 logger.info(f"Watching {directory} using native file system events") 

594 

595 # Create stability tracker 

596 tracker = FileStabilityTracker(stability_delay=stability_delay) 

597 

598 # Calculate timeouts 

599 stability_timeout_ms = int(stability_delay * 1000) 

600 testing_timeout_ms = int(self.testing_timeout_s * 1000) 

601 

602 def cap_for_rescan(ms: int) -> int: 

603 """ 

604 Ensure the watch loop wakes often enough to run the rescan. 

605 

606 ``watch()`` blocks for up to ``rust_timeout``, so the rescan can 

607 only run that often. A timeout of 0 means "wait indefinitely", 

608 which would never wake to rescan; cap it at the rescan interval. 

609 """ 

610 if rescan_timeout_ms <= 0: 

611 return ms 

612 if ms <= 0: 

613 return rescan_timeout_ms 

614 return min(ms, rescan_timeout_ms) 

615 

616 # Calculate appropriate timeout for watch loop 

617 # In polling mode, rust_timeout must be significantly longer than poll_delay_ms 

618 # to ensure poll cycles can complete before timing out 

619 if is_testing: 

620 if use_polling: 

621 # For polling: timeout must be at least 3x the poll interval to allow 

622 # multiple poll cycles. This prevents timeouts from interfering with 

623 # the polling mechanism. 

624 min_polling_timeout_ms = poll_delay_ms * 3 

625 timeout_ms = max(min_polling_timeout_ms, testing_timeout_ms) 

626 else: 

627 # For native watching, use short timeout to check stop flag 

628 timeout_ms = testing_timeout_ms 

629 else: 

630 # Not testing, wait indefinitely for first event 

631 timeout_ms = 0 

632 

633 timeout_ms = cap_for_rescan(timeout_ms) 

634 

635 self.stop_flag.clear() 

636 

637 while not self.stop_flag.is_set(): 

638 try: 

639 for changes in watch( 

640 directory, 

641 watch_filter=consumer_filter, 

642 rust_timeout=timeout_ms, 

643 yield_on_timeout=True, 

644 force_polling=use_polling, 

645 poll_delay_ms=poll_delay_ms, 

646 recursive=recursive, 

647 stop_event=self.stop_flag, 

648 ): 

649 # Process each change 

650 for change_type, path in changes: 

651 path = Path(path).resolve() 

652 if change_type == Change.deleted: 

653 # Consumed (or otherwise removed); a later file 

654 # reusing this name must not be skipped as 

655 # already-queued. 

656 queued.discard(path) 

657 if not path.is_file(): 

658 continue 

659 if path in queued: 

660 # Already queued and awaiting consumption; a stray 

661 # event (NAS metadata touch, AV scan, etc.) while 

662 # the file sits on disk mid-consumption must not 

663 # cause it to be queued a second time (GH #13511). 

664 logger.debug(f"Ignoring event for queued file: {path}") 

665 continue 

666 logger.debug(f"Event: {change_type.name} for {path}") 

667 tracker.track(path, change_type) 

668 

669 # Check for stable files 

670 for stable_path in tracker.get_stable_files(): 

671 # Only remember files that were actually queued, so the 

672 # rescan does not re-queue them while the consume task 

673 # has yet to remove them from disk, but does retry a 

674 # failed publish instead of stranding it 

675 if _consume_file( 

676 filepath=stable_path, 

677 consumption_dir=directory, 

678 subdirs_as_tags=subdirs_as_tags, 

679 ): 

680 queued.add(stable_path) 

681 

682 # Exit watch loop to reconfigure timeout 

683 break 

684 

685 # Periodic full-glob safety net for files the watcher missed 

686 if rescan_timeout_ms > 0 and ( 

687 monotonic() - last_rescan >= rescan_interval_s 

688 ): 

689 self._rescan_existing_files( 

690 directory=directory, 

691 recursive=recursive, 

692 consumer_filter=consumer_filter, 

693 tracker=tracker, 

694 queued=queued, 

695 ) 

696 last_rescan = monotonic() 

697 

698 # Determine next timeout 

699 if tracker.has_pending_files(): 

700 # Check pending files at stability interval 

701 timeout_ms = stability_timeout_ms 

702 elif is_testing: 

703 # In testing, use appropriate timeout based on watch mode 

704 if use_polling: 

705 # For polling: ensure timeout allows polls to complete 

706 min_polling_timeout_ms = poll_delay_ms * 3 

707 timeout_ms = max(min_polling_timeout_ms, testing_timeout_ms) 

708 else: 

709 # For native watching, use short timeout to check stop flag 

710 timeout_ms = testing_timeout_ms 

711 else: # pragma: nocover 

712 # No pending files, wait indefinitely 

713 timeout_ms = 0 

714 

715 timeout_ms = cap_for_rescan(timeout_ms) 

716 

717 except KeyboardInterrupt: # pragma: nocover 

718 logger.info("Received interrupt, stopping consumer") 

719 self.stop_flag.set()