#!/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 from datetime import datetime, timedelta, timezone from typing import Any 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 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") # Per-integration watch-folder cache file prefix _USER_WF_CACHE_PREFIX = os.path.join(settings.workdir, "user_wf_") _CACHE_RETENTION_DAYS = 30 # Maximum length to store as last_error on UserIntegration to prevent DB bloat _MAX_ERROR_LENGTH = 500 # Database session factory (imported lazily to avoid circular imports) _db_session_factory = None def _get_db_session(): """Return a new SQLAlchemy session (lazy import to avoid startup issues).""" global _db_session_factory # noqa: PLW0603 if _db_session_factory is None: from app.database import SessionLocal _db_session_factory = SessionLocal return _db_session_factory() # --------------------------------------------------------------------------- # 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, owner_id: str | None = None) -> None: """Enqueue a local file path for document processing. Args: file_path: Absolute path to the file on disk. filename: Optional display filename (defaults to basename of *file_path*). owner_id: Optional user identifier forwarded to ``process_document`` / ``convert_to_pdf`` for multi-tenant attribution. """ 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, owner_id=owner_id) logger.info("Enqueued for processing: %s", fname) else: convert_to_pdf.delay(file_path, owner_id=owner_id) 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 as close_exc: logger.debug("SFTP ingest: error closing SSH after failed connection: %s", close_exc) 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 # --------------------------------------------------------------------------- # Dropbox watch folder scanning # --------------------------------------------------------------------------- # Additional cache files for cloud providers DROPBOX_INGEST_CACHE_FILE = os.path.join(settings.workdir, "dropbox_ingest_processed.json") GDRIVE_INGEST_CACHE_FILE = os.path.join(settings.workdir, "gdrive_ingest_processed.json") ONEDRIVE_INGEST_CACHE_FILE = os.path.join(settings.workdir, "onedrive_ingest_processed.json") NEXTCLOUD_INGEST_CACHE_FILE = os.path.join(settings.workdir, "nextcloud_ingest_processed.json") S3_INGEST_CACHE_FILE = os.path.join(settings.workdir, "s3_ingest_processed.json") WEBDAV_INGEST_CACHE_FILE = os.path.join(settings.workdir, "webdav_ingest_processed.json") def _scan_dropbox_folder(folder_path: str, cache: dict[str, str], delete_after: bool) -> int: """ List *folder_path* in Dropbox and download new allowed files to workdir. Returns the number of files newly enqueued. """ try: from app.tasks.upload_to_dropbox import get_dropbox_client except ImportError as exc: logger.error("Dropbox ingest: dropbox SDK not installed: %s", exc) return 0 try: dbx = get_dropbox_client() except Exception as exc: logger.error("Dropbox ingest: authentication failed: %s", exc) return 0 count = 0 try: result = dbx.files_list_folder(folder_path) entries = result.entries while result.has_more: result = dbx.files_list_folder_continue(result.cursor) entries.extend(result.entries) except Exception as exc: logger.error("Dropbox ingest: cannot list folder %s: %s", folder_path, exc) return 0 for entry in entries: # Only process files, not sub-folders import dropbox as dropbox_module if not isinstance(entry, dropbox_module.files.FileMetadata): continue filename = entry.name if not _is_allowed_file(filename): logger.debug("Dropbox ingest: skipping %s (unsupported type)", filename) continue cache_key = f"dropbox:{entry.id}" if cache_key in cache: logger.debug("Dropbox ingest: already processed %s", filename) continue dest_path = os.path.join(settings.workdir, f"dropbox_{filename}") if os.path.exists(dest_path): base, ext2 = os.path.splitext(f"dropbox_{filename}") dest_path = os.path.join(settings.workdir, f"{base}_{int(datetime.now().timestamp())}{ext2}") try: _meta, response = dbx.files_download(entry.path_lower) with open(dest_path, "wb") as f: f.write(response.content) logger.info("Dropbox ingest: downloaded %s to %s", filename, dest_path) except Exception as exc: logger.error("Dropbox 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: dbx.files_delete_v2(entry.path_lower) logger.info("Dropbox ingest: deleted %s after ingestion", entry.path_lower) except Exception as exc: logger.warning("Dropbox ingest: could not delete %s: %s", entry.path_lower, exc) return count # --------------------------------------------------------------------------- # Google Drive watch folder scanning # --------------------------------------------------------------------------- def _scan_google_drive_folder(folder_id: str, cache: dict[str, str], delete_after: bool) -> int: """ List files in *folder_id* on Google Drive and download new allowed files to workdir. Returns the number of files newly enqueued. """ try: from app.tasks.upload_to_google_drive import get_google_drive_service except ImportError as exc: logger.error("Google Drive ingest: google-api SDK not installed: %s", exc) return 0 service = get_google_drive_service() if service is None: logger.error("Google Drive ingest: could not authenticate.") return 0 count = 0 query = f"'{folder_id}' in parents and trashed = false and mimeType != 'application/vnd.google-apps.folder'" page_token = None while True: try: params: dict = { "q": query, "fields": "nextPageToken, files(id, name, mimeType)", "pageSize": 100, } if page_token: params["pageToken"] = page_token response = service.files().list(**params).execute() except Exception as exc: logger.error("Google Drive ingest: listing folder %s failed: %s", folder_id, exc) break for file_meta in response.get("files", []): file_id_gd = file_meta["id"] filename = file_meta["name"] if not _is_allowed_file(filename): logger.debug("Google Drive ingest: skipping %s (unsupported type)", filename) continue cache_key = f"gdrive:{file_id_gd}" if cache_key in cache: logger.debug("Google Drive ingest: already processed %s", filename) continue dest_path = os.path.join(settings.workdir, f"gdrive_{filename}") if os.path.exists(dest_path): base, ext2 = os.path.splitext(f"gdrive_{filename}") dest_path = os.path.join(settings.workdir, f"{base}_{int(datetime.now().timestamp())}{ext2}") try: import io from googleapiclient.http import MediaIoBaseDownload request = service.files().get_media(fileId=file_id_gd) buf = io.BytesIO() downloader = MediaIoBaseDownload(buf, request) done = False while not done: _, done = downloader.next_chunk() with open(dest_path, "wb") as f: f.write(buf.getvalue()) logger.info("Google Drive ingest: downloaded %s to %s", filename, dest_path) except Exception as exc: logger.error("Google Drive 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: service.files().delete(fileId=file_id_gd).execute() logger.info("Google Drive ingest: deleted %s after ingestion", filename) except Exception as exc: logger.warning("Google Drive ingest: could not delete %s: %s", filename, exc) page_token = response.get("nextPageToken") if not page_token: break return count # --------------------------------------------------------------------------- # OneDrive watch folder scanning # --------------------------------------------------------------------------- def _scan_onedrive_folder(folder_path: str, cache: dict[str, str], delete_after: bool) -> int: """ List files in *folder_path* on OneDrive (Microsoft Graph) and download new allowed files to workdir. Returns the number of files newly enqueued. """ import requests as req_lib try: from app.tasks.upload_to_onedrive import get_onedrive_token except ImportError as exc: logger.error("OneDrive ingest: msal not installed: %s", exc) return 0 try: token = get_onedrive_token() except Exception as exc: logger.error("OneDrive ingest: authentication failed: %s", exc) return 0 headers = {"Authorization": f"Bearer {token}"} # URL-encode the path and construct the Graph API endpoint import urllib.parse encoded_path = urllib.parse.quote(folder_path.lstrip("/")) list_url = f"https://graph.microsoft.com/v1.0/me/drive/root:/{encoded_path}:/children" count = 0 while list_url: try: resp = req_lib.get(list_url, headers=headers, timeout=getattr(settings, "http_request_timeout", 120)) resp.raise_for_status() data = resp.json() except Exception as exc: logger.error("OneDrive ingest: listing folder %s failed: %s", folder_path, exc) break for item in data.get("value", []): # Skip folders if "folder" in item: continue filename = item["name"] item_id = item["id"] if not _is_allowed_file(filename): logger.debug("OneDrive ingest: skipping %s (unsupported type)", filename) continue cache_key = f"onedrive:{item_id}" if cache_key in cache: logger.debug("OneDrive ingest: already processed %s", filename) continue # Get download URL download_url = item.get("@microsoft.graph.downloadUrl") if not download_url: logger.warning("OneDrive ingest: no download URL for %s — skipping.", filename) continue dest_path = os.path.join(settings.workdir, f"onedrive_{filename}") if os.path.exists(dest_path): base, ext2 = os.path.splitext(f"onedrive_{filename}") dest_path = os.path.join(settings.workdir, f"{base}_{int(datetime.now().timestamp())}{ext2}") try: dl_resp = req_lib.get( download_url, headers=headers, timeout=getattr(settings, "http_request_timeout", 120) ) dl_resp.raise_for_status() with open(dest_path, "wb") as f: f.write(dl_resp.content) logger.info("OneDrive ingest: downloaded %s to %s", filename, dest_path) except Exception as exc: logger.error("OneDrive 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: del_resp = req_lib.delete( f"https://graph.microsoft.com/v1.0/me/drive/items/{item_id}", headers=headers, timeout=getattr(settings, "http_request_timeout", 120), ) del_resp.raise_for_status() logger.info("OneDrive ingest: deleted %s after ingestion", filename) except Exception as exc: logger.warning("OneDrive ingest: could not delete %s: %s", filename, exc) list_url = data.get("@odata.nextLink") return count # --------------------------------------------------------------------------- # Nextcloud watch folder scanning (WebDAV) # --------------------------------------------------------------------------- def _scan_nextcloud_folder(folder_path: str, cache: dict[str, str], delete_after: bool) -> int: """ List files in *folder_path* on Nextcloud (via WebDAV PROPFIND) and download new allowed files to workdir. Returns the number of files newly enqueued. """ import defusedxml.ElementTree as ET import requests as req_lib from requests.auth import HTTPBasicAuth nc_url: str | None = getattr(settings, "nextcloud_upload_url", None) nc_user: str | None = getattr(settings, "nextcloud_username", None) nc_pass: str | None = getattr(settings, "nextcloud_password", None) if not (nc_url and nc_user and nc_pass): logger.warning("Nextcloud ingest: connection settings incomplete — skipping.") return 0 auth = HTTPBasicAuth(nc_user, nc_pass) timeout = getattr(settings, "http_request_timeout", 120) # Build the WebDAV PROPFIND URL base = nc_url.rstrip("/") folder = folder_path.strip("/") propfind_url = f"{base}/{folder}/" if folder else f"{base}/" try: resp = req_lib.request( "PROPFIND", propfind_url, auth=auth, headers={"Depth": "1", "Content-Type": "application/xml"}, timeout=timeout, ) resp.raise_for_status() except Exception as exc: logger.error("Nextcloud ingest: PROPFIND on %s failed: %s", propfind_url, exc) return 0 count = 0 # Parse WebDAV multistatus response using defusedxml (safe against XML bomb attacks) try: root = ET.fromstring(resp.text) # noqa: S314 — defusedxml is safe except Exception as exc: logger.error("Nextcloud ingest: failed to parse PROPFIND response: %s", exc) return 0 ns = {"d": "DAV:"} for response_el in root.findall("d:response", ns): href_el = response_el.find("d:href", ns) if href_el is None or href_el.text is None: continue href = href_el.text # Skip the folder itself if href.rstrip("/").endswith(folder.rstrip("/")): continue filename = href.rstrip("/").split("/")[-1] import urllib.parse filename = urllib.parse.unquote(filename) if not _is_allowed_file(filename): logger.debug("Nextcloud ingest: skipping %s (unsupported type)", filename) continue # Use the href as cache key (stable across runs) cache_key = f"nextcloud:{href}" if cache_key in cache: logger.debug("Nextcloud ingest: already processed %s", filename) continue # Build absolute download URL if href.startswith("http"): file_url = href else: from urllib.parse import urlparse parsed = urlparse(nc_url) file_url = f"{parsed.scheme}://{parsed.netloc}{href}" dest_path = os.path.join(settings.workdir, f"nc_{filename}") if os.path.exists(dest_path): base_name, ext2 = os.path.splitext(f"nc_{filename}") dest_path = os.path.join(settings.workdir, f"{base_name}_{int(datetime.now().timestamp())}{ext2}") try: dl = req_lib.get(file_url, auth=auth, timeout=timeout) dl.raise_for_status() with open(dest_path, "wb") as f: f.write(dl.content) logger.info("Nextcloud ingest: downloaded %s to %s", filename, dest_path) except Exception as exc: logger.error("Nextcloud 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: del_resp = req_lib.request("DELETE", file_url, auth=auth, timeout=timeout) del_resp.raise_for_status() logger.info("Nextcloud ingest: deleted %s after ingestion", filename) except Exception as exc: logger.warning("Nextcloud ingest: could not delete %s: %s", filename, exc) return count # --------------------------------------------------------------------------- # S3 watch folder scanning # --------------------------------------------------------------------------- def _scan_s3_prefix(prefix: str, cache: dict[str, str], delete_after: bool) -> int: """ List objects under *prefix* in the configured S3 bucket and download new allowed files to workdir. Returns the number of files newly enqueued. """ try: import boto3 from botocore.exceptions import ClientError except ImportError as exc: logger.error("S3 ingest: boto3 not installed: %s", exc) return 0 bucket = getattr(settings, "s3_bucket_name", None) if not bucket: logger.warning("S3 ingest: S3_BUCKET_NAME not set — skipping.") return 0 try: s3 = boto3.client( "s3", region_name=getattr(settings, "aws_region", "us-east-1"), aws_access_key_id=getattr(settings, "aws_access_key_id", None), aws_secret_access_key=getattr(settings, "aws_secret_access_key", None), ) except Exception as exc: logger.error("S3 ingest: failed to create S3 client: %s", exc) return 0 count = 0 paginator = s3.get_paginator("list_objects_v2") try: pages = paginator.paginate(Bucket=bucket, Prefix=prefix) except Exception as exc: logger.error("S3 ingest: failed to list objects in %s/%s: %s", bucket, prefix, exc) return 0 for page in pages: for obj in page.get("Contents", []): key = obj["Key"] filename = key.split("/")[-1] # Skip zero-byte "folder marker" objects and unsupported types if not filename or not _is_allowed_file(filename): logger.debug("S3 ingest: skipping %s (unsupported or empty)", key) continue cache_key = f"s3:{bucket}/{key}" if cache_key in cache: logger.debug("S3 ingest: already processed %s", key) continue dest_path = os.path.join(settings.workdir, f"s3_{filename}") if os.path.exists(dest_path): base2, ext2 = os.path.splitext(f"s3_{filename}") dest_path = os.path.join(settings.workdir, f"{base2}_{int(datetime.now().timestamp())}{ext2}") try: s3.download_file(bucket, key, dest_path) logger.info("S3 ingest: downloaded s3://%s/%s to %s", bucket, key, dest_path) except ClientError as exc: logger.error("S3 ingest: failed to download %s: %s", key, 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: s3.delete_object(Bucket=bucket, Key=key) logger.info("S3 ingest: deleted s3://%s/%s after ingestion", bucket, key) except Exception as exc: logger.warning("S3 ingest: could not delete s3://%s/%s: %s", bucket, key, exc) return count # --------------------------------------------------------------------------- # WebDAV watch folder scanning # --------------------------------------------------------------------------- def _scan_webdav_folder(folder_path: str, cache: dict[str, str], delete_after: bool) -> int: """ List files in *folder_path* on a WebDAV server (PROPFIND) and download new allowed files to workdir. Returns the number of files newly enqueued. """ import defusedxml.ElementTree as ET import requests as req_lib from requests.auth import HTTPBasicAuth webdav_url: str | None = getattr(settings, "webdav_url", None) webdav_user: str | None = getattr(settings, "webdav_username", None) webdav_pass: str | None = getattr(settings, "webdav_password", None) verify_ssl: bool = getattr(settings, "webdav_verify_ssl", True) timeout = getattr(settings, "http_request_timeout", 120) if not webdav_url: logger.warning("WebDAV ingest: WEBDAV_URL not configured — skipping.") return 0 from urllib.parse import unquote, urlparse base = webdav_url.rstrip("/") folder = folder_path.strip("/") propfind_url = f"{base}/{folder}/" if folder else f"{base}/" auth = HTTPBasicAuth(webdav_user, webdav_pass) if webdav_user else None try: resp = req_lib.request( "PROPFIND", propfind_url, auth=auth, headers={"Depth": "1"}, verify=verify_ssl, timeout=timeout, ) resp.raise_for_status() except Exception as exc: logger.error("WebDAV ingest: PROPFIND on %s failed: %s", propfind_url, exc) return 0 count = 0 try: root = ET.fromstring(resp.text) # noqa: S314 — defusedxml is safe except Exception as exc: logger.error("WebDAV ingest: failed to parse PROPFIND response: %s", exc) return 0 ns = {"d": "DAV:"} for response_el in root.findall("d:response", ns): href_el = response_el.find("d:href", ns) if href_el is None or href_el.text is None: continue href = href_el.text # Skip the folder itself and anything that looks like a directory if href.endswith("/"): continue filename = unquote(href.split("/")[-1]) if not _is_allowed_file(filename): logger.debug("WebDAV ingest: skipping %s (unsupported type)", filename) continue cache_key = f"webdav:{href}" if cache_key in cache: logger.debug("WebDAV ingest: already processed %s", filename) continue # Build absolute URL if href.startswith("http"): file_url = href else: parsed = urlparse(webdav_url) file_url = f"{parsed.scheme}://{parsed.netloc}{href}" dest_path = os.path.join(settings.workdir, f"webdav_{filename}") if os.path.exists(dest_path): base2, ext2 = os.path.splitext(f"webdav_{filename}") dest_path = os.path.join(settings.workdir, f"{base2}_{int(datetime.now().timestamp())}{ext2}") try: dl = req_lib.get(file_url, auth=auth, verify=verify_ssl, timeout=timeout) dl.raise_for_status() with open(dest_path, "wb") as f: f.write(dl.content) logger.info("WebDAV ingest: downloaded %s to %s", filename, dest_path) except Exception as exc: logger.error("WebDAV 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: del_resp = req_lib.request("DELETE", file_url, auth=auth, verify=verify_ssl, timeout=timeout) del_resp.raise_for_status() logger.info("WebDAV ingest: deleted %s after ingestion", filename) except Exception as exc: logger.warning("WebDAV ingest: could not delete %s: %s", filename, 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 as exc: logger.debug("FTP ingest: error during FTP quit: %s", exc) _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 as exc: logger.debug("SFTP ingest: error closing SFTP channel: %s", exc) try: ssh.close() except Exception as exc: logger.debug("SFTP ingest: error closing SSH connection: %s", exc) _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_dropbox_watch_folder() -> dict: """ Celery task: scan the configured Dropbox ingest folder for new files. Uses the existing Dropbox OAuth credentials and DROPBOX_INGEST_FOLDER to poll for new documents. """ if not getattr(settings, "dropbox_ingest_enabled", False): return {"status": "skipped", "reason": "DROPBOX_INGEST_ENABLED is False"} ingest_folder: str | None = getattr(settings, "dropbox_ingest_folder", None) if not ingest_folder: logger.warning("Dropbox ingest enabled but DROPBOX_INGEST_FOLDER is not set — skipping.") return {"status": "skipped", "reason": "DROPBOX_INGEST_FOLDER not configured"} delete_after: bool = getattr(settings, "dropbox_ingest_delete_after_process", False) cache = _load_cache(DROPBOX_INGEST_CACHE_FILE) n = _scan_dropbox_folder(ingest_folder, cache, delete_after) _save_cache(DROPBOX_INGEST_CACHE_FILE, cache) logger.info("Dropbox ingest: %d new file(s) enqueued from %s.", n, ingest_folder) return {"status": "ok", "files_enqueued": n, "folder": ingest_folder} @shared_task def scan_google_drive_watch_folder() -> dict: """ Celery task: scan the configured Google Drive ingest folder for new files. Uses the existing Google Drive credentials and GOOGLE_DRIVE_INGEST_FOLDER_ID to poll for new documents. """ if not getattr(settings, "google_drive_ingest_enabled", False): return {"status": "skipped", "reason": "GOOGLE_DRIVE_INGEST_ENABLED is False"} folder_id: str | None = getattr(settings, "google_drive_ingest_folder_id", None) if not folder_id: logger.warning("Google Drive ingest enabled but GOOGLE_DRIVE_INGEST_FOLDER_ID is not set — skipping.") return {"status": "skipped", "reason": "GOOGLE_DRIVE_INGEST_FOLDER_ID not configured"} delete_after: bool = getattr(settings, "google_drive_ingest_delete_after_process", False) cache = _load_cache(GDRIVE_INGEST_CACHE_FILE) n = _scan_google_drive_folder(folder_id, cache, delete_after) _save_cache(GDRIVE_INGEST_CACHE_FILE, cache) logger.info("Google Drive ingest: %d new file(s) enqueued from folder %s.", n, folder_id) return {"status": "ok", "files_enqueued": n, "folder_id": folder_id} @shared_task def scan_onedrive_watch_folder() -> dict: """ Celery task: scan the configured OneDrive ingest folder for new files. Uses the existing OneDrive MSAL credentials and ONEDRIVE_INGEST_FOLDER_PATH to poll for new documents. """ if not getattr(settings, "onedrive_ingest_enabled", False): return {"status": "skipped", "reason": "ONEDRIVE_INGEST_ENABLED is False"} folder_path: str | None = getattr(settings, "onedrive_ingest_folder_path", None) if not folder_path: logger.warning("OneDrive ingest enabled but ONEDRIVE_INGEST_FOLDER_PATH is not set — skipping.") return {"status": "skipped", "reason": "ONEDRIVE_INGEST_FOLDER_PATH not configured"} delete_after: bool = getattr(settings, "onedrive_ingest_delete_after_process", False) cache = _load_cache(ONEDRIVE_INGEST_CACHE_FILE) n = _scan_onedrive_folder(folder_path, cache, delete_after) _save_cache(ONEDRIVE_INGEST_CACHE_FILE, cache) logger.info("OneDrive ingest: %d new file(s) enqueued from %s.", n, folder_path) return {"status": "ok", "files_enqueued": n, "folder": folder_path} @shared_task def scan_nextcloud_watch_folder() -> dict: """ Celery task: scan the configured Nextcloud ingest folder for new files. Uses the existing Nextcloud WebDAV credentials and NEXTCLOUD_INGEST_FOLDER to poll for new documents. """ if not getattr(settings, "nextcloud_ingest_enabled", False): return {"status": "skipped", "reason": "NEXTCLOUD_INGEST_ENABLED is False"} ingest_folder: str | None = getattr(settings, "nextcloud_ingest_folder", None) if not ingest_folder: logger.warning("Nextcloud ingest enabled but NEXTCLOUD_INGEST_FOLDER is not set — skipping.") return {"status": "skipped", "reason": "NEXTCLOUD_INGEST_FOLDER not configured"} delete_after: bool = getattr(settings, "nextcloud_ingest_delete_after_process", False) cache = _load_cache(NEXTCLOUD_INGEST_CACHE_FILE) n = _scan_nextcloud_folder(ingest_folder, cache, delete_after) _save_cache(NEXTCLOUD_INGEST_CACHE_FILE, cache) logger.info("Nextcloud ingest: %d new file(s) enqueued from %s.", n, ingest_folder) return {"status": "ok", "files_enqueued": n, "folder": ingest_folder} @shared_task def scan_s3_watch_folder() -> dict: """ Celery task: scan the configured S3 ingest prefix for new objects. Uses the existing S3/AWS credentials and S3_INGEST_PREFIX to poll for new documents in the configured S3 bucket. """ if not getattr(settings, "s3_ingest_enabled", False): return {"status": "skipped", "reason": "S3_INGEST_ENABLED is False"} ingest_prefix: str | None = getattr(settings, "s3_ingest_prefix", None) if not ingest_prefix: logger.warning("S3 ingest enabled but S3_INGEST_PREFIX is not set — skipping.") return {"status": "skipped", "reason": "S3_INGEST_PREFIX not configured"} delete_after: bool = getattr(settings, "s3_ingest_delete_after_process", False) cache = _load_cache(S3_INGEST_CACHE_FILE) n = _scan_s3_prefix(ingest_prefix, cache, delete_after) _save_cache(S3_INGEST_CACHE_FILE, cache) logger.info("S3 ingest: %d new file(s) enqueued from prefix %s.", n, ingest_prefix) return {"status": "ok", "files_enqueued": n, "prefix": ingest_prefix} @shared_task def scan_webdav_watch_folder() -> dict: """ Celery task: scan the configured WebDAV ingest folder for new files. Uses the existing WebDAV URL/credentials and WEBDAV_INGEST_FOLDER to poll for new documents. """ if not getattr(settings, "webdav_ingest_enabled", False): return {"status": "skipped", "reason": "WEBDAV_INGEST_ENABLED is False"} ingest_folder: str | None = getattr(settings, "webdav_ingest_folder", None) if not ingest_folder: logger.warning("WebDAV ingest enabled but WEBDAV_INGEST_FOLDER is not set — skipping.") return {"status": "skipped", "reason": "WEBDAV_INGEST_FOLDER not configured"} delete_after: bool = getattr(settings, "webdav_ingest_delete_after_process", False) cache = _load_cache(WEBDAV_INGEST_CACHE_FILE) n = _scan_webdav_folder(ingest_folder, cache, delete_after) _save_cache(WEBDAV_INGEST_CACHE_FILE, cache) logger.info("WebDAV ingest: %d new file(s) enqueued from %s.", n, ingest_folder) return {"status": "ok", "files_enqueued": n, "folder": ingest_folder} # --------------------------------------------------------------------------- # Per-user watch folder integration scanning # --------------------------------------------------------------------------- def _is_safe_watch_path(folder_path: str) -> bool: """Validate that a user-configured watch folder path is safe. Rejects paths that attempt directory traversal (``..``), use relative references, or are not absolute. This prevents a malicious user from configuring a watch folder that could escape its intended directory. Args: folder_path: The path to validate. Returns: ``True`` if the path is considered safe, ``False`` otherwise. """ if not folder_path: return False # Must be absolute if not os.path.isabs(folder_path): logger.warning("Rejecting non-absolute watch folder path: %s", folder_path) return False # Resolve to canonical path and ensure no traversal components exist resolved = os.path.realpath(folder_path) if ".." in folder_path.split(os.sep): logger.warning("Rejecting path with traversal components: %s", folder_path) return False # Ensure resolved path matches the original intent (no symlink escapes) if resolved != os.path.normpath(folder_path): logger.warning( "Watch folder path resolves differently (possible symlink escape): %s -> %s", folder_path, resolved, ) return False return True def _scan_user_watch_folder( folder_path: str, cache: dict[str, str], delete_after: bool, owner_id: str, ) -> int: """Scan a user-configured local watch folder, attributing files to *owner_id*. Delegates to the same file-scanning logic as :func:`_scan_local_folder` but passes ``owner_id`` to :func:`_enqueue_file` so that ingested documents are correctly attributed to the user. Args: folder_path: Absolute directory path to scan. cache: In-memory dict of already-processed file keys. delete_after: Whether to remove the source file after ingestion. owner_id: The user to attribute ingested documents to. Returns: Number of files newly enqueued. """ if not os.path.isdir(folder_path): logger.warning("User 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 user 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): continue abs_path = entry.path if abs_path in cache: continue dest_filename = f"uwf_{owner_id}_{entry.name}" dest_path = os.path.join(settings.workdir, dest_filename) 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, owner_id=owner_id) _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 # --------------------------------------------------------------------------- # Per-user cloud source watch folder scanning # --------------------------------------------------------------------------- # Maps source_type values to their per-user scan functions. _USER_WF_CLOUD_HANDLERS: dict[str, Any] = {} # populated after function defs def _scan_user_s3_folder( cfg: dict, creds: dict, cache: dict[str, str], delete_after: bool, owner_id: str, ) -> int: """Scan an S3 bucket prefix using per-user credentials. Args: cfg: Integration config with ``bucket``, ``region``, ``prefix``, ``endpoint_url``. creds: Decrypted credentials with ``access_key_id``, ``secret_access_key``. cache: In-memory dict of already-processed file keys. delete_after: Whether to remove the source object after ingestion. owner_id: The user to attribute ingested documents to. Returns: Number of files newly enqueued. """ try: import boto3 from botocore.exceptions import ClientError except ImportError as exc: logger.error("User S3 watch folder: boto3 not installed: %s", exc) return 0 bucket = cfg.get("bucket", "") prefix = cfg.get("prefix", "") region = cfg.get("region", "us-east-1") endpoint_url = cfg.get("endpoint_url") or None if not bucket: logger.warning("User S3 watch folder: bucket not configured.") return 0 access_key = creds.get("access_key_id", "") secret_key = creds.get("secret_access_key", "") if not (access_key and secret_key): logger.warning("User S3 watch folder: credentials incomplete.") return 0 try: client_kwargs: dict = { "region_name": region, "aws_access_key_id": access_key, "aws_secret_access_key": secret_key, } if endpoint_url: client_kwargs["endpoint_url"] = endpoint_url s3 = boto3.client("s3", **client_kwargs) except Exception as exc: logger.error("User S3 watch folder: failed to create client: %s", exc) return 0 count = 0 paginator = s3.get_paginator("list_objects_v2") try: pages = paginator.paginate(Bucket=bucket, Prefix=prefix) except Exception as exc: logger.error("User S3 watch folder: failed to list %s/%s: %s", bucket, prefix, exc) return 0 for page in pages: for obj in page.get("Contents", []): key = obj["Key"] filename = key.split("/")[-1] if not filename or not _is_allowed_file(filename): continue cache_key = f"s3:{bucket}/{key}" if cache_key in cache: continue dest_path = os.path.join(settings.workdir, f"uwf_s3_{owner_id}_{filename}") if os.path.exists(dest_path): base2, ext2 = os.path.splitext(f"uwf_s3_{owner_id}_{filename}") dest_path = os.path.join(settings.workdir, f"{base2}_{int(datetime.now().timestamp())}{ext2}") try: s3.download_file(bucket, key, dest_path) logger.info("User S3 watch folder: downloaded s3://%s/%s", bucket, key) except ClientError as exc: logger.error("User S3 watch folder: failed to download %s: %s", key, exc) if os.path.exists(dest_path): os.remove(dest_path) continue _enqueue_file(dest_path, filename=filename, owner_id=owner_id) _mark_processed(cache, cache_key) count += 1 if delete_after: try: s3.delete_object(Bucket=bucket, Key=key) logger.info("User S3 watch folder: deleted s3://%s/%s", bucket, key) except Exception as exc: logger.warning("User S3 watch folder: could not delete %s: %s", key, exc) return count def _scan_user_dropbox_folder( cfg: dict, creds: dict, cache: dict[str, str], delete_after: bool, owner_id: str, ) -> int: """Scan a Dropbox folder using per-user credentials. Args: cfg: Integration config with ``folder_path``. creds: Decrypted credentials with ``refresh_token``, ``app_key``, ``app_secret``. cache: In-memory dict of already-processed file keys. delete_after: Whether to remove the source file after ingestion. owner_id: The user to attribute ingested documents to. Returns: Number of files newly enqueued. """ try: import dropbox as dropbox_module except ImportError as exc: logger.error("User Dropbox watch folder: dropbox SDK not installed: %s", exc) return 0 refresh_token = creds.get("refresh_token", "") app_key = creds.get("app_key", "") app_secret = creds.get("app_secret", "") if not (refresh_token and app_key and app_secret): logger.warning("User Dropbox watch folder: credentials incomplete.") return 0 try: dbx = dropbox_module.Dropbox( oauth2_refresh_token=refresh_token, app_key=app_key, app_secret=app_secret, ) except Exception as exc: logger.error("User Dropbox watch folder: auth failed: %s", exc) return 0 folder_path = cfg.get("folder_path", "") if not folder_path: logger.warning("User Dropbox watch folder: folder_path not configured.") return 0 count = 0 try: result = dbx.files_list_folder(folder_path) entries = list(result.entries) while result.has_more: result = dbx.files_list_folder_continue(result.cursor) entries.extend(result.entries) except Exception as exc: logger.error("User Dropbox watch folder: cannot list %s: %s", folder_path, exc) return 0 for entry in entries: if not isinstance(entry, dropbox_module.files.FileMetadata): continue filename = entry.name if not _is_allowed_file(filename): continue cache_key = f"dropbox:{entry.id}" if cache_key in cache: continue dest_path = os.path.join(settings.workdir, f"uwf_dbx_{owner_id}_{filename}") if os.path.exists(dest_path): base2, ext2 = os.path.splitext(f"uwf_dbx_{owner_id}_{filename}") dest_path = os.path.join(settings.workdir, f"{base2}_{int(datetime.now().timestamp())}{ext2}") try: _meta, response = dbx.files_download(entry.path_lower) with open(dest_path, "wb") as f: f.write(response.content) logger.info("User Dropbox watch folder: downloaded %s", filename) except Exception as exc: logger.error("User Dropbox watch folder: failed to download %s: %s", filename, exc) if os.path.exists(dest_path): os.remove(dest_path) continue _enqueue_file(dest_path, filename=filename, owner_id=owner_id) _mark_processed(cache, cache_key) count += 1 if delete_after: try: dbx.files_delete_v2(entry.path_lower) logger.info("User Dropbox watch folder: deleted %s", entry.path_lower) except Exception as exc: logger.warning("User Dropbox watch folder: could not delete %s: %s", entry.path_lower, exc) return count def _scan_user_google_drive_folder( cfg: dict, creds: dict, cache: dict[str, str], delete_after: bool, owner_id: str, ) -> int: """Scan a Google Drive folder using per-user service-account credentials. Args: cfg: Integration config with ``folder_id``. creds: Decrypted credentials with ``credentials_json``. cache: In-memory dict of already-processed file keys. delete_after: Whether to remove the source file after ingestion. owner_id: The user to attribute ingested documents to. Returns: Number of files newly enqueued. """ try: from google.oauth2 import service_account from googleapiclient.discovery import build except ImportError as exc: logger.error("User Google Drive watch folder: SDK not installed: %s", exc) return 0 creds_json = creds.get("credentials_json", "") if not creds_json: logger.warning("User Google Drive watch folder: credentials_json not provided.") return 0 folder_id = cfg.get("folder_id", "") if not folder_id: logger.warning("User Google Drive watch folder: folder_id not configured.") return 0 try: import json as _json info = _json.loads(creds_json) if isinstance(creds_json, str) else creds_json credentials = service_account.Credentials.from_service_account_info( info, scopes=["https://www.googleapis.com/auth/drive"], ) service = build("drive", "v3", credentials=credentials) except Exception as exc: logger.error("User Google Drive watch folder: auth failed: %s", exc) return 0 count = 0 query = f"'{folder_id}' in parents and trashed = false and mimeType != 'application/vnd.google-apps.folder'" page_token = None while True: try: params: dict = { "q": query, "fields": "nextPageToken, files(id, name, mimeType)", "pageSize": 100, } if page_token: params["pageToken"] = page_token response = service.files().list(**params).execute() except Exception as exc: logger.error("User Google Drive watch folder: listing %s failed: %s", folder_id, exc) break for file_meta in response.get("files", []): file_id_gd = file_meta["id"] filename = file_meta["name"] if not _is_allowed_file(filename): continue cache_key = f"gdrive:{file_id_gd}" if cache_key in cache: continue dest_path = os.path.join(settings.workdir, f"uwf_gd_{owner_id}_{filename}") if os.path.exists(dest_path): base2, ext2 = os.path.splitext(f"uwf_gd_{owner_id}_{filename}") dest_path = os.path.join(settings.workdir, f"{base2}_{int(datetime.now().timestamp())}{ext2}") try: import io from googleapiclient.http import MediaIoBaseDownload request = service.files().get_media(fileId=file_id_gd) buf = io.BytesIO() downloader = MediaIoBaseDownload(buf, request) done = False while not done: _, done = downloader.next_chunk() with open(dest_path, "wb") as f: f.write(buf.getvalue()) logger.info("User Google Drive watch folder: downloaded %s", filename) except Exception as exc: logger.error("User Google Drive watch folder: download %s failed: %s", filename, exc) if os.path.exists(dest_path): os.remove(dest_path) continue _enqueue_file(dest_path, filename=filename, owner_id=owner_id) _mark_processed(cache, cache_key) count += 1 if delete_after: try: service.files().delete(fileId=file_id_gd).execute() logger.info("User Google Drive watch folder: deleted %s", filename) except Exception as exc: logger.warning("User Google Drive watch folder: could not delete %s: %s", filename, exc) page_token = response.get("nextPageToken") if not page_token: break return count def _scan_user_onedrive_folder( cfg: dict, creds: dict, cache: dict[str, str], delete_after: bool, owner_id: str, ) -> int: """Scan a OneDrive folder using per-user OAuth credentials. Args: cfg: Integration config with ``folder_path``. creds: Decrypted credentials with ``refresh_token``, ``client_id``, ``client_secret``. cache: In-memory dict of already-processed file keys. delete_after: Whether to remove the source file after ingestion. owner_id: The user to attribute ingested documents to. Returns: Number of files newly enqueued. """ import requests as req_lib refresh_token = creds.get("refresh_token", "") client_id = creds.get("client_id", "") client_secret = creds.get("client_secret", "") if not (refresh_token and client_id and client_secret): logger.warning("User OneDrive watch folder: credentials incomplete.") return 0 folder_path = cfg.get("folder_path", "") if not folder_path: logger.warning("User OneDrive watch folder: folder_path not configured.") return 0 # Exchange refresh token for an access token try: token_resp = req_lib.post( "https://login.microsoftonline.com/common/oauth2/v2.0/token", data={ "grant_type": "refresh_token", "refresh_token": refresh_token, "client_id": client_id, "client_secret": client_secret, "scope": "https://graph.microsoft.com/.default", }, timeout=getattr(settings, "http_request_timeout", 120), ) token_resp.raise_for_status() access_token = token_resp.json()["access_token"] except Exception as exc: logger.error("User OneDrive watch folder: token exchange failed: %s", exc) return 0 headers = {"Authorization": f"Bearer {access_token}"} import urllib.parse encoded_path = urllib.parse.quote(folder_path.lstrip("/")) list_url: str | None = f"https://graph.microsoft.com/v1.0/me/drive/root:/{encoded_path}:/children" count = 0 timeout = getattr(settings, "http_request_timeout", 120) while list_url: try: resp = req_lib.get(list_url, headers=headers, timeout=timeout) resp.raise_for_status() data = resp.json() except Exception as exc: logger.error("User OneDrive watch folder: listing %s failed: %s", folder_path, exc) break for item in data.get("value", []): if "folder" in item: continue filename = item["name"] item_id = item["id"] if not _is_allowed_file(filename): continue cache_key = f"onedrive:{item_id}" if cache_key in cache: continue download_url = item.get("@microsoft.graph.downloadUrl") if not download_url: continue dest_path = os.path.join(settings.workdir, f"uwf_od_{owner_id}_{filename}") if os.path.exists(dest_path): base2, ext2 = os.path.splitext(f"uwf_od_{owner_id}_{filename}") dest_path = os.path.join(settings.workdir, f"{base2}_{int(datetime.now().timestamp())}{ext2}") try: dl_resp = req_lib.get(download_url, headers=headers, timeout=timeout) dl_resp.raise_for_status() with open(dest_path, "wb") as f: f.write(dl_resp.content) logger.info("User OneDrive watch folder: downloaded %s", filename) except Exception as exc: logger.error("User OneDrive watch folder: download %s failed: %s", filename, exc) if os.path.exists(dest_path): os.remove(dest_path) continue _enqueue_file(dest_path, filename=filename, owner_id=owner_id) _mark_processed(cache, cache_key) count += 1 if delete_after: try: del_resp = req_lib.delete( f"https://graph.microsoft.com/v1.0/me/drive/items/{item_id}", headers=headers, timeout=timeout, ) del_resp.raise_for_status() logger.info("User OneDrive watch folder: deleted %s", filename) except Exception as exc: logger.warning("User OneDrive watch folder: could not delete %s: %s", filename, exc) list_url = data.get("@odata.nextLink") return count def _scan_user_nextcloud_folder( cfg: dict, creds: dict, cache: dict[str, str], delete_after: bool, owner_id: str, ) -> int: """Scan a Nextcloud folder using per-user WebDAV credentials. Args: cfg: Integration config with ``url``, ``folder_path``. creds: Decrypted credentials with ``username``, ``password``. cache: In-memory dict of already-processed file keys. delete_after: Whether to remove the source file after ingestion. owner_id: The user to attribute ingested documents to. Returns: Number of files newly enqueued. """ import defusedxml.ElementTree as ET import requests as req_lib from requests.auth import HTTPBasicAuth nc_url = cfg.get("url", "") folder_path = cfg.get("folder_path", "") nc_user = creds.get("username", "") nc_pass = creds.get("password", "") if not (nc_url and nc_user and nc_pass): logger.warning("User Nextcloud watch folder: connection settings incomplete.") return 0 auth = HTTPBasicAuth(nc_user, nc_pass) timeout = getattr(settings, "http_request_timeout", 120) base = nc_url.rstrip("/") folder = folder_path.strip("/") propfind_url = f"{base}/{folder}/" if folder else f"{base}/" try: resp = req_lib.request( "PROPFIND", propfind_url, auth=auth, headers={"Depth": "1", "Content-Type": "application/xml"}, timeout=timeout, ) resp.raise_for_status() except Exception as exc: logger.error("User Nextcloud watch folder: PROPFIND on %s failed: %s", propfind_url, exc) return 0 count = 0 try: root = ET.fromstring(resp.text) # noqa: S314 — defusedxml is safe except Exception as exc: logger.error("User Nextcloud watch folder: failed to parse response: %s", exc) return 0 ns = {"d": "DAV:"} for response_el in root.findall("d:response", ns): href_el = response_el.find("d:href", ns) if href_el is None or href_el.text is None: continue href = href_el.text if href.rstrip("/").endswith(folder.rstrip("/")): continue import urllib.parse filename = urllib.parse.unquote(href.rstrip("/").split("/")[-1]) if not _is_allowed_file(filename): continue cache_key = f"nextcloud:{href}" if cache_key in cache: continue if href.startswith("http"): file_url = href else: from urllib.parse import urlparse parsed = urlparse(nc_url) file_url = f"{parsed.scheme}://{parsed.netloc}{href}" dest_path = os.path.join(settings.workdir, f"uwf_nc_{owner_id}_{filename}") if os.path.exists(dest_path): base_name, ext2 = os.path.splitext(f"uwf_nc_{owner_id}_{filename}") dest_path = os.path.join(settings.workdir, f"{base_name}_{int(datetime.now().timestamp())}{ext2}") try: dl = req_lib.get(file_url, auth=auth, timeout=timeout) dl.raise_for_status() with open(dest_path, "wb") as f: f.write(dl.content) logger.info("User Nextcloud watch folder: downloaded %s", filename) except Exception as exc: logger.error("User Nextcloud watch folder: download %s failed: %s", filename, exc) if os.path.exists(dest_path): os.remove(dest_path) continue _enqueue_file(dest_path, filename=filename, owner_id=owner_id) _mark_processed(cache, cache_key) count += 1 if delete_after: try: del_resp = req_lib.request("DELETE", file_url, auth=auth, timeout=timeout) del_resp.raise_for_status() logger.info("User Nextcloud watch folder: deleted %s", filename) except Exception as exc: logger.warning("User Nextcloud watch folder: could not delete %s: %s", filename, exc) return count def _scan_user_webdav_folder( cfg: dict, creds: dict, cache: dict[str, str], delete_after: bool, owner_id: str, ) -> int: """Scan a WebDAV folder using per-user credentials. Args: cfg: Integration config with ``url``, ``folder_path``. creds: Decrypted credentials with ``username``, ``password``. cache: In-memory dict of already-processed file keys. delete_after: Whether to remove the source file after ingestion. owner_id: The user to attribute ingested documents to. Returns: Number of files newly enqueued. """ import defusedxml.ElementTree as ET import requests as req_lib from requests.auth import HTTPBasicAuth webdav_url = cfg.get("url", "") folder_path = cfg.get("folder_path", "") dav_user = creds.get("username", "") dav_pass = creds.get("password", "") if not webdav_url: logger.warning("User WebDAV watch folder: URL not configured.") return 0 base = webdav_url.rstrip("/") folder = folder_path.strip("/") propfind_url = f"{base}/{folder}/" if folder else f"{base}/" auth = HTTPBasicAuth(dav_user, dav_pass) if dav_user else None timeout = getattr(settings, "http_request_timeout", 120) try: resp = req_lib.request( "PROPFIND", propfind_url, auth=auth, headers={"Depth": "1"}, timeout=timeout, ) resp.raise_for_status() except Exception as exc: logger.error("User WebDAV watch folder: PROPFIND on %s failed: %s", propfind_url, exc) return 0 count = 0 try: root = ET.fromstring(resp.text) # noqa: S314 — defusedxml is safe except Exception as exc: logger.error("User WebDAV watch folder: failed to parse response: %s", exc) return 0 from urllib.parse import unquote, urlparse ns = {"d": "DAV:"} for response_el in root.findall("d:response", ns): href_el = response_el.find("d:href", ns) if href_el is None or href_el.text is None: continue href = href_el.text if href.endswith("/"): continue filename = unquote(href.split("/")[-1]) if not _is_allowed_file(filename): continue cache_key = f"webdav:{href}" if cache_key in cache: continue if href.startswith("http"): file_url = href else: parsed = urlparse(webdav_url) file_url = f"{parsed.scheme}://{parsed.netloc}{href}" dest_path = os.path.join(settings.workdir, f"uwf_dav_{owner_id}_{filename}") if os.path.exists(dest_path): base2, ext2 = os.path.splitext(f"uwf_dav_{owner_id}_{filename}") dest_path = os.path.join(settings.workdir, f"{base2}_{int(datetime.now().timestamp())}{ext2}") try: dl = req_lib.get(file_url, auth=auth, timeout=timeout) dl.raise_for_status() with open(dest_path, "wb") as f: f.write(dl.content) logger.info("User WebDAV watch folder: downloaded %s", filename) except Exception as exc: logger.error("User WebDAV watch folder: download %s failed: %s", filename, exc) if os.path.exists(dest_path): os.remove(dest_path) continue _enqueue_file(dest_path, filename=filename, owner_id=owner_id) _mark_processed(cache, cache_key) count += 1 if delete_after: try: del_resp = req_lib.request("DELETE", file_url, auth=auth, timeout=timeout) del_resp.raise_for_status() logger.info("User WebDAV watch folder: deleted %s", filename) except Exception as exc: logger.warning("User WebDAV watch folder: could not delete %s: %s", filename, exc) return count # Populate the cloud handler dispatch table _USER_WF_CLOUD_HANDLERS.update( { "s3": _scan_user_s3_folder, "dropbox": _scan_user_dropbox_folder, "google_drive": _scan_user_google_drive_folder, "onedrive": _scan_user_onedrive_folder, "nextcloud": _scan_user_nextcloud_folder, "webdav": _scan_user_webdav_folder, } ) def _pull_user_integration_watch_folders() -> dict: """Iterate over all active WATCH_FOLDER UserIntegrations and scan their sources. Polls the ``user_integrations`` table for records with ``integration_type='WATCH_FOLDER'``, ``direction='SOURCE'``, and ``is_active=True``. Each integration's config is decoded and the configured source is scanned for new files, which are enqueued with the owning user's ``owner_id``. Path traversal protection is enforced on local filesystem paths. Cloud source types (S3, Dropbox, Google Drive, OneDrive, Nextcloud, WebDAV) are dispatched to their per-user scanning helpers. Individual integration failures are caught and recorded without crashing the polling loop. Returns: Summary dict with ``status`` and ``integrations_processed`` count. """ try: import json as _json from app.models import IntegrationDirection, IntegrationType, UserIntegration from app.utils.encryption import decrypt_value db = _get_db_session() try: integrations = ( db.query(UserIntegration) .filter( UserIntegration.integration_type == IntegrationType.WATCH_FOLDER, UserIntegration.direction == IntegrationDirection.SOURCE, UserIntegration.is_active.is_(True), ) .all() ) logger.info("Processing %d WATCH_FOLDER UserIntegration(s)", len(integrations)) total_files = 0 for integ in integrations: try: cfg = _json.loads(integ.config) if integ.config else {} delete_after = cfg.get("delete_after_process", False) source_type = cfg.get("source_type", "local") cache_file = f"{_USER_WF_CACHE_PREFIX}{integ.id}.json" cache = _load_cache(cache_file) if source_type == "local": # Local filesystem watch folder (original behaviour) folder_path = cfg.get("folder_path", "") if not folder_path: logger.warning( "Watch folder integration %d (owner %s) has no folder_path — skipping.", integ.id, integ.owner_id, ) continue if not _is_safe_watch_path(folder_path): error_msg = f"Unsafe watch folder path rejected: {folder_path}" logger.error( "Watch folder integration %d (owner %s): %s", integ.id, integ.owner_id, error_msg, ) integ.last_error = error_msg[:_MAX_ERROR_LENGTH] db.commit() continue n = _scan_user_watch_folder(folder_path, cache, delete_after, integ.owner_id) elif source_type in _USER_WF_CLOUD_HANDLERS: # Cloud source — decrypt per-user credentials and delegate raw_creds = decrypt_value(integ.credentials) if integ.credentials else None creds = _json.loads(raw_creds) if raw_creds else {} handler = _USER_WF_CLOUD_HANDLERS[source_type] n = handler(cfg, creds, cache, delete_after, integ.owner_id) else: logger.warning( "Watch folder integration %d (owner %s): unknown source_type '%s' — skipping.", integ.id, integ.owner_id, source_type, ) continue _save_cache(cache_file, cache) total_files += n integ.last_used_at = datetime.now(timezone.utc) integ.last_error = None db.commit() logger.info( "Watch folder integration %d (owner %s): %d file(s) enqueued (source=%s)", integ.id, integ.owner_id, n, source_type, ) except Exception as exc: # noqa: BLE001 error_msg = str(exc)[:_MAX_ERROR_LENGTH] logger.error( "Error scanning watch folder integration %d (owner %s): %s", integ.id, integ.owner_id, error_msg, ) try: integ.last_used_at = datetime.now(timezone.utc) integ.last_error = error_msg db.commit() except Exception: # noqa: BLE001 db.rollback() return {"status": "ok", "integrations_processed": len(integrations), "files_enqueued": total_files} finally: db.close() except Exception as exc: # noqa: BLE001 logger.error("Failed to process WATCH_FOLDER UserIntegrations: %s", exc) return {"status": "error", "reason": str(exc)[:_MAX_ERROR_LENGTH]} @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 scans: 1. Local filesystem watch folders 2. FTP ingest folder (if enabled) 3. SFTP ingest folder (if enabled) 4. Dropbox ingest folder (if enabled) 5. Google Drive ingest folder (if enabled) 6. OneDrive ingest folder (if enabled) 7. Nextcloud ingest folder (if enabled) 8. Amazon S3 ingest prefix (if enabled) 9. WebDAV ingest folder (if enabled) 10. Per-user WATCH_FOLDER integrations from the database """ 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() results["dropbox"] = scan_dropbox_watch_folder() results["google_drive"] = scan_google_drive_watch_folder() results["onedrive"] = scan_onedrive_watch_folder() results["nextcloud"] = scan_nextcloud_watch_folder() results["s3"] = scan_s3_watch_folder() results["webdav"] = scan_webdav_watch_folder() results["user_watch_folders"] = _pull_user_integration_watch_folders() finally: _release_lock(WATCH_FOLDER_LOCK_KEY) return {"status": "ok", "results": results}