diff --git a/.env.demo b/.env.demo index 8e58c068..5a486aa4 100644 --- a/.env.demo +++ b/.env.demo @@ -206,7 +206,30 @@ EMAIL_USE_TLS=True EMAIL_SENDER=DocuElevate System EMAIL_DEFAULT_RECIPIENT=recipient@example.com +# **Watch Folder Ingestion** +# DocuElevate can automatically monitor directories (local filesystem, FTP, SFTP) for new files. +# +# Local watch folders — works with any mounted path (SMB/CIFS, NFS, local disk, etc.) +# Set WATCH_FOLDERS to a comma-separated list of absolute paths inside the container. +WATCH_FOLDERS= +WATCH_FOLDER_POLL_INTERVAL=1 +WATCH_FOLDER_DELETE_AFTER_PROCESS=false + +# FTP ingest — poll an FTP directory for new files (uses FTP connection settings above) +FTP_INGEST_ENABLED=false +FTP_INGEST_FOLDER= +FTP_INGEST_DELETE_AFTER_PROCESS=false + +# SFTP ingest — poll an SFTP directory for new files (uses SFTP connection settings above) +SFTP_INGEST_ENABLED=false +SFTP_INGEST_FOLDER= +SFTP_INGEST_DELETE_AFTER_PROCESS=false + # **IMAP Settings** +# DocuElevate polls these mailboxes for new email attachments and automatically ingests them. +# No manual forwarding required — DocuElevate acts as an IMAP *client*. +# For HP Scanners / Scan-to-Email: configure the scanner to send to a dedicated mailbox, +# then point DocuElevate at that mailbox using the settings below. IMAP1_HOST=mail.example.com IMAP1_PORT=993 IMAP1_USERNAME= diff --git a/app/celery_worker.py b/app/celery_worker.py index 509f6967..934ef132 100644 --- a/app/celery_worker.py +++ b/app/celery_worker.py @@ -17,6 +17,7 @@ from app.tasks.extract_metadata_with_gpt import extract_metadata_with_gpt # noq from app.tasks.finalize_document_storage import finalize_document_storage # noqa: F401 from app.tasks.imap_tasks import pull_all_inboxes # noqa: F401 from app.tasks.monitor_stalled_steps import monitor_stalled_steps # noqa: F401 +from app.tasks.watch_folder_tasks import scan_all_watch_folders # noqa: F401 # **Ensure all tasks are imported before Celery starts** from app.tasks.process_document import process_document # noqa: F401 @@ -97,6 +98,21 @@ celery.conf.beat_schedule = { "schedule": crontab(minute="*/1"), # Every minute "options": {"expires": 55}, # Must complete within 55 seconds }, + # Watch folder scanning — polls local paths, FTP, and SFTP ingest folders. + # Schedule is controlled by WATCH_FOLDER_POLL_INTERVAL (default: 1 minute). + "scan-watch-folders": ( + { + "task": "app.tasks.watch_folder_tasks.scan_all_watch_folders", + "schedule": crontab(minute=f"*/{max(1, settings.watch_folder_poll_interval)}"), + "options": {"expires": 55}, + } + if ( + settings.watch_folders + or settings.ftp_ingest_enabled + or settings.sftp_ingest_enabled + ) + else None + ), # Backfill embeddings for files that were processed before the # embedding pipeline was enabled, or where the embedding task failed. "backfill-missing-embeddings": { diff --git a/app/config.py b/app/config.py index 848846a6..3680f633 100644 --- a/app/config.py +++ b/app/config.py @@ -190,6 +190,83 @@ class Settings(BaseSettings): stripe_success_url: Optional[str] = None # e.g. https://app.example.com/billing/success stripe_cancel_url: Optional[str] = None # e.g. https://app.example.com/pricing + # --------------------------------------------------------------------------- + # Watch Folder Ingestion + # --------------------------------------------------------------------------- + # Local filesystem watch folders (comma-separated list of absolute paths). + # DocuElevate will poll each path for new files and automatically ingest them. + # Works with any mounted path, including SMB/CIFS (via system mount), NFS, etc. + # Example: /watchfolders/scanner,/mnt/shared/inbox + watch_folders: Optional[str] = Field( + default=None, + description=( + "Comma-separated list of local filesystem paths (absolute) that DocuElevate will " + "poll for new files to ingest. Each file found is enqueued for document processing. " + "Works with any mounted path including SMB/CIFS (mounted via system) and NFS. " + "Example: /watchfolders/scanner,/mnt/shared/inbox" + ), + ) + watch_folder_poll_interval: int = Field( + default=1, + description=( + "Poll interval in minutes for local watch folder scanning. Default: 1 minute." + ), + ) + watch_folder_delete_after_process: bool = Field( + default=False, + description=( + "Delete files from local watch folders after they have been successfully enqueued " + "for processing. When False (default), files are left in place and tracked via a " + "cache file to avoid re-ingesting them." + ), + ) + + # FTP Ingest / Watch Folder + # Uses the existing FTP credentials (ftp_host, ftp_username, ftp_password) to poll + # a source folder on the FTP server for new files to ingest. + ftp_ingest_folder: Optional[str] = Field( + default=None, + description=( + "FTP folder path to monitor for new files to ingest. " + "Uses the existing FTP connection settings (FTP_HOST, FTP_USERNAME, FTP_PASSWORD). " + "When set, DocuElevate will periodically poll this folder and download new files for processing." + ), + ) + ftp_ingest_enabled: bool = Field( + default=False, + description="Enable FTP watch folder ingestion. Requires FTP_INGEST_FOLDER and FTP connection settings.", + ) + ftp_ingest_delete_after_process: bool = Field( + default=False, + description=( + "Delete files from the FTP ingest folder after they have been successfully downloaded " + "and enqueued for processing. Default: False (files are left in place)." + ), + ) + + # SFTP Ingest / Watch Folder + # Uses the existing SFTP credentials (sftp_host, sftp_username, sftp_password/sftp_private_key) + # to poll a source folder on the SFTP server for new files to ingest. + sftp_ingest_folder: Optional[str] = Field( + default=None, + description=( + "SFTP folder path to monitor for new files to ingest. " + "Uses the existing SFTP connection settings (SFTP_HOST, SFTP_USERNAME, SFTP_PASSWORD/SFTP_PRIVATE_KEY). " + "When set, DocuElevate will periodically poll this folder and download new files for processing." + ), + ) + sftp_ingest_enabled: bool = Field( + default=False, + description="Enable SFTP watch folder ingestion. Requires SFTP_INGEST_FOLDER and SFTP connection settings.", + ) + sftp_ingest_delete_after_process: bool = Field( + default=False, + description=( + "Delete files from the SFTP ingest folder after they have been successfully downloaded " + "and enqueued for processing. Default: False (files are left in place)." + ), + ) + # IMAP 1 imap1_host: Optional[str] = None imap1_port: Optional[int] = 993 diff --git a/app/tasks/watch_folder_tasks.py b/app/tasks/watch_folder_tasks.py new file mode 100644 index 00000000..a699484c --- /dev/null +++ b/app/tasks/watch_folder_tasks.py @@ -0,0 +1,566 @@ +#!/usr/bin/env python3 +""" +Watch Folder Ingestion Tasks + +Periodically scans configured directories (local filesystem, FTP, SFTP) for new files +and enqueues them for document processing. + +Supported sources: +- Local filesystem paths (works with any mounted path: SMB/CIFS, NFS, etc.) +- FTP server directories (uses existing FTP connection settings) +- SFTP server directories (uses existing SFTP connection settings) +""" + +import ftplib # nosec B402 - FTP usage is intentional for legacy server support +import json +import logging +import os +import tempfile +from datetime import datetime, timedelta, timezone + +import redis +from celery import shared_task + +from app.config import settings +from app.tasks.convert_to_pdf import convert_to_pdf +from app.tasks.process_document import process_document +from app.utils.allowed_types import ALLOWED_EXTENSIONS, ALLOWED_MIME_TYPES + +logger = logging.getLogger(__name__) + +redis_client = redis.StrictRedis.from_url(settings.redis_url, decode_responses=True) + +WATCH_FOLDER_LOCK_KEY = "watch_folder_lock" +WATCH_FOLDER_LOCK_EXPIRE = 300 # 5 minutes + +# Cache file for tracking already-ingested files (local watch folders) +WATCH_FOLDER_CACHE_FILE = os.path.join(settings.workdir, "watch_folder_processed.json") +# Cache file for tracking already-ingested files (FTP watch folder) +FTP_INGEST_CACHE_FILE = os.path.join(settings.workdir, "ftp_ingest_processed.json") +# Cache file for tracking already-ingested files (SFTP watch folder) +SFTP_INGEST_CACHE_FILE = os.path.join(settings.workdir, "sftp_ingest_processed.json") + +_CACHE_RETENTION_DAYS = 30 + + +# --------------------------------------------------------------------------- +# Locking helpers +# --------------------------------------------------------------------------- + + +def _acquire_lock(lock_key: str, expire: int = WATCH_FOLDER_LOCK_EXPIRE) -> bool: + """Acquire a Redis-based distributed lock. Returns True if acquired.""" + acquired = redis_client.setnx(lock_key, "locked") + if acquired: + redis_client.expire(lock_key, expire) + logger.debug("Lock acquired: %s", lock_key) + return True + logger.debug("Lock already held: %s — skipping.", lock_key) + return False + + +def _release_lock(lock_key: str) -> None: + """Release a Redis-based distributed lock.""" + redis_client.delete(lock_key) + logger.debug("Lock released: %s", lock_key) + + +# --------------------------------------------------------------------------- +# Cache helpers +# --------------------------------------------------------------------------- + + +def _load_cache(cache_file: str) -> dict[str, str]: + """Load the set of already-processed file identifiers from a JSON cache file.""" + if os.path.exists(cache_file): + try: + with open(cache_file) as f: + data: dict[str, str] = json.load(f) + return _evict_old_entries(data) + except (json.JSONDecodeError, OSError): + logger.warning("Failed to read cache file %s — starting fresh.", cache_file) + return {} + + +def _save_cache(cache_file: str, data: dict[str, str]) -> None: + """Persist the processed-file cache to disk.""" + try: + with open(cache_file, "w") as f: + json.dump(data, f, indent=2) + except OSError as exc: + logger.error("Failed to write cache file %s: %s", cache_file, exc) + + +def _evict_old_entries(data: dict[str, str]) -> dict[str, str]: + """Remove entries older than CACHE_RETENTION_DAYS to prevent unbounded growth.""" + cutoff = datetime.now(timezone.utc) - timedelta(days=_CACHE_RETENTION_DAYS) + result = {} + for key, date_str in data.items(): + try: + dt = datetime.fromisoformat(date_str) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + if dt > cutoff: + result[key] = date_str + except (ValueError, TypeError): + pass # Skip entries with malformed timestamps + return result + + +def _mark_processed(cache: dict[str, str], key: str) -> None: + """Add a file identifier to the in-memory cache dict.""" + cache[key] = datetime.now(timezone.utc).isoformat() + + +# --------------------------------------------------------------------------- +# File type helpers +# --------------------------------------------------------------------------- + + +def _is_allowed_file(filename: str) -> bool: + """Return True if the file should be ingested based on its name/extension.""" + ext = os.path.splitext(filename)[1].lower() + return ext in ALLOWED_EXTENSIONS or filename.lower().endswith(".pdf") + + +def _enqueue_file(file_path: str, *, filename: str | None = None) -> None: + """Enqueue a local file path for document processing.""" + fname = filename or os.path.basename(file_path) + _, ext = os.path.splitext(fname) + mime_check = ext.lower() in {".pdf"} + + if mime_check or fname.lower().endswith(".pdf"): + process_document.delay(file_path) + logger.info("Enqueued for processing: %s", fname) + else: + convert_to_pdf.delay(file_path) + logger.info("Enqueued for PDF conversion: %s", fname) + + +# --------------------------------------------------------------------------- +# Local filesystem watch folder scanning +# --------------------------------------------------------------------------- + + +def _scan_local_folder(folder_path: str, cache: dict[str, str], delete_after: bool) -> int: + """ + Scan a single local directory for new, allowed files. + + Files already present in *cache* (keyed by their absolute path) are skipped. + Returns the number of files newly enqueued. + """ + if not os.path.isdir(folder_path): + logger.warning("Watch folder does not exist or is not a directory: %s", folder_path) + return 0 + + count = 0 + try: + entries = os.scandir(folder_path) + except PermissionError as exc: + logger.error("Cannot scan watch folder %s: %s", folder_path, exc) + return 0 + + for entry in entries: + if not entry.is_file(follow_symlinks=True): + continue + if not _is_allowed_file(entry.name): + logger.debug("Skipping unsupported file type: %s", entry.name) + continue + + abs_path = entry.path + if abs_path in cache: + logger.debug("Already processed: %s", abs_path) + continue + + # Copy file to workdir before enqueueing so the original isn't locked + dest_filename = f"wf_{entry.name}" + dest_path = os.path.join(settings.workdir, dest_filename) + + # Avoid overwriting if a file with the same name is already there + if os.path.exists(dest_path): + base, ext2 = os.path.splitext(dest_filename) + dest_path = os.path.join(settings.workdir, f"{base}_{int(datetime.now().timestamp())}{ext2}") + + try: + import shutil + + shutil.copy2(abs_path, dest_path) + except OSError as exc: + logger.error("Failed to copy %s to workdir: %s", abs_path, exc) + continue + + _enqueue_file(dest_path) + _mark_processed(cache, abs_path) + count += 1 + + if delete_after: + try: + os.remove(abs_path) + logger.info("Deleted source file after ingestion: %s", abs_path) + except OSError as exc: + logger.warning("Could not delete source file %s: %s", abs_path, exc) + + return count + + +# --------------------------------------------------------------------------- +# FTP watch folder scanning +# --------------------------------------------------------------------------- + + +def _connect_ftp() -> ftplib.FTP | None: + """Establish an FTP/FTPS connection using the global FTP settings.""" + host = getattr(settings, "ftp_host", None) + port = getattr(settings, "ftp_port", 21) or 21 + username = getattr(settings, "ftp_username", None) + password = getattr(settings, "ftp_password", None) + + if not (host and username and password): + logger.warning("FTP ingest: connection settings incomplete — skipping.") + return None + + use_tls: bool = getattr(settings, "ftp_use_tls", True) + allow_plaintext: bool = getattr(settings, "ftp_allow_plaintext", True) + + if use_tls: + try: + ftp = ftplib.FTP_TLS() # noqa: S321 + ftp.connect(host=host, port=port) + ftp.login(user=username, passwd=password) + ftp.prot_p() + logger.debug("FTP ingest: connected via FTPS to %s:%s", host, port) + return ftp + except Exception as exc: + if not allow_plaintext: + logger.error("FTP ingest: FTPS failed and plaintext not allowed: %s", exc) + return None + logger.warning("FTP ingest: FTPS failed, falling back to plain FTP: %s", exc) + + try: + ftp = ftplib.FTP() # nosec B321 # noqa: S321 + ftp.connect(host=host, port=port) + ftp.login(user=username, passwd=password) + logger.debug("FTP ingest: connected via plain FTP to %s:%s", host, port) + return ftp + except Exception as exc: + logger.error("FTP ingest: connection failed: %s", exc) + return None + + +def _scan_ftp_folder(ftp: ftplib.FTP, remote_folder: str, cache: dict[str, str], delete_after: bool) -> int: + """ + List *remote_folder* on the FTP server and download new allowed files to workdir. + Returns the number of files newly enqueued. + """ + count = 0 + try: + ftp.cwd(remote_folder) + except ftplib.error_perm as exc: + logger.error("FTP ingest: cannot CWD to %s: %s", remote_folder, exc) + return 0 + + try: + filenames: list[str] = ftp.nlst() + except Exception as exc: + logger.error("FTP ingest: NLST failed on %s: %s", remote_folder, exc) + return 0 + + for filename in filenames: + if not _is_allowed_file(filename): + logger.debug("FTP ingest: skipping %s (unsupported type)", filename) + continue + + cache_key = f"ftp:{remote_folder}/{filename}" + if cache_key in cache: + logger.debug("FTP ingest: already processed %s", filename) + continue + + dest_path = os.path.join(settings.workdir, f"ftp_{filename}") + if os.path.exists(dest_path): + base, ext2 = os.path.splitext(f"ftp_{filename}") + dest_path = os.path.join(settings.workdir, f"{base}_{int(datetime.now().timestamp())}{ext2}") + + try: + with open(dest_path, "wb") as local_file: + ftp.retrbinary(f"RETR {filename}", local_file.write) + logger.info("FTP ingest: downloaded %s to %s", filename, dest_path) + except Exception as exc: + logger.error("FTP ingest: failed to download %s: %s", filename, exc) + if os.path.exists(dest_path): + os.remove(dest_path) + continue + + _enqueue_file(dest_path, filename=filename) + _mark_processed(cache, cache_key) + count += 1 + + if delete_after: + try: + ftp.delete(filename) + logger.info("FTP ingest: deleted remote file %s after ingestion", filename) + except Exception as exc: + logger.warning("FTP ingest: could not delete remote file %s: %s", filename, exc) + + return count + + +# --------------------------------------------------------------------------- +# SFTP watch folder scanning +# --------------------------------------------------------------------------- + + +def _get_sftp_connection(): + """ + Establish an SFTP connection using the global SFTP settings. + Returns (ssh_client, sftp_client) tuple, or (None, None) on failure. + """ + import paramiko + + host = getattr(settings, "sftp_host", None) + port = getattr(settings, "sftp_port", 22) or 22 + username = getattr(settings, "sftp_username", None) + + if not (host and username): + logger.warning("SFTP ingest: connection settings incomplete — skipping.") + return None, None + + ssh = paramiko.SSHClient() + + if getattr(settings, "sftp_disable_host_key_verification", False): + logger.warning( + "SFTP ingest: host key verification is DISABLED — connections are vulnerable to MITM attacks. " + "Set SFTP_DISABLE_HOST_KEY_VERIFICATION=False for production use." + ) + ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy()) # nosec B507 # noqa: S507 + else: + ssh.load_system_host_keys() + ssh.set_missing_host_key_policy(paramiko.RejectPolicy()) + + connect_kwargs: dict = {"hostname": host, "port": port, "username": username} + + sftp_key_path = getattr(settings, "sftp_private_key", None) + sftp_key_passphrase = getattr(settings, "sftp_private_key_passphrase", None) + sftp_password = getattr(settings, "sftp_password", None) + + if sftp_key_path and os.path.exists(sftp_key_path): + connect_kwargs["key_filename"] = sftp_key_path + if sftp_key_passphrase: + connect_kwargs["passphrase"] = sftp_key_passphrase + elif sftp_password: + connect_kwargs["password"] = sftp_password + else: + logger.warning("SFTP ingest: no password or private key configured — skipping.") + return None, None + + try: + ssh.connect(**connect_kwargs) + sftp = ssh.open_sftp() + logger.debug("SFTP ingest: connected to %s:%s", host, port) + return ssh, sftp + except Exception as exc: + logger.error("SFTP ingest: connection failed: %s", exc) + try: + ssh.close() + except Exception: + pass + return None, None + + +def _scan_sftp_folder(sftp, remote_folder: str, cache: dict[str, str], delete_after: bool) -> int: + """ + List *remote_folder* on the SFTP server and download new allowed files to workdir. + Returns the number of files newly enqueued. + """ + import stat as stat_module + + count = 0 + try: + entries = sftp.listdir_attr(remote_folder) + except Exception as exc: + logger.error("SFTP ingest: cannot list %s: %s", remote_folder, exc) + return 0 + + for attr in entries: + # Skip directories + if stat_module.S_ISDIR(attr.st_mode or 0): + continue + + filename = attr.filename + if not _is_allowed_file(filename): + logger.debug("SFTP ingest: skipping %s (unsupported type)", filename) + continue + + remote_path = f"{remote_folder}/{filename}".replace("//", "/") + cache_key = f"sftp:{remote_path}" + if cache_key in cache: + logger.debug("SFTP ingest: already processed %s", filename) + continue + + dest_path = os.path.join(settings.workdir, f"sftp_{filename}") + if os.path.exists(dest_path): + base, ext2 = os.path.splitext(f"sftp_{filename}") + dest_path = os.path.join(settings.workdir, f"{base}_{int(datetime.now().timestamp())}{ext2}") + + try: + sftp.get(remote_path, dest_path) + logger.info("SFTP ingest: downloaded %s to %s", remote_path, dest_path) + except Exception as exc: + logger.error("SFTP ingest: failed to download %s: %s", remote_path, exc) + if os.path.exists(dest_path): + os.remove(dest_path) + continue + + _enqueue_file(dest_path, filename=filename) + _mark_processed(cache, cache_key) + count += 1 + + if delete_after: + try: + sftp.remove(remote_path) + logger.info("SFTP ingest: deleted remote file %s after ingestion", remote_path) + except Exception as exc: + logger.warning("SFTP ingest: could not delete remote file %s: %s", remote_path, exc) + + return count + + +# --------------------------------------------------------------------------- +# Main Celery tasks +# --------------------------------------------------------------------------- + + +@shared_task +def scan_local_watch_folders() -> dict: + """ + Celery task: scan all configured local filesystem watch folders for new files. + + Reads the WATCH_FOLDERS setting (comma-separated list of absolute directory paths) + and enqueues any new allowed files for document processing. Already-processed files + are tracked in a JSON cache file to prevent duplicate ingestion. + """ + watch_folders_raw: str | None = getattr(settings, "watch_folders", None) + if not watch_folders_raw: + logger.debug("No local watch folders configured — nothing to scan.") + return {"status": "skipped", "reason": "WATCH_FOLDERS not configured"} + + folder_paths = [p.strip() for p in watch_folders_raw.split(",") if p.strip()] + if not folder_paths: + return {"status": "skipped", "reason": "WATCH_FOLDERS is empty"} + + delete_after: bool = getattr(settings, "watch_folder_delete_after_process", False) + cache = _load_cache(WATCH_FOLDER_CACHE_FILE) + total = 0 + + for folder_path in folder_paths: + logger.info("Scanning local watch folder: %s", folder_path) + n = _scan_local_folder(folder_path, cache, delete_after) + logger.info("Local watch folder %s: %d new file(s) enqueued.", folder_path, n) + total += n + + _save_cache(WATCH_FOLDER_CACHE_FILE, cache) + return {"status": "ok", "files_enqueued": total, "folders_scanned": len(folder_paths)} + + +@shared_task +def scan_ftp_watch_folder() -> dict: + """ + Celery task: scan the configured FTP ingest folder for new files. + + Uses the existing FTP connection settings and FTP_INGEST_FOLDER to poll for + new documents to ingest. Downloaded files are enqueued for processing and + (optionally) deleted from the server after download. + """ + if not getattr(settings, "ftp_ingest_enabled", False): + logger.debug("FTP ingest is disabled — skipping.") + return {"status": "skipped", "reason": "FTP_INGEST_ENABLED is False"} + + ingest_folder: str | None = getattr(settings, "ftp_ingest_folder", None) + if not ingest_folder: + logger.warning("FTP ingest enabled but FTP_INGEST_FOLDER is not set — skipping.") + return {"status": "skipped", "reason": "FTP_INGEST_FOLDER not configured"} + + delete_after: bool = getattr(settings, "ftp_ingest_delete_after_process", False) + cache = _load_cache(FTP_INGEST_CACHE_FILE) + + ftp = _connect_ftp() + if ftp is None: + return {"status": "error", "reason": "FTP connection failed"} + + try: + n = _scan_ftp_folder(ftp, ingest_folder, cache, delete_after) + finally: + try: + ftp.quit() + except Exception: + pass + + _save_cache(FTP_INGEST_CACHE_FILE, cache) + logger.info("FTP ingest: %d new file(s) enqueued from %s.", n, ingest_folder) + return {"status": "ok", "files_enqueued": n, "folder": ingest_folder} + + +@shared_task +def scan_sftp_watch_folder() -> dict: + """ + Celery task: scan the configured SFTP ingest folder for new files. + + Uses the existing SFTP connection settings and SFTP_INGEST_FOLDER to poll for + new documents to ingest. Downloaded files are enqueued for processing and + (optionally) deleted from the server after download. + """ + if not getattr(settings, "sftp_ingest_enabled", False): + logger.debug("SFTP ingest is disabled — skipping.") + return {"status": "skipped", "reason": "SFTP_INGEST_ENABLED is False"} + + ingest_folder: str | None = getattr(settings, "sftp_ingest_folder", None) + if not ingest_folder: + logger.warning("SFTP ingest enabled but SFTP_INGEST_FOLDER is not set — skipping.") + return {"status": "skipped", "reason": "SFTP_INGEST_FOLDER not configured"} + + delete_after: bool = getattr(settings, "sftp_ingest_delete_after_process", False) + cache = _load_cache(SFTP_INGEST_CACHE_FILE) + + ssh, sftp = _get_sftp_connection() + if sftp is None: + return {"status": "error", "reason": "SFTP connection failed"} + + try: + n = _scan_sftp_folder(sftp, ingest_folder, cache, delete_after) + finally: + try: + sftp.close() + except Exception: + pass + try: + ssh.close() + except Exception: + pass + + _save_cache(SFTP_INGEST_CACHE_FILE, cache) + logger.info("SFTP ingest: %d new file(s) enqueued from %s.", n, ingest_folder) + return {"status": "ok", "files_enqueued": n, "folder": ingest_folder} + + +@shared_task +def scan_all_watch_folders() -> dict: + """ + Main periodic Celery task that runs all watch-folder scans. + + Acquires a Redis lock to prevent concurrent runs, then sequentially: + 1. Scans local filesystem watch folders + 2. Scans FTP ingest folder (if enabled) + 3. Scans SFTP ingest folder (if enabled) + """ + if not _acquire_lock(WATCH_FOLDER_LOCK_KEY): + logger.info("Watch folder scan already running — skipping this cycle.") + return {"status": "skipped", "reason": "lock held"} + + results: dict = {} + try: + results["local"] = scan_local_watch_folders() + results["ftp"] = scan_ftp_watch_folder() + results["sftp"] = scan_sftp_watch_folder() + finally: + _release_lock(WATCH_FOLDER_LOCK_KEY) + + return {"status": "ok", "results": results} diff --git a/app/utils/config_validator/settings_display.py b/app/utils/config_validator/settings_display.py index 1e417166..c8df7cbe 100644 --- a/app/utils/config_validator/settings_display.py +++ b/app/utils/config_validator/settings_display.py @@ -106,6 +106,17 @@ def get_settings_for_display(show_values: bool = False) -> dict[str, list[dict[s "email_sender", "email_default_recipient", ], + "Watch Folders": [ + "watch_folders", + "watch_folder_poll_interval", + "watch_folder_delete_after_process", + "ftp_ingest_enabled", + "ftp_ingest_folder", + "ftp_ingest_delete_after_process", + "sftp_ingest_enabled", + "sftp_ingest_folder", + "sftp_ingest_delete_after_process", + ], "IMAP": [ "imap1_host", "imap1_port", diff --git a/app/utils/settings_service.py b/app/utils/settings_service.py index eb622f1b..f6646833 100644 --- a/app/utils/settings_service.py +++ b/app/utils/settings_service.py @@ -971,6 +971,94 @@ SETTING_METADATA = { "required": False, "restart_required": False, }, + # Watch Folder / Ingest Settings + "watch_folders": { + "category": "Watch Folders", + "description": ( + "Comma-separated list of absolute local filesystem paths that DocuElevate will " + "poll for new files to ingest. Works with any mounted path (SMB/CIFS, NFS, etc.). " + "Example: /watchfolders/scanner,/mnt/shared/inbox" + ), + "type": "string", + "sensitive": False, + "required": False, + "restart_required": False, + }, + "watch_folder_poll_interval": { + "category": "Watch Folders", + "description": "Poll interval in minutes for local watch folder scanning (default: 1)", + "type": "integer", + "sensitive": False, + "required": False, + "restart_required": False, + }, + "watch_folder_delete_after_process": { + "category": "Watch Folders", + "description": ( + "Delete files from local watch folders after they have been enqueued for processing. " + "When False (default), processed files are tracked via cache to avoid re-ingestion." + ), + "type": "boolean", + "sensitive": False, + "required": False, + "restart_required": False, + }, + # FTP Ingest + "ftp_ingest_enabled": { + "category": "Watch Folders", + "description": "Enable FTP watch folder ingestion (requires FTP_INGEST_FOLDER and FTP connection settings)", + "type": "boolean", + "sensitive": False, + "required": False, + "restart_required": False, + }, + "ftp_ingest_folder": { + "category": "Watch Folders", + "description": ( + "FTP folder path to poll for new documents to ingest. " + "Uses the existing FTP connection settings (FTP_HOST, FTP_USERNAME, FTP_PASSWORD)." + ), + "type": "string", + "sensitive": False, + "required": False, + "restart_required": False, + }, + "ftp_ingest_delete_after_process": { + "category": "Watch Folders", + "description": "Delete files from the FTP ingest folder after they have been downloaded and enqueued", + "type": "boolean", + "sensitive": False, + "required": False, + "restart_required": False, + }, + # SFTP Ingest + "sftp_ingest_enabled": { + "category": "Watch Folders", + "description": "Enable SFTP watch folder ingestion (requires SFTP_INGEST_FOLDER and SFTP connection settings)", + "type": "boolean", + "sensitive": False, + "required": False, + "restart_required": False, + }, + "sftp_ingest_folder": { + "category": "Watch Folders", + "description": ( + "SFTP folder path to poll for new documents to ingest. " + "Uses the existing SFTP connection settings (SFTP_HOST, SFTP_USERNAME, SFTP_PASSWORD/SFTP_PRIVATE_KEY)." + ), + "type": "string", + "sensitive": False, + "required": False, + "restart_required": False, + }, + "sftp_ingest_delete_after_process": { + "category": "Watch Folders", + "description": "Delete files from the SFTP ingest folder after they have been downloaded and enqueued", + "type": "boolean", + "sensitive": False, + "required": False, + "restart_required": False, + }, # IMAP Settings - Account 1 "imap1_host": { "category": "IMAP", diff --git a/docs/ConfigurationGuide.md b/docs/ConfigurationGuide.md index 863afb09..b4f63b15 100644 --- a/docs/ConfigurationGuide.md +++ b/docs/ConfigurationGuide.md @@ -121,9 +121,97 @@ MAX_SINGLE_FILE_SIZE=524288000 - **With splitting**: Recommended for servers with limited memory or when processing very large scanned documents - **Higher limits**: For environments specifically designed to handle large architectural plans, books, or scanned archives -### IMAP Configuration +### Watch Folder Ingestion -DocuElevate can monitor multiple IMAP mailboxes for document attachments. Each mailbox uses a numbered prefix (e.g., `IMAP1_`, `IMAP2_`). +DocuElevate can automatically monitor directories for new files and ingest them without any manual action. +This works for: +- **Local filesystem paths** — including SMB/CIFS shares, NFS mounts, or any path accessible to the Docker container +- **FTP server directories** — using the configured FTP connection credentials +- **SFTP server directories** — using the configured SFTP connection credentials + +#### Local Watch Folders + +Mount the share or directory into the Docker container and configure one or more paths to watch. + +| **Variable** | **Description** | **Default** | +|-------------------------------------|----------------------------------------------------------------------------------------------|-------------| +| `WATCH_FOLDERS` | Comma-separated list of **absolute** local filesystem paths to poll for new files. | *(empty)* | +| `WATCH_FOLDER_POLL_INTERVAL` | How often to scan the folders, in minutes. | `1` | +| `WATCH_FOLDER_DELETE_AFTER_PROCESS` | Delete source files from the watch folder after they are successfully enqueued. When `false`, processed files are tracked in a cache file to prevent re-ingestion. | `false` | + +**Example (docker-compose.yaml):** + +```yaml +services: + worker: + volumes: + - /mnt/smb/scanner:/watchfolders/scanner # SMB/CIFS share mounted on the host + - /mnt/nfs/inbox:/watchfolders/inbox # NFS mount + environment: + WATCH_FOLDERS: /watchfolders/scanner,/watchfolders/inbox + WATCH_FOLDER_POLL_INTERVAL: 1 + WATCH_FOLDER_DELETE_AFTER_PROCESS: false +``` + +> **Tip for HP Scanners and MFPs**: Configure your scanner's "Scan to Network Folder" to point at an SMB share that is also mounted into the DocuElevate worker container. DocuElevate will pick up the scan files automatically every minute. No email forwarding is required. + +#### FTP Ingest (Watch Folder) + +DocuElevate can poll an FTP server directory for new files. It reuses the FTP connection settings already configured for uploads. + +| **Variable** | **Description** | **Default** | +|-----------------------------------|------------------------------------------------------------------------------------------------|-------------| +| `FTP_INGEST_ENABLED` | Enable FTP folder watching (`true`/`false`). | `false` | +| `FTP_INGEST_FOLDER` | Path on the FTP server to poll (e.g. `/incoming`). Uses the existing FTP connection settings. | *(empty)* | +| `FTP_INGEST_DELETE_AFTER_PROCESS` | Delete files from the FTP server after they are downloaded and enqueued. | `false` | + +**Example:** + +```dotenv +# Existing FTP upload settings (also used for ingest) +FTP_HOST=ftp.example.com +FTP_USERNAME=docuelevate +FTP_PASSWORD=secret + +# FTP ingest configuration +FTP_INGEST_ENABLED=true +FTP_INGEST_FOLDER=/incoming +FTP_INGEST_DELETE_AFTER_PROCESS=false +``` + +#### SFTP Ingest (Watch Folder) + +DocuElevate can poll an SFTP server directory for new files. It reuses the SFTP connection settings already configured for uploads. + +| **Variable** | **Description** | **Default** | +|------------------------------------|-------------------------------------------------------------------------------------------------|-------------| +| `SFTP_INGEST_ENABLED` | Enable SFTP folder watching (`true`/`false`). | `false` | +| `SFTP_INGEST_FOLDER` | Path on the SFTP server to poll (e.g. `/uploads/inbox`). Uses the existing SFTP connection settings. | *(empty)* | +| `SFTP_INGEST_DELETE_AFTER_PROCESS` | Delete files from the SFTP server after they are downloaded and enqueued. | `false` | + +**Example:** + +```dotenv +# Existing SFTP upload settings (also used for ingest) +SFTP_HOST=sftp.example.com +SFTP_USERNAME=docuelevate +SFTP_PRIVATE_KEY=/run/secrets/sftp_key + +# SFTP ingest configuration +SFTP_INGEST_ENABLED=true +SFTP_INGEST_FOLDER=/uploads/inbox +SFTP_INGEST_DELETE_AFTER_PROCESS=false +``` + +#### Supported File Types for Watch Folders + +Watch folder ingestion accepts the same file types as the web upload interface: PDF, Word, Excel, PowerPoint, images (JPEG, PNG, TIFF, BMP, GIF), plain text, CSV, RTF, and more. Unsupported files (executables, archives, etc.) are silently skipped. + +### IMAP Email Ingestion + +DocuElevate can automatically pull document attachments from IMAP mailboxes — no need to forward emails manually. Configure one or two mailboxes and DocuElevate polls them on the schedule you set. + +> **For HP Scanners (Scan to Email)**: If your scanner is set up to email scanned documents to a dedicated mailbox, configure that mailbox in DocuElevate using the settings below. DocuElevate will automatically retrieve the scanned PDFs from the inbox and process them. You do **not** need to configure DocuElevate as an email server — it acts as an email *client* that reads from your existing mailbox. | **Variable** | **Description** | **Example** | |-------------------------------|--------------------------------------------------------------|-------------------| diff --git a/docs/UserGuide.md b/docs/UserGuide.md index 399791ea..59ffdc40 100644 --- a/docs/UserGuide.md +++ b/docs/UserGuide.md @@ -74,14 +74,37 @@ For even more convenience, you can upload files directly from the **Files** page This feature allows you to quickly add new files without navigating away from your document management view. -### Email Attachments +### Email Attachments (IMAP Ingestion) -If configured, DocuElevate can automatically fetch documents from email attachments: +DocuElevate acts as an email *client* that automatically retrieves document attachments from one or more IMAP mailboxes. You do not need to set up DocuElevate as an email server — it simply polls an existing mailbox that you designate for document delivery. -1. Send an email with attachments to the configured email account -2. DocuElevate will poll the mailbox at the configured interval -3. Attachments will be automatically downloaded and processed -4. No further action is required +**How it works:** +1. A document is sent as an email attachment to the configured mailbox (e.g. from a scanner, a colleague, or any email client) +2. DocuElevate polls the mailbox at the configured interval (typically every 1–5 minutes) +3. Email attachments in supported formats are automatically downloaded and enqueued for processing +4. Processed emails are marked with a label or star (Gmail) or tracked locally, so they are not re-processed + +> **HP Scanners and MFPs (Scan to Email)**: Configure your scanner's "Scan to Email" feature to send scanned documents to a dedicated email account. Point DocuElevate at that mailbox using the IMAP settings. DocuElevate will retrieve the scanned PDFs automatically — no manual forwarding required. + +### Watch Folders (Automatic Folder Ingestion) + +Watch folders allow DocuElevate to automatically monitor directories for new files and ingest them without any manual action. + +#### Local Watch Folders (including SMB/CIFS and NFS) + +Mount a network share or local directory into the DocuElevate worker container and configure the path in `WATCH_FOLDERS`. DocuElevate scans the folder every minute (configurable via `WATCH_FOLDER_POLL_INTERVAL`) and enqueues any new documents it finds. + +This is the recommended approach for: +- **HP Scanners / MFPs** using "Scan to Network Folder" — point the scanner at a shared folder that DocuElevate also has access to +- **SMB/CIFS shares** — mount the Windows/Samba share and add the path to `WATCH_FOLDERS` +- **NFS mounts** — works identically, just configure the mount path +- **Any local directory** on the server running DocuElevate + +#### FTP / SFTP Watch Folders + +DocuElevate can poll an FTP or SFTP directory for new files. Enable this with `FTP_INGEST_ENABLED` or `SFTP_INGEST_ENABLED` and set the corresponding ingest folder. DocuElevate downloads new files, enqueues them for processing, and optionally deletes them from the remote server. + +See [Configuration Guide — Watch Folder Ingestion](ConfigurationGuide.md#watch-folder-ingestion) for full setup instructions. ## Managing Documents diff --git a/frontend/templates/settings.html b/frontend/templates/settings.html index 40311ea1..bc61f82a 100644 --- a/frontend/templates/settings.html +++ b/frontend/templates/settings.html @@ -9,6 +9,7 @@ 'OCR Engines': 'fas fa-file-alt', 'Storage Providers': 'fas fa-cloud-upload-alt', 'Email': 'fas fa-envelope', + 'Watch Folders': 'fas fa-folder-open', 'IMAP': 'fas fa-inbox', 'Monitoring': 'fas fa-chart-line', 'Processing': 'fas fa-cogs', diff --git a/tests/test_watch_folder_tasks.py b/tests/test_watch_folder_tasks.py new file mode 100644 index 00000000..d5e3e6e6 --- /dev/null +++ b/tests/test_watch_folder_tasks.py @@ -0,0 +1,494 @@ +"""Tests for app/tasks/watch_folder_tasks.py module.""" + +import json +import os +import tempfile +from datetime import datetime, timedelta, timezone +from unittest.mock import MagicMock, patch + +import pytest + + +@pytest.mark.unit +class TestCacheHelpers: + """Tests for the cache helper functions in watch_folder_tasks.""" + + def test_evict_old_entries_removes_old(self): + """Entries older than CACHE_RETENTION_DAYS should be removed.""" + from app.tasks.watch_folder_tasks import _evict_old_entries + + old_dt = (datetime.now(timezone.utc) - timedelta(days=40)).isoformat() + recent_dt = datetime.now(timezone.utc).isoformat() + data = {"old_key": old_dt, "new_key": recent_dt} + result = _evict_old_entries(data) + assert "old_key" not in result + assert "new_key" in result + + def test_evict_old_entries_empty(self): + """Empty cache should return empty dict.""" + from app.tasks.watch_folder_tasks import _evict_old_entries + + assert _evict_old_entries({}) == {} + + def test_evict_old_entries_skips_malformed_dates(self): + """Entries with malformed timestamps should be silently dropped.""" + from app.tasks.watch_folder_tasks import _evict_old_entries + + data = {"bad_key": "not-a-date", "good_key": datetime.now(timezone.utc).isoformat()} + result = _evict_old_entries(data) + assert "bad_key" not in result + assert "good_key" in result + + def test_mark_processed_adds_entry(self): + """_mark_processed should add an ISO-formatted timestamp for the key.""" + from app.tasks.watch_folder_tasks import _mark_processed + + cache: dict = {} + _mark_processed(cache, "/some/file.pdf") + assert "/some/file.pdf" in cache + # Timestamp should be parseable + datetime.fromisoformat(cache["/some/file.pdf"]) + + def test_load_cache_returns_empty_when_no_file(self): + """_load_cache should return {} when the cache file does not exist.""" + from app.tasks.watch_folder_tasks import _load_cache + + result = _load_cache("/tmp/does_not_exist_xyz.json") + assert result == {} + + def test_save_and_load_roundtrip(self, tmp_path): + """Saving and loading cache should preserve entries.""" + from app.tasks.watch_folder_tasks import _load_cache, _save_cache + + cache_file = str(tmp_path / "cache.json") + data = {"key1": datetime.now(timezone.utc).isoformat()} + _save_cache(cache_file, data) + loaded = _load_cache(cache_file) + assert "key1" in loaded + + def test_load_cache_handles_invalid_json(self, tmp_path): + """_load_cache should return {} for corrupted JSON files.""" + from app.tasks.watch_folder_tasks import _load_cache + + cache_file = str(tmp_path / "bad.json") + with open(cache_file, "w") as f: + f.write("not valid json{{{{") + result = _load_cache(cache_file) + assert result == {} + + +@pytest.mark.unit +class TestIsAllowedFile: + """Tests for the _is_allowed_file helper.""" + + def test_pdf_is_allowed(self): + from app.tasks.watch_folder_tasks import _is_allowed_file + + assert _is_allowed_file("document.pdf") is True + assert _is_allowed_file("DOCUMENT.PDF") is True + + def test_docx_is_allowed(self): + from app.tasks.watch_folder_tasks import _is_allowed_file + + assert _is_allowed_file("report.docx") is True + + def test_exe_is_not_allowed(self): + from app.tasks.watch_folder_tasks import _is_allowed_file + + assert _is_allowed_file("malware.exe") is False + + def test_zip_is_not_allowed(self): + from app.tasks.watch_folder_tasks import _is_allowed_file + + assert _is_allowed_file("archive.zip") is False + + +@pytest.mark.unit +class TestScanLocalFolder: + """Tests for _scan_local_folder.""" + + def test_nonexistent_folder_returns_zero(self): + from app.tasks.watch_folder_tasks import _scan_local_folder + + count = _scan_local_folder("/tmp/does_not_exist_xyz_abc", {}, False) + assert count == 0 + + def test_empty_folder_returns_zero(self, tmp_path): + from app.tasks.watch_folder_tasks import _scan_local_folder + + count = _scan_local_folder(str(tmp_path), {}, False) + assert count == 0 + + def test_new_pdf_is_enqueued(self, tmp_path): + """A new PDF in the watch folder should be enqueued for processing.""" + from app.tasks.watch_folder_tasks import _scan_local_folder + + pdf_file = tmp_path / "test.pdf" + pdf_file.write_bytes(b"%PDF-1.4 test") + + cache: dict = {} + with patch("app.tasks.watch_folder_tasks.process_document") as mock_proc, patch( + "app.tasks.watch_folder_tasks.settings" + ) as mock_settings: + mock_settings.workdir = str(tmp_path / "workdir") + os.makedirs(mock_settings.workdir, exist_ok=True) + count = _scan_local_folder(str(tmp_path), cache, False) + + assert count == 1 + assert str(pdf_file) in cache + + def test_already_cached_file_is_skipped(self, tmp_path): + """Files already in cache should not be re-processed.""" + from app.tasks.watch_folder_tasks import _scan_local_folder + + pdf_file = tmp_path / "already.pdf" + pdf_file.write_bytes(b"%PDF-1.4 test") + + cache = {str(pdf_file): datetime.now(timezone.utc).isoformat()} + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.workdir = str(tmp_path / "workdir") + os.makedirs(mock_settings.workdir, exist_ok=True) + count = _scan_local_folder(str(tmp_path), cache, False) + + assert count == 0 + + def test_unsupported_file_is_skipped(self, tmp_path): + """Unsupported file types should not be enqueued.""" + from app.tasks.watch_folder_tasks import _scan_local_folder + + exe_file = tmp_path / "bad.exe" + exe_file.write_bytes(b"MZ malware") + + cache: dict = {} + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.workdir = str(tmp_path / "workdir") + count = _scan_local_folder(str(tmp_path), cache, False) + + assert count == 0 + assert str(exe_file) not in cache + + def test_delete_after_process_removes_source(self, tmp_path): + """When delete_after_process=True, source files should be deleted.""" + from app.tasks.watch_folder_tasks import _scan_local_folder + + pdf_file = tmp_path / "invoice.pdf" + pdf_file.write_bytes(b"%PDF-1.4") + + cache: dict = {} + with patch("app.tasks.watch_folder_tasks.process_document"), patch( + "app.tasks.watch_folder_tasks.settings" + ) as mock_settings: + mock_settings.workdir = str(tmp_path / "workdir") + os.makedirs(mock_settings.workdir, exist_ok=True) + _scan_local_folder(str(tmp_path), cache, delete_after=True) + + assert not pdf_file.exists() + + +@pytest.mark.unit +class TestAcquireReleaseLock: + """Tests for the Redis-based locking helpers.""" + + def test_acquire_lock_succeeds(self): + from app.tasks.watch_folder_tasks import _acquire_lock + + mock_redis = MagicMock() + mock_redis.setnx.return_value = True + with patch("app.tasks.watch_folder_tasks.redis_client", mock_redis): + result = _acquire_lock("test_lock") + assert result is True + + def test_acquire_lock_fails_when_held(self): + from app.tasks.watch_folder_tasks import _acquire_lock + + mock_redis = MagicMock() + mock_redis.setnx.return_value = False + with patch("app.tasks.watch_folder_tasks.redis_client", mock_redis): + result = _acquire_lock("test_lock") + assert result is False + + def test_release_lock_deletes_key(self): + from app.tasks.watch_folder_tasks import _release_lock + + mock_redis = MagicMock() + with patch("app.tasks.watch_folder_tasks.redis_client", mock_redis): + _release_lock("test_lock") + mock_redis.delete.assert_called_once_with("test_lock") + + +@pytest.mark.unit +class TestScanLocalWatchFoldersTask: + """Tests for the scan_local_watch_folders Celery task.""" + + def test_returns_skipped_when_no_folders_configured(self): + from app.tasks.watch_folder_tasks import scan_local_watch_folders + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.watch_folders = None + result = scan_local_watch_folders() + assert result["status"] == "skipped" + + def test_returns_skipped_for_empty_string(self): + from app.tasks.watch_folder_tasks import scan_local_watch_folders + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.watch_folders = "" + result = scan_local_watch_folders() + assert result["status"] == "skipped" + + def test_scans_configured_folder(self, tmp_path): + """With a valid folder configured, the task should scan it.""" + from app.tasks.watch_folder_tasks import scan_local_watch_folders + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings, patch( + "app.tasks.watch_folder_tasks._load_cache", return_value={} + ), patch("app.tasks.watch_folder_tasks._save_cache"), patch( + "app.tasks.watch_folder_tasks._scan_local_folder", return_value=0 + ) as mock_scan: + mock_settings.watch_folders = str(tmp_path) + mock_settings.watch_folder_delete_after_process = False + result = scan_local_watch_folders() + + mock_scan.assert_called_once() + assert result["status"] == "ok" + assert result["files_enqueued"] == 0 + + +@pytest.mark.unit +class TestScanFtpWatchFolderTask: + """Tests for the scan_ftp_watch_folder Celery task.""" + + def test_returns_skipped_when_disabled(self): + from app.tasks.watch_folder_tasks import scan_ftp_watch_folder + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.ftp_ingest_enabled = False + result = scan_ftp_watch_folder() + assert result["status"] == "skipped" + + def test_returns_skipped_when_no_folder_configured(self): + from app.tasks.watch_folder_tasks import scan_ftp_watch_folder + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.ftp_ingest_enabled = True + mock_settings.ftp_ingest_folder = None + result = scan_ftp_watch_folder() + assert result["status"] == "skipped" + + def test_returns_error_when_connection_fails(self): + from app.tasks.watch_folder_tasks import scan_ftp_watch_folder + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings, patch( + "app.tasks.watch_folder_tasks._connect_ftp", return_value=None + ): + mock_settings.ftp_ingest_enabled = True + mock_settings.ftp_ingest_folder = "/inbox" + mock_settings.ftp_ingest_delete_after_process = False + result = scan_ftp_watch_folder() + assert result["status"] == "error" + + def test_successful_scan(self): + from app.tasks.watch_folder_tasks import scan_ftp_watch_folder + + mock_ftp = MagicMock() + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings, patch( + "app.tasks.watch_folder_tasks._connect_ftp", return_value=mock_ftp + ), patch("app.tasks.watch_folder_tasks._load_cache", return_value={}), patch( + "app.tasks.watch_folder_tasks._save_cache" + ), patch( + "app.tasks.watch_folder_tasks._scan_ftp_folder", return_value=2 + ): + mock_settings.ftp_ingest_enabled = True + mock_settings.ftp_ingest_folder = "/inbox" + mock_settings.ftp_ingest_delete_after_process = False + result = scan_ftp_watch_folder() + + assert result["status"] == "ok" + assert result["files_enqueued"] == 2 + + +@pytest.mark.unit +class TestScanSftpWatchFolderTask: + """Tests for the scan_sftp_watch_folder Celery task.""" + + def test_returns_skipped_when_disabled(self): + from app.tasks.watch_folder_tasks import scan_sftp_watch_folder + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.sftp_ingest_enabled = False + result = scan_sftp_watch_folder() + assert result["status"] == "skipped" + + def test_returns_skipped_when_no_folder_configured(self): + from app.tasks.watch_folder_tasks import scan_sftp_watch_folder + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.sftp_ingest_enabled = True + mock_settings.sftp_ingest_folder = None + result = scan_sftp_watch_folder() + assert result["status"] == "skipped" + + def test_returns_error_when_connection_fails(self): + from app.tasks.watch_folder_tasks import scan_sftp_watch_folder + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings, patch( + "app.tasks.watch_folder_tasks._get_sftp_connection", return_value=(None, None) + ): + mock_settings.sftp_ingest_enabled = True + mock_settings.sftp_ingest_folder = "/upload" + mock_settings.sftp_ingest_delete_after_process = False + result = scan_sftp_watch_folder() + assert result["status"] == "error" + + def test_successful_scan(self): + from app.tasks.watch_folder_tasks import scan_sftp_watch_folder + + mock_ssh = MagicMock() + mock_sftp = MagicMock() + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings, patch( + "app.tasks.watch_folder_tasks._get_sftp_connection", return_value=(mock_ssh, mock_sftp) + ), patch("app.tasks.watch_folder_tasks._load_cache", return_value={}), patch( + "app.tasks.watch_folder_tasks._save_cache" + ), patch( + "app.tasks.watch_folder_tasks._scan_sftp_folder", return_value=3 + ): + mock_settings.sftp_ingest_enabled = True + mock_settings.sftp_ingest_folder = "/upload" + mock_settings.sftp_ingest_delete_after_process = False + result = scan_sftp_watch_folder() + + assert result["status"] == "ok" + assert result["files_enqueued"] == 3 + + +@pytest.mark.unit +class TestScanAllWatchFolders: + """Tests for the scan_all_watch_folders orchestrator task.""" + + def test_skips_when_lock_held(self): + from app.tasks.watch_folder_tasks import scan_all_watch_folders + + with patch("app.tasks.watch_folder_tasks._acquire_lock", return_value=False): + result = scan_all_watch_folders() + assert result["status"] == "skipped" + + def test_runs_all_scans_and_releases_lock(self): + from app.tasks.watch_folder_tasks import scan_all_watch_folders + + with patch("app.tasks.watch_folder_tasks._acquire_lock", return_value=True), patch( + "app.tasks.watch_folder_tasks._release_lock" + ) as mock_release, patch( + "app.tasks.watch_folder_tasks.scan_local_watch_folders", return_value={"status": "ok"} + ), patch( + "app.tasks.watch_folder_tasks.scan_ftp_watch_folder", return_value={"status": "skipped"} + ), patch( + "app.tasks.watch_folder_tasks.scan_sftp_watch_folder", return_value={"status": "skipped"} + ): + result = scan_all_watch_folders() + + assert result["status"] == "ok" + assert "results" in result + mock_release.assert_called_once() + + def test_lock_released_even_on_exception(self): + """Lock must be released even if a sub-scan raises an exception.""" + from app.tasks.watch_folder_tasks import scan_all_watch_folders + + with patch("app.tasks.watch_folder_tasks._acquire_lock", return_value=True), patch( + "app.tasks.watch_folder_tasks._release_lock" + ) as mock_release, patch( + "app.tasks.watch_folder_tasks.scan_local_watch_folders", side_effect=RuntimeError("boom") + ): + with pytest.raises(RuntimeError): + scan_all_watch_folders() + + mock_release.assert_called_once() + + +@pytest.mark.unit +class TestConnectFtp: + """Tests for the _connect_ftp helper.""" + + def test_returns_none_when_settings_incomplete(self): + from app.tasks.watch_folder_tasks import _connect_ftp + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings: + mock_settings.ftp_host = None + mock_settings.ftp_username = "user" + mock_settings.ftp_password = "pass" # noqa: S105 + result = _connect_ftp() + assert result is None + + def test_returns_none_when_connection_fails(self): + from app.tasks.watch_folder_tasks import _connect_ftp + + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings, patch( + "ftplib.FTP_TLS" + ) as mock_ftps_cls, patch("ftplib.FTP") as mock_ftp_cls: + mock_settings.ftp_host = "ftp.example.com" + mock_settings.ftp_port = 21 + mock_settings.ftp_username = "user" + mock_settings.ftp_password = "pass" # noqa: S105 + mock_settings.ftp_use_tls = False + mock_settings.ftp_allow_plaintext = True + + mock_ftp_cls.return_value.connect.side_effect = ConnectionRefusedError("refused") + result = _connect_ftp() + assert result is None + + +@pytest.mark.unit +class TestScanFtpFolder: + """Tests for the _scan_ftp_folder helper.""" + + def test_cwd_failure_returns_zero(self): + import ftplib + + from app.tasks.watch_folder_tasks import _scan_ftp_folder + + mock_ftp = MagicMock() + mock_ftp.cwd.side_effect = ftplib.error_perm("550 no such directory") + count = _scan_ftp_folder(mock_ftp, "/missing", {}, False) + assert count == 0 + + def test_skips_disallowed_files(self): + from app.tasks.watch_folder_tasks import _scan_ftp_folder + + mock_ftp = MagicMock() + mock_ftp.cwd.return_value = None + mock_ftp.nlst.return_value = ["photo.exe", "virus.bat"] + count = _scan_ftp_folder(mock_ftp, "/inbox", {}, False) + assert count == 0 + + def test_downloads_new_allowed_file(self, tmp_path): + from app.tasks.watch_folder_tasks import _scan_ftp_folder + + mock_ftp = MagicMock() + mock_ftp.cwd.return_value = None + mock_ftp.nlst.return_value = ["invoice.pdf"] + + cache: dict = {} + with patch("app.tasks.watch_folder_tasks.settings") as mock_settings, patch( + "app.tasks.watch_folder_tasks.process_document" + ): + mock_settings.workdir = str(tmp_path) + # Simulate retrbinary writing bytes + def fake_retrbinary(cmd, callback): + callback(b"%PDF-1.4") + + mock_ftp.retrbinary.side_effect = fake_retrbinary + count = _scan_ftp_folder(mock_ftp, "/inbox", cache, False) + + assert count == 1 + assert "ftp:/inbox/invoice.pdf" in cache + + def test_already_cached_file_is_skipped(self, tmp_path): + from app.tasks.watch_folder_tasks import _scan_ftp_folder + + mock_ftp = MagicMock() + mock_ftp.cwd.return_value = None + mock_ftp.nlst.return_value = ["invoice.pdf"] + + cache = {"ftp:/inbox/invoice.pdf": datetime.now(timezone.utc).isoformat()} + count = _scan_ftp_folder(mock_ftp, "/inbox", cache, False) + assert count == 0