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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
1"""
2Document consumer management command.
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"""
9from __future__ import annotations
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
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
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
35if TYPE_CHECKING:
36 from collections.abc import Iterator
39logger = logging.getLogger("paperless.management.consumer")
42@dataclass
43class TrackedFile:
44 """Represents a file being tracked for stability."""
46 path: Path
47 last_event_time: float
48 last_mtime: float | None = None
49 last_size: int | None = None
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
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
75class FileStabilityTracker:
76 """
77 Tracks file events and determines when files are stable for consumption.
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
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 """
91 def __init__(self, stability_delay: float = 1.0) -> None:
92 """
93 Initialize the tracker.
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] = {}
102 def track(self, path: Path, change: Change) -> None:
103 """
104 Register a file event.
106 Args:
107 path: The file path that changed.
108 change: The type of change (added, modified, deleted).
109 """
110 path = path.resolve()
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}")
131 def get_stable_files(self) -> Iterator[Path]:
132 """
133 Yield files that have been stable for the configured delay.
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] = []
141 for path, tracked in self._tracked.items():
142 time_since_event = current_time - tracked.last_event_time
144 if time_since_event < self.stability_delay:
145 continue
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
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
168 # File is stable, we can return it
169 to_yield.append(path)
170 logger.info(f"File is stable: {path}")
172 # Remove files that are no longer valid
173 for path in to_remove:
174 self._tracked.pop(path, None)
176 # Remove and yield stable files
177 for path in to_yield:
178 self._tracked.pop(path, None)
179 yield path
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
185 def has_pending_files(self) -> bool:
186 """Check if there are files waiting for stability check."""
187 return len(self._tracked) > 0
189 @property
190 def pending_count(self) -> int:
191 """Number of files being tracked."""
192 return len(self._tracked)
195class ConsumerFilter(DefaultFilter):
196 """
197 Filter for watchfiles that accepts only supported document types
198 and ignores system files/directories.
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
204 We add custom logic for file extension filtering (only accept supported
205 document types), which the library doesn't provide.
206 """
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 )
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 )
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.
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
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)
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)
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 )
268 def __call__(self, change: Change, path: str) -> bool:
269 """
270 Filter function for watchfiles.
272 Returns True if the path should be watched, False to ignore.
274 The parent DefaultFilter handles:
275 - Hidden files/directories (starting with .)
276 - Directories in ignore_dirs
277 - Files/directories matching ignore_entity_patterns
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
285 path_obj = Path(path)
287 # For directories, parent filter already handled everything
288 if path_obj.is_dir():
289 return True
291 # For files, check extension
292 return self._has_supported_extension(path_obj)
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
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.
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
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)
318 return list(tag_ids)
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.
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.
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
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}")
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
374 return True
377class Command(BaseCommand):
378 """
379 Watch a consumption directory and queue new documents for processing.
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 """
386 help = "Watch the consumption directory for new documents"
388 # For testing - allows tests to stop the consumer
389 stop_flag: Event = Event()
391 # Testing timeout in seconds
392 testing_timeout_s: Final[float] = 0.5
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
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 )
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")
431 directory = Path(directory).resolve()
433 if not directory.exists():
434 raise CommandError(f"Consumption directory does not exist: {directory}")
436 if not directory.is_dir():
437 raise CommandError(f"Consumption path is not a directory: {directory}")
439 # Ensure scratch directory exists
440 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True)
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)
452 # Create filter
453 consumer_filter = ConsumerFilter(
454 ignore_patterns=ignore_patterns,
455 ignore_dirs=ignore_dirs,
456 )
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 )
466 if is_oneshot:
467 logger.info("Oneshot mode: processed existing files, exiting")
468 return
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 )
482 logger.debug("Consumer exiting")
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.
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}")
501 glob_pattern = "**/*" if recursive else "*"
502 queued: set[Path] = set()
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
509 if not consumer_filter(Change.added, str(filepath)):
510 continue
512 if _consume_file(
513 filepath=filepath,
514 consumption_dir=directory,
515 subdirs_as_tags=subdirs_as_tags,
516 ):
517 queued.add(filepath.resolve())
519 return queued
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.
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)
544 glob_pattern = "**/*" if recursive else "*"
546 for filepath in directory.glob(glob_pattern):
547 if not filepath.is_file():
548 continue
550 if not consumer_filter(Change.added, str(filepath)):
551 continue
553 resolved = filepath.resolve()
554 if tracker.is_tracking(resolved) or resolved in queued:
555 continue
557 logger.debug(f"Rescan found untracked file: {resolved}")
558 tracker.track(resolved, Change.added)
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
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
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()
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")
595 # Create stability tracker
596 tracker = FileStabilityTracker(stability_delay=stability_delay)
598 # Calculate timeouts
599 stability_timeout_ms = int(stability_delay * 1000)
600 testing_timeout_ms = int(self.testing_timeout_s * 1000)
602 def cap_for_rescan(ms: int) -> int:
603 """
604 Ensure the watch loop wakes often enough to run the rescan.
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)
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
633 timeout_ms = cap_for_rescan(timeout_ms)
635 self.stop_flag.clear()
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)
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)
682 # Exit watch loop to reconfigure timeout
683 break
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()
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
715 timeout_ms = cap_for_rescan(timeout_ms)
717 except KeyboardInterrupt: # pragma: nocover
718 logger.info("Received interrupt, stopping consumer")
719 self.stop_flag.set()