From e7d1289f1b96168460e60612d23bb9e0b644d6e4 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Mon, 9 Mar 2026 02:15:03 +0000 Subject: [PATCH] feat(integrations): add IMAP Gmail labels/delete and Watch Folder cloud source support MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add delete_after_process and gmail_apply_labels checkboxes to IMAP integration UI - Add source_type selector (Local/S3/Dropbox/Google Drive/OneDrive/Nextcloud/WebDAV) to Watch Folder UI - Add per-source-type config and credential fields for watch folders - Add delete_after_process checkbox to Watch Folder settings - Fix field name mismatch (path → folder_path) in Watch Folder UI - Update pull_inbox() to accept and respect gmail_apply_labels parameter - Update _pull_user_integration_imap() to read gmail_apply_labels from config - Add per-user cloud scanning functions for S3, Dropbox, Google Drive, OneDrive, Nextcloud, WebDAV - Update _pull_user_integration_watch_folders() to dispatch based on source_type - Update IMAP and WATCH_FOLDER config shape documentation in models.py Co-authored-by: christianlouis <361235+christianlouis@users.noreply.github.com> --- app/models.py | 48 +- app/tasks/imap_tasks.py | 20 +- app/tasks/watch_folder_tasks.py | 782 +++++++++++++++++- .../templates/integrations_dashboard.html | 223 ++++- 4 files changed, 1037 insertions(+), 36 deletions(-) diff --git a/app/models.py b/app/models.py index e782c494..4e5bd6ef 100644 --- a/app/models.py +++ b/app/models.py @@ -544,9 +544,55 @@ class UserIntegration(Base): IMAP: config = {"host": "imap.example.com", "port": 993, "username": "user@example.com", "use_ssl": true, - "delete_after_process": false} + "delete_after_process": false, + "gmail_apply_labels": true} credentials = {"password": "secret"} + WATCH_FOLDER (local): + config = {"source_type": "local", + "folder_path": "/data/inbox", + "delete_after_process": false} + + WATCH_FOLDER (s3): + config = {"source_type": "s3", "bucket": "my-bucket", + "region": "us-east-1", "prefix": "inbox/", + "endpoint_url": null, "delete_after_process": false} + credentials = {"access_key_id": "AKI…", "secret_access_key": "…"} + + WATCH_FOLDER (dropbox): + config = {"source_type": "dropbox", + "folder_path": "/Inbox/Scanner", + "delete_after_process": false} + credentials = {"refresh_token": "…", "app_key": "…", + "app_secret": "…"} + + WATCH_FOLDER (google_drive): + config = {"source_type": "google_drive", + "folder_id": "1abc…", + "delete_after_process": false} + credentials = {"credentials_json": "{…service-account…}"} + + WATCH_FOLDER (onedrive): + config = {"source_type": "onedrive", + "folder_path": "/Documents/Inbox", + "delete_after_process": false} + credentials = {"refresh_token": "…", "client_id": "…", + "client_secret": "…"} + + WATCH_FOLDER (nextcloud): + config = {"source_type": "nextcloud", + "url": "https://cloud.example.com", + "folder_path": "/Documents/Inbox", + "delete_after_process": false} + credentials = {"username": "user", "password": "secret"} + + WATCH_FOLDER (webdav): + config = {"source_type": "webdav", + "url": "https://webdav.example.com/dav/", + "folder_path": "/remote.php/webdav/Inbox", + "delete_after_process": false} + credentials = {"username": "user", "password": "secret"} + S3: config = {"bucket": "my-bucket", "region": "us-east-1", "endpoint_url": null, "folder_prefix": ""} diff --git a/app/tasks/imap_tasks.py b/app/tasks/imap_tasks.py index 2ab827cf..27c42d86 100644 --- a/app/tasks/imap_tasks.py +++ b/app/tasks/imap_tasks.py @@ -249,6 +249,7 @@ def _pull_user_integration_imap() -> None: password = creds.get("password") use_ssl = cfg.get("use_ssl", True) delete_after = cfg.get("delete_after_process", False) + gmail_labels = cfg.get("gmail_apply_labels", True) if not (host and username and password): logger.warning( @@ -267,6 +268,7 @@ def _pull_user_integration_imap() -> None: use_ssl=use_ssl, delete_after_process=delete_after, owner_id=integ.owner_id, + gmail_apply_labels=gmail_labels, ) integ.last_used_at = datetime.now(timezone.utc) integ.last_error = None @@ -317,7 +319,17 @@ def check_and_pull_mailbox( ) -def pull_inbox(mailbox_key, host, port, username, password, use_ssl, delete_after_process, owner_id=None): +def pull_inbox( + mailbox_key, + host, + port, + username, + password, + use_ssl, + delete_after_process, + owner_id=None, + gmail_apply_labels=True, +): """ Connects to the IMAP inbox, fetches new unread emails from the last 3 days, and processes attachments while preserving the original unread status. @@ -331,6 +343,8 @@ def pull_inbox(mailbox_key, host, port, username, password, use_ssl, delete_afte Args: owner_id: Optional user identifier. When provided, ingested documents are attributed to this user via ``process_document`` / ``convert_to_pdf``. + gmail_apply_labels: Whether to apply Gmail-specific labels and stars to + processed emails. Only relevant for Gmail hosts. Defaults to True. """ logger.info("Connecting to %s at %s:%s (SSL=%s)", mailbox_key, host, port, use_ssl) processed_emails = load_processed_emails() @@ -386,7 +400,7 @@ def pull_inbox(mailbox_key, host, port, username, password, use_ssl, delete_afte continue # For Gmail, check if the email already has the "Ingested" label. - if is_gmail_host: + if is_gmail_host and gmail_apply_labels: if email_already_has_label(mail, num, "Ingested"): logger.info("Skipping email %s in %s, already labeled 'Ingested'.", msg_id, mailbox_key) continue @@ -398,7 +412,7 @@ def pull_inbox(mailbox_key, host, port, username, password, use_ssl, delete_afte if settings.imap_readonly_mode: logger.info("Readonly mode: skipping mailbox modifications for %s in %s", msg_id, mailbox_key) else: - if is_gmail_host: + if is_gmail_host and gmail_apply_labels: mark_as_processed_with_star(mail, num) mark_as_processed_with_label(mail, num, label="Ingested") diff --git a/app/tasks/watch_folder_tasks.py b/app/tasks/watch_folder_tasks.py index 4b1dfef0..f01feabd 100644 --- a/app/tasks/watch_folder_tasks.py +++ b/app/tasks/watch_folder_tasks.py @@ -1424,16 +1424,724 @@ def _scan_user_watch_folder( 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, callable] = {} # 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 paths. + """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 ``folder_path`` is scanned for new files, which are enqueued + configured source is scanned for new files, which are enqueued with the owning user's ``owner_id``. - Path traversal protection is enforced on the configured path. + 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. @@ -1445,6 +2153,7 @@ def _pull_user_integration_watch_folders() -> dict: import json as _json from app.models import IntegrationDirection, IntegrationType, UserIntegration + from app.utils.encryption import decrypt_value db = _get_db_session() try: @@ -1462,32 +2171,51 @@ def _pull_user_integration_watch_folders() -> dict: for integ in integrations: try: cfg = _json.loads(integ.config) if integ.config else {} - folder_path = cfg.get("folder_path", "") delete_after = cfg.get("delete_after_process", False) - - 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 + source_type = cfg.get("source_type", "local") cache_file = f"{_USER_WF_CACHE_PREFIX}{integ.id}.json" cache = _load_cache(cache_file) - n = _scan_user_watch_folder(folder_path, cache, delete_after, integ.owner_id) + + 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 @@ -1495,11 +2223,11 @@ def _pull_user_integration_watch_folders() -> dict: integ.last_error = None db.commit() logger.info( - "Watch folder integration %d (owner %s): %d file(s) enqueued from %s", + "Watch folder integration %d (owner %s): %d file(s) enqueued (source=%s)", integ.id, integ.owner_id, n, - folder_path, + source_type, ) except Exception as exc: # noqa: BLE001 error_msg = str(exc)[:_MAX_ERROR_LENGTH] diff --git a/frontend/templates/integrations_dashboard.html b/frontend/templates/integrations_dashboard.html index 5b71aaa7..ea530404 100644 --- a/frontend/templates/integrations_dashboard.html +++ b/frontend/templates/integrations_dashboard.html @@ -461,16 +461,28 @@ -
+ + Gmail labels & star: when enabled, processed emails are starred and tagged with an "Ingested" label in Gmail. Only applies to Gmail servers. +
Watch Folder Settings