diff --git a/.env.example b/.env.example index 2ce63eb..55b3bf5 100644 --- a/.env.example +++ b/.env.example @@ -30,6 +30,11 @@ ACME_HTTP_PORT=80 # WEB_USAGE_SCAN_INTERVAL_SECONDS=900 # WEB_DIRECTORY_CACHE_SECONDS=300 # WEB_LOGIN_ATTEMPTS_PER_5_MIN=10 +# DOCUMENT_SCAN_SECONDS=30 +# DOCUMENT_OCR_LANGUAGE=deu+eng +# DOCUMENT_OCR_TIMEOUT_SECONDS=600 +# DOCUMENT_MAX_FILE_MB=512 +# DOCUMENT_MAX_PDF_PAGES=500 # TRASH_RETENTION_DAYS=7 # LDAP_URI=ldaps://example.com # LDAP_BASE_DN=DC=example,DC=com diff --git a/Dockerfile b/Dockerfile index 2a068b9..4b7534a 100644 --- a/Dockerfile +++ b/Dockerfile @@ -15,6 +15,11 @@ RUN apt-get update \ libpam-winbind \ python3 \ python3-samba \ + python3-pil \ + ocrmypdf \ + poppler-utils \ + tesseract-ocr-deu \ + tesseract-ocr-eng \ p7zip-full \ rclone \ rsync \ @@ -25,6 +30,8 @@ RUN apt-get update \ winbind \ && rm -rf /var/lib/apt/lists/* +RUN useradd --system --no-create-home --shell /usr/sbin/nologin document-ocr + RUN mkdir -p /app /data/private /data/fslogix /data/groups/data /data/groups/archive /state COPY --from=step-cli /usr/local/bin/step /usr/local/bin/step @@ -38,6 +45,9 @@ COPY app/audit_store.py /app/audit_store.py COPY app/audit_collector.py /app/audit_collector.py COPY app/trash.py /app/trash.py COPY app/web_ui.py /app/web_ui.py +COPY app/documents.py /app/documents.py +COPY app/document_index.py /app/document_index.py +COPY app/extract_document.py /app/extract_document.py COPY app/web /app/web COPY app/init.sh /app/init.sh COPY etc/samba/smb.conf /app/smb.conf.template diff --git a/README.md b/README.md index 09b5f1c..06c312f 100644 --- a/README.md +++ b/README.md @@ -25,8 +25,9 @@ This repository provides a production-oriented Samba file server container that - Samba `full_audit` records successful and failed reads, writes, renames, and deletions on all three shares; FSLogix events are retained in a separate indexed activity stream. - A collector normalizes those four actions and persists them in indexed SQLite tables; activity is never automatically deleted. - Samba retains deleted files for seven days in per-user recycle repositories on the same data volumes. -- A plain HTTPS administration console manages Data folders and individual user access, and provides statistics, logs, trash downloads/restores, manual backups, and share reconciliation. It also includes a fully client-side Typst PDF report. -- Web sign-in validates the submitted username/password with Kerberos, permits only users whose winbind group SID set contains `DOMAIN_ADMINS_SID`, and issues an expiring JWT in a Secure, HttpOnly, SameSite=Strict cookie. The browser does not use NTLM/SPNEGO or Kerberos negotiation. +- The HTTPS file interface at `/` lets AD users search accessible Data folders and their own Private folder. An independent background worker indexes text and recognizes scanned PDFs without changing original files. +- The administration console at `/admin/...` manages Data folders and individual user access, and provides statistics, logs, trash downloads/restores, manual backups, and share reconciliation. It also includes a fully client-side Typst PDF report. +- Web sign-in validates the submitted username/password with Kerberos and issues an expiring JWT in a Secure, HttpOnly, SameSite=Strict cookie. Human AD users can sign in; `MSOL_*` and `krbtgt` accounts are excluded. Administration pages and APIs require membership in `DOMAIN_ADMINS_SID`. The browser does not use NTLM/SPNEGO or Kerberos negotiation. - HTTPS certificates are requested from a configured local Smallstep CA and renewed automatically. Pre-issued certificate files are also supported. - Optional remote backups run when `BACKUP_DESTINATION` and `BACKUP_ARCHIVE_PASSWORD` are configured; each active or archived group folder is uploaded as its own encrypted, non-solid 7z archive. - Private home creation skips well-known/service accounts by default (including `krbtgt`, `msol_*`, `FileShare_ServiceAcc`). @@ -68,6 +69,34 @@ A new installation starts without assignments. Existing untracked Data directori The Data share uses Samba Windows ACL checks instead of POSIX ACLs. Keep its data volumes private to the container; direct local or NFS access is outside this permission model. The protected TDB store avoids requiring extra container capabilities. After restoring Data to different filesystem inodes, run reconciliation with `REPAIR_DATA_ACLS=1` to rebuild ACL records. Trash restoration applies the current folder policy before exposing the restored file. +## File Search and PDF Recognition + +Users sign in at `/` with their existing AD credentials using `username`, `DOMAIN\username`, or `username@domain`. Domain-qualified names accept the configured AD DNS domain/realm; the short NetBIOS domain is also accepted after `@`. The file view provides live search by filename and content, folder and file-type filters, grid/list views, previews, and downloads of the original files. Only Domain Admins see the administration link and can use management pages and APIs under `/admin/...`; old administration URLs redirect there. Administrative APIs are under `/admin/api/...`, while sign-in, session, and document endpoints remain under `/api/...`. This upgrade expires existing web sessions. + +The catalog includes active Data folders and each user's own Private folder. Every search, detail, preview, and download request checks current individual folder permissions. Revoking access or archiving a folder takes effect without waiting for reindexing. Domain Admins can see all active Data folders, but the user file view still shows only their own Private folder. Archived folders, recycle repositories, FSLogix profiles, symlinks, nested mounts, and special files are excluded. Private files additionally require ownership and read/traverse permissions for the signed-in user's mapped Unix identity. + +Linux filesystem events update the filename catalog as files change, with periodic scans recovering missed events. Text extraction and OCR run through a persistent background queue; files become searchable by filename before content extraction finishes. Search updates while typing, and open views refresh every two seconds to show completed extraction. Full-text queries match word prefixes and ignore accents; filename-only queries also match substrings. + +Content extraction supports PDFs, UTF-8/UTF-16 text files, modern Office formats (`docx`, `xlsx`, `pptx`), and OpenDocument formats (`odt`, `ods`, `odp`). PDFs and common images receive a first-page/image thumbnail. PDFs open in a fullscreen dialog with the locally bundled Mozilla PDF.js viewer, including page navigation, zoom, thumbnails and search within native PDF text. PDF loading and byte-range requests check current folder access. Other formats, including legacy binary Office files, remain searchable by filename. Extracted text is limited to 2 MiB per file; encrypted, damaged, or oversized documents may have no searchable content. + +PDFs containing raster images and pages with little extracted text are queued for local OCR using OCRmyPDF and Tesseract, with German and English enabled by default. Native text remains searchable while OCR is pending. OCR text is used for search. No file type displays a separate extracted document-text panel in its dialog. The PDF viewer displays the original PDF. Downloads always return the original file. No OCR replacement or separate OCR PDF is published. Failed jobs retry with backoff, and pending work survives restarts. + +Parsers run as an unprivileged service account on temporary copies, with filesystem/network restrictions and resource limits. This requires a Linux kernel with Landlock enabled (Linux 5.13 or newer) on x86-64 or ARM64. If isolation is unavailable, extraction fails closed while filename search remains available; the worker logs the reason. OCR needs no external service or additional container capabilities. + +The derived index and previews live in `/state/documents` (`DOCUMENT_STATE_ROOT` override), separate from the managed-access database. The directory is accessible only to the service. Backups retain all original files and access records, and omit this rebuildable cache. Restoring a backup rebuilds the catalog automatically. Allow space for extracted text, previews, and temporary OCR work; the first indexing pass runs asynchronously. + +Domain Admins monitor the worker at `/admin/documents`: current file and phase, completed/total files, text/OCR queues, retries, heartbeat, last catalog scan, and recent activity. Pause and resume controls stop both catalog updates and extraction. Pausing interrupts an active parser and its child processes, discards its temporary workspace, and keeps the original queue phase for resumption. Existing indexed files remain searchable subject to current access rules. The pause flag survives container restarts. Control events record the administrator in the document activity history; status and control APIs are restricted to administrators. This derived worker history and its pause flag are omitted with the document cache from backups. + +| Variable | Default | Purpose | +| --- | --- | --- | +| `DOCUMENT_SCAN_SECONDS` | `30` | Recovery scan interval; filesystem events handle intervening changes | +| `DOCUMENT_OCR_LANGUAGE` | `deu+eng` | Installed Tesseract languages used for OCR | +| `DOCUMENT_OCR_TIMEOUT_SECONDS` | `600` | Maximum runtime for an OCR job | +| `DOCUMENT_MAX_FILE_MB` | `512` | Largest file copied for extraction; larger files use filename search | +| `DOCUMENT_MAX_PDF_PAGES` | `500` | Maximum PDF page count for extraction | + +These settings are passed through the existing Compose environment file. Setting `WEB_ENABLED=false` also disables the document worker. + ## Shared SQLite State Database The default database path is `/state/shares.db`; `STATE_DB_PATH` can override it. `SHARE_DB_PATH` remains a backward-compatible fallback. @@ -129,6 +158,9 @@ Kerberos requires close time alignment. - `app/audit_store.py` - `app/state_db.py` - `app/web_ui.py` +- `app/documents.py` +- `app/document_index.py` +- `app/extract_document.py` - `app/web/` - `etc/samba/smb.conf` - `dev/` (disposable AD DC, backup target, SMB client, seed data, and E2E assertions) diff --git a/app/backup_to_destination.py b/app/backup_to_destination.py index 60d8ca2..8dc22d6 100644 --- a/app/backup_to_destination.py +++ b/app/backup_to_destination.py @@ -959,7 +959,12 @@ def prepare_state_snapshot( temporary_root = tempfile.mkdtemp(prefix="backup-state-") staged_state = os.path.join(temporary_root, "state") + # The document index and previews are derived from backed-up originals. + # Rebuild them on restore rather than copy a live second WAL database. + document_cache = os.path.abspath(os.getenv("DOCUMENT_STATE_ROOT", os.path.join(state_root, "documents"))) + document_relative = os.path.relpath(document_cache, state_root) excluded = { + document_relative, database_relative, f"{database_relative}-wal", f"{database_relative}-shm", diff --git a/app/document_index.py b/app/document_index.py new file mode 100644 index 0000000..22e8ff7 --- /dev/null +++ b/app/document_index.py @@ -0,0 +1,386 @@ +#!/usr/bin/env python3 +"""Durable extraction/OCR queue with immediate Linux file events and periodic recovery.""" + +import ctypes +import fcntl +import hashlib +import json +import os +from pathlib import Path +import pwd +import select +import shutil +import signal +import sqlite3 +import struct +import subprocess +import sys +import tempfile +import threading +import time + +try: + from . import documents, reconcile_shares as directory +except ImportError: + import documents + import reconcile_shares as directory + + +def log(message): + print(f'[documents] {message}', flush=True) + + +def setting(name, default, minimum=1, maximum=3600): + try: + return min(maximum, max(minimum, int(os.getenv(name, default)))) + except ValueError: + return default + + +class JobPaused(Exception): + """Leave the durable queue unchanged when extraction is interrupted.""" + + +def worker_state(conn, state, current=''): + conn.execute('UPDATE document_worker SET state=?,current_id=? WHERE id=1', (state, current)) + conn.commit() + + +class FileEvents: + """A separate watcher keeps metadata responsive while OCR is busy.""" + MASK = 0x00000004 | 0x00000008 | 0x00000040 | 0x00000080 | 0x00000100 | 0x00000200 | 0x00000400 | 0x00000800 + def __init__(self): + self.libc = ctypes.CDLL(None, use_errno=True) + self.fd = self.libc.inotify_init1(os.O_NONBLOCK | os.O_CLOEXEC) + if self.fd < 0: + raise OSError(ctypes.get_errno(), 'inotify unavailable') + self.watches = {} + self.paths = {} + + def add(self, path, source='', relative=''): + if path in self.paths: + wd = self.paths[path] + else: + wd = self.libc.inotify_add_watch(self.fd, os.fsencode(path), self.MASK | 0x02000000) # DONT_FOLLOW + if wd < 0: + return + self.paths[path] = wd + self.watches[wd] = (path, source, relative) + + def read(self, timeout=0.5): + if not select.select([self.fd], [], [], timeout)[0]: + return [] + try: + data = os.read(self.fd, 1024 * 1024) + except BlockingIOError: + return [] + events, offset = [], 0 + while offset + 16 <= len(data): + wd, mask, cookie, length = struct.unpack_from('iIII', data, offset) + name = os.fsdecode(data[offset+16:offset+16+length].split(b'\0', 1)[0]) + offset += 16 + length + if mask & 0x4000: # Queue overflow: full rescan is mandatory. + events.append(('', '', True)) + continue + watch = self.watches.get(wd) + if not watch: + continue + path, source, relative = watch + if mask & 0x8000: + self.paths.pop(path, None) + self.watches.pop(wd, None) + if name == '.trash': + continue + is_directory = bool(mask & (0x40000000 | 0x400 | 0x800 | 0x8000)) + if source and name and not is_directory: + events.append((source, '/'.join(filter(None, (relative, name))), False)) + else: + events.append((source, '', True)) + return events + + def close(self): + os.close(self.fd) + + +def update_sources(conn): + current = {source['id']: source for source in documents.source_inventory()} + old = {row['id']: dict(row) for row in conn.execute('SELECT * FROM document_sources')} + for key in old.keys() - current.keys(): + conn.execute('DELETE FROM document_sources WHERE id=?', (key,)) + new = set() + for key, source in current.items(): + if old.get(key) != source: + new.add(key) + if key in old and old[key]['root'] != source['root']: + conn.execute('DELETE FROM documents WHERE source_id=?', (key,)) + conn.execute('''INSERT INTO document_sources VALUES(:id,:kind,:label,:root,:owner_sid) + ON CONFLICT(id) DO UPDATE SET kind=excluded.kind,label=excluded.label, + root=excluded.root,owner_sid=excluded.owner_sid''', source) + conn.commit() + return current, new + + +def watch_loop(stop): + conn = documents.connect() + watcher = None + try: + documents.ensure_schema(conn) + try: + watcher = FileEvents() + except OSError: + log('File events unavailable; using periodic scans') + next_sources, next_scan = 0, 0 + sources = {} + last_heartbeat = 0 + def cancelled(): + nonlocal last_heartbeat + now = time.monotonic() + if now - last_heartbeat >= 1: + conn.execute('UPDATE document_worker SET heartbeat=? WHERE id=1', (time.time(),)) + conn.commit() + last_heartbeat = now + return stop.is_set() or documents.worker_paused(conn) + while not stop.is_set(): + try: + if cancelled(): + conn.execute('UPDATE document_worker SET scanning=0 WHERE id=1') + conn.commit() + next_sources = next_scan = 0 + stop.wait(0.5) + continue + now = time.monotonic() + full = now >= next_scan + rescan = set() + if now >= next_sources or full: + sources, added = update_sources(conn) + rescan.update(added) + next_sources = now + 2 + if watcher: + for path in (directory.GROUP_ROOT, directory.PRIVATE_ROOT): + watcher.add(path) + if full: + rescan.update(sources) + next_scan = now + setting('DOCUMENT_SCAN_SECONDS', 30, 2, 3600) + changed = set() + events = watcher.read() if watcher else [] + for source_id, relative, directory_change in events: + if directory_change: + if source_id: + rescan.add(source_id) + else: + next_sources = 0 + if full or relative == '': + next_scan = 0 + else: + changed.add((source_id, relative)) + for source_id in rescan: + if cancelled(): + break + if source_id in sources: + conn.execute('UPDATE document_worker SET scanning=1 WHERE id=1') + conn.commit() + complete = documents.walk_source(conn, sources[source_id], watcher.add if watcher else None, cancelled) + conn.execute('UPDATE document_worker SET scanning=0,last_scan=CASE WHEN ? THEN ? ELSE last_scan END WHERE id=1', (int(complete),time.time())) + for source_id, relative in changed: + if cancelled(): + break + if source_id in sources and source_id not in rescan: + documents.discover_file(conn, sources[source_id], relative) + conn.commit() + except Exception as exc: + conn.rollback() + log(f'Catalog refresh failed: {exc}') + stop.wait(2) + if not watcher: + stop.wait(1) + finally: + if watcher: + watcher.close() + conn.close() + + +def run_job(row, source, phase, cancelled=None): + timeout = setting('DOCUMENT_OCR_TIMEOUT_SECONDS', 600, 30, 3600) if phase == 'ocr' else 120 + job_root = '/tmp/adfs-document-jobs' + os.makedirs(job_root, mode=0o700, exist_ok=True) + user = pwd.getpwnam('document-ocr') + with tempfile.TemporaryDirectory(prefix='job-', dir=job_root) as workspace: + # Only the extractor's own directory is traversable; the private index and + # original trees stay inaccessible. Landlock further confines all children. + os.chmod(job_root, 0o711) + os.chown(workspace, user.pw_uid, user.pw_gid) + config = {'extension': row['extension'], 'phase': phase, + 'language': os.getenv('DOCUMENT_OCR_LANGUAGE', 'deu+eng'), + 'max_pages': setting('DOCUMENT_MAX_PDF_PAGES', 500, 1, 5000), 'timeout': timeout} + for name, content in [('config.json', json.dumps(config))]: + path = os.path.join(workspace, name) + with open(path, 'w') as handle: + handle.write(content) + os.chown(path, user.pw_uid, user.pw_gid) + os.chmod(path, 0o600) + with documents.open_file(source['root'], row['path']) as (handle, info): + if documents.fingerprint(info) != row['fingerprint']: + return None + target = os.path.join(workspace, 'input') + with open(target, 'xb') as copied: + remaining = row['size'] + while remaining: + if cancelled and cancelled(): + raise JobPaused() + chunk = handle.read(min(remaining, 1024 * 1024)) + if not chunk: + return None + copied.write(chunk) + remaining -= len(chunk) + if handle.read(1): + return None + if documents.fingerprint(os.fstat(handle.fileno())) != row['fingerprint']: + return None + os.chown(target, user.pw_uid, user.pw_gid) + os.chmod(target, 0o400) + env = {'PATH': '/usr/local/bin:/usr/bin:/bin', 'LANG': 'C.UTF-8', 'HOME': workspace, + 'TMPDIR': workspace, 'OMP_THREAD_LIMIT': '1', 'PYTHONDONTWRITEBYTECODE': '1'} + script = str(Path(__file__).with_name('extract_document.py')) + process = subprocess.Popen([sys.executable, script, workspace], user=user.pw_uid, group=user.pw_gid, + extra_groups=[], env=env, stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, start_new_session=True, umask=0o077) + try: + deadline = time.monotonic() + timeout + while process.poll() is None: + if cancelled and cancelled(): + raise JobPaused() + if time.monotonic() >= deadline: + raise RuntimeError('Extraction timed out') + try: + process.wait(timeout=0.5) + except subprocess.TimeoutExpired: + pass + finally: + # OCR child processes must not survive cancellation or a parser exit. + try: + os.killpg(process.pid, signal.SIGKILL) + except ProcessLookupError: + pass + process.wait() + with documents.open_file(workspace, 'result.json') as (handle, info): + if info.st_size > documents.MAX_TEXT * 6 + 4096: + raise RuntimeError('Extraction response too large') + result = json.load(handle) + if process.returncode or result.get('error'): + raise RuntimeError(result.get('error', 'Extraction failed')) + result['preview'] = '' + with documents.open_file(source['root'], row['path']) as (_, info): + if documents.fingerprint(info) != row['fingerprint']: + return None + if os.path.isfile(os.path.join(workspace, 'preview.jpg')): + with documents.open_file(workspace, 'preview.jpg') as (handle, info): + if info.st_size > 16 * 1024 * 1024 or handle.read(2) != b'\xff\xd8': + raise RuntimeError('Invalid preview') + handle.seek(0) + previews = os.path.join(documents.SEARCH_ROOT, 'previews') + os.makedirs(previews, mode=0o700, exist_ok=True) + name = uuid_name(row) + with os.fdopen(os.open(os.path.join(previews, name), os.O_WRONLY | os.O_CREAT | os.O_TRUNC | os.O_NOFOLLOW, 0o600), 'wb') as output: + shutil.copyfileobj(handle, output) + os.chmod(os.path.join(previews, name), 0o600) + result['preview'] = name + return result + + +def uuid_name(row): + return hashlib.sha256((row['id'] + row['fingerprint']).encode()).hexdigest()[:32] + '.jpg' + + +def process_next(conn, prefer_ocr=False, stop=None): + if documents.worker_paused(conn): + worker_state(conn, 'paused') + return False + priority = "CASE state WHEN 'ocr' THEN 0 WHEN 'pending' THEN 1 ELSE 2 END" if prefer_ocr else "CASE state WHEN 'pending' THEN 0 WHEN 'ocr' THEN 1 ELSE 2 END" + row = conn.execute(f"SELECT * FROM documents WHERE state IN ('pending','ocr','failed') AND retry_at<=? ORDER BY {priority},modified DESC LIMIT 1", (time.time(),)).fetchone() + if row is None: + worker_state(conn, 'idle') + return False + source = conn.execute('SELECT * FROM document_sources WHERE id=?', (row['source_id'],)).fetchone() + if source is None: + return False + phase = 'ocr' if row['state'] == 'ocr' else 'text' + source_label = source['label'] + (' / ' + source['id'][8:] if source['kind'] == 'private' else '') + worker_state(conn, phase, row['id']) + supported = {'.pdf'} | documents.TEXT_SUFFIXES | documents.OFFICE_SUFFIXES | documents.IMAGE_SUFFIXES + if row['extension'] not in supported or row['size'] > setting('DOCUMENT_MAX_FILE_MB', 512, 1, 4096) * 1024 * 1024: + conn.execute("UPDATE documents SET state='name-only',indexed_at=? WHERE id=? AND fingerprint=?", (time.time(), row['id'], row['fingerprint'])) + conn.commit() + worker_state(conn, 'idle') + return True + documents.worker_event(conn, phase, row['path'], source_label) + conn.commit() + try: + result = run_job(row, source, phase, lambda: documents.worker_paused(conn) or (stop and stop.is_set())) + if result is None: + documents.discover_file(conn, source, row['path']) + documents.worker_event(conn, 'changed', row['path'], source_label) + else: + state = 'ocr' if result.get('needsOcr') else 'ready' + conn.execute('''UPDATE documents SET body=?,state=?,pages=?,preview=CASE WHEN ?='' THEN preview ELSE ? END, + indexed_at=?,attempts=0,retry_at=0 WHERE id=? AND fingerprint=?''', + (str(result.get('body', ''))[:documents.MAX_TEXT], state, result.get('pages', 0), + result.get('preview', ''), result.get('preview', ''), time.time(), row['id'], row['fingerprint'])) + documents.worker_event(conn, 'queued-ocr' if state == 'ocr' else 'complete', row['path'], source_label) + conn.execute('UPDATE document_worker SET processed=processed+1 WHERE id=1') + except JobPaused: + documents.worker_event(conn, 'interrupted', row['path'], source_label) + conn.commit() + worker_state(conn, 'paused' if documents.worker_paused(conn) else 'idle') + return False + except Exception as exc: + attempts = row['attempts'] + 1 + # Preserve the OCR phase across retries and process restarts. + retry_state = 'ocr' if phase == 'ocr' else 'failed' + conn.execute('UPDATE documents SET state=?,attempts=?,retry_at=? WHERE id=? AND fingerprint=?', + (retry_state, attempts, time.time() + min(3600, 30 * 2**min(attempts,7)), row['id'], row['fingerprint'])) + log(f'Extraction deferred for {row["id"]}: {exc}') + documents.worker_event(conn, 'failed', row['path'], source_label, message=exc) + conn.commit() + worker_state(conn, 'idle') + return True + + +def main(): + conn = documents.connect() + documents.ensure_schema(conn) + lock = open(os.path.join(documents.SEARCH_ROOT, 'worker.lock'), 'a') + try: + fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + return 0 + worker_state(conn, 'paused' if documents.worker_paused(conn) else 'idle') + documents.worker_event(conn, 'started') + conn.commit() + stop = threading.Event() + for sig in (signal.SIGTERM, signal.SIGINT): + signal.signal(sig, lambda *_: stop.set()) + watcher = threading.Thread(target=watch_loop, args=(stop,), name='document-watcher', daemon=True) + watcher.start() + log('Watching Data and Private; extraction queue is ready') + processed = 0 + try: + while not stop.is_set(): + try: + # Give scans regular turns while native files are initially indexed. + processed += 1 + if not process_next(conn, prefer_ocr=processed % 5 == 0, stop=stop): + stop.wait(0.5) + except sqlite3.Error as exc: + conn.rollback() + log(f'Queue retry: {exc}') + stop.wait(1) + finally: + stop.set() + watcher.join(timeout=3) + conn.close() + lock.close() + return 0 + + +if __name__ == '__main__': + sys.exit(main()) diff --git a/app/documents.py b/app/documents.py new file mode 100644 index 0000000..f88c247 --- /dev/null +++ b/app/documents.py @@ -0,0 +1,451 @@ +"""Permission-filtered document catalog. Originals are only ever opened for reading.""" + +import contextlib +import hashlib +import json +import os +import re +import sqlite3 +import stat +import time +import uuid + +try: + from . import access_control, reconcile_shares as directory + from .account_policy import account_name, is_excluded_user + from .state_db import connect_state_db +except ImportError: + import access_control + import reconcile_shares as directory + from account_policy import account_name, is_excluded_user + from state_db import connect_state_db + +SEARCH_ROOT = os.getenv('DOCUMENT_STATE_ROOT', os.path.join(os.getenv('STATE_ROOT', '/state'), 'documents')) +SEARCH_DB = os.path.join(SEARCH_ROOT, 'search.db') +MAX_TEXT = 2 * 1024 * 1024 +TEXT_SUFFIXES = {'.txt', '.md', '.csv', '.tsv', '.json', '.xml', '.html', '.htm', '.log', '.ini', '.yaml', '.yml'} +OFFICE_SUFFIXES = {'.docx', '.xlsx', '.pptx', '.odt', '.ods', '.odp'} +IMAGE_SUFFIXES = {'.jpg', '.jpeg', '.png', '.tif', '.tiff', '.webp', '.bmp'} + + +def connect(): + os.makedirs(SEARCH_ROOT, mode=0o700, exist_ok=True) + os.chmod(SEARCH_ROOT, 0o700) + return connect_state_db(SEARCH_DB) + + +def ensure_schema(conn): + conn.executescript(''' + CREATE TABLE IF NOT EXISTS document_sources ( + id TEXT PRIMARY KEY, kind TEXT NOT NULL, label TEXT NOT NULL, + root TEXT NOT NULL, owner_sid TEXT NOT NULL DEFAULT '' + ); + CREATE TABLE IF NOT EXISTS documents ( + rowid INTEGER PRIMARY KEY, id TEXT NOT NULL UNIQUE, + source_id TEXT NOT NULL REFERENCES document_sources(id) ON DELETE CASCADE, + path TEXT NOT NULL, name TEXT NOT NULL, extension TEXT NOT NULL, + size INTEGER NOT NULL, modified REAL NOT NULL, fingerprint TEXT NOT NULL, + body TEXT NOT NULL DEFAULT '', state TEXT NOT NULL DEFAULT 'pending', + preview TEXT NOT NULL DEFAULT '', pages INTEGER NOT NULL DEFAULT 0, + attempts INTEGER NOT NULL DEFAULT 0, retry_at REAL NOT NULL DEFAULT 0, + indexed_at REAL NOT NULL DEFAULT 0, seen INTEGER NOT NULL DEFAULT 0, + UNIQUE(source_id,path) + ); + CREATE INDEX IF NOT EXISTS document_queue ON documents(state,retry_at); + CREATE INDEX IF NOT EXISTS document_source ON documents(source_id,modified); + CREATE TABLE IF NOT EXISTS document_worker ( + id INTEGER PRIMARY KEY CHECK(id=1), paused INTEGER NOT NULL DEFAULT 0, + heartbeat REAL NOT NULL DEFAULT 0, state TEXT NOT NULL DEFAULT 'idle', + current_id TEXT NOT NULL DEFAULT '', processed INTEGER NOT NULL DEFAULT 0, + scanning INTEGER NOT NULL DEFAULT 0, last_scan REAL NOT NULL DEFAULT 0 + ); + INSERT OR IGNORE INTO document_worker(id) VALUES(1); + CREATE TABLE IF NOT EXISTS document_activity ( + id INTEGER PRIMARY KEY, timestamp REAL NOT NULL, action TEXT NOT NULL, + path TEXT NOT NULL DEFAULT '', source TEXT NOT NULL DEFAULT '', + actor TEXT NOT NULL DEFAULT '', message TEXT NOT NULL DEFAULT '' + ); + CREATE VIRTUAL TABLE IF NOT EXISTS document_fts USING fts5( + name,path,body,content='documents',content_rowid='rowid', + tokenize='unicode61 remove_diacritics 2',prefix='2 3 4' + ); + CREATE TRIGGER IF NOT EXISTS document_insert AFTER INSERT ON documents BEGIN + INSERT INTO document_fts(rowid,name,path,body) VALUES(new.rowid,new.name,new.path,new.body); + END; + CREATE TRIGGER IF NOT EXISTS document_delete AFTER DELETE ON documents BEGIN + INSERT INTO document_fts(document_fts,rowid,name,path,body) VALUES('delete',old.rowid,old.name,old.path,old.body); + END; + CREATE TRIGGER IF NOT EXISTS document_update AFTER UPDATE OF name,path,body ON documents BEGIN + INSERT INTO document_fts(document_fts,rowid,name,path,body) VALUES('delete',old.rowid,old.name,old.path,old.body); + INSERT INTO document_fts(rowid,name,path,body) VALUES(new.rowid,new.name,new.path,new.body); + END; + ''') + + +def worker_paused(conn): + return bool(conn.execute('SELECT paused FROM document_worker WHERE id=1').fetchone()[0]) + + +def worker_event(conn, action, path='', source='', actor='', message=''): + conn.execute('INSERT INTO document_activity(timestamp,action,path,source,actor,message) VALUES(?,?,?,?,?,?)', + (time.time(), action, path, source, actor, str(message)[:200])) + + +def control_worker(conn, action, actor): + if action not in {'pause', 'resume'}: + raise ValueError('Ungültige Aktion') + conn.execute('BEGIN IMMEDIATE') + try: + paused = action == 'pause' + if worker_paused(conn) != paused: + conn.execute('UPDATE document_worker SET paused=? WHERE id=1', (int(paused),)) + worker_event(conn, action, actor=actor) + conn.commit() + except Exception: + conn.rollback() + raise + return worker_snapshot(conn) + + +def worker_snapshot(conn): + result = dict(conn.execute('SELECT * FROM document_worker WHERE id=1').fetchone()) + result['paused'] = bool(result['paused']) + result['online'] = result['heartbeat'] > time.time() - 15 + counts = dict(conn.execute('SELECT state,count(*) FROM documents GROUP BY state').fetchall()) + total = sum(counts.values()) + ready = counts.get('ready', 0) + counts.get('name-only', 0) + result['counts'] = {'total': total, 'complete': ready, 'pending': counts.get('pending', 0), + 'ocr': counts.get('ocr', 0), 'nameOnly': counts.get('name-only', 0), + 'errors': conn.execute("SELECT count(*) FROM documents WHERE state='failed' OR attempts>0").fetchone()[0]} + result['progress'] = round(100 * ready / total) if total else 0 + row = conn.execute('''SELECT d.name,d.path,d.pages,CASE WHEN s.kind='private' THEN s.label || ' / ' || substr(s.id,9) ELSE s.label END AS source,s.kind FROM documents d + JOIN document_sources s ON s.id=d.source_id WHERE d.id=?''', (result['current_id'],)).fetchone() + result['current'] = dict(row) if row else None + result['activity'] = [dict(row) for row in conn.execute('SELECT * FROM document_activity ORDER BY id DESC LIMIT 100')] + result['failures'] = [dict(row) for row in conn.execute('''SELECT d.name,d.path,CASE WHEN s.kind='private' THEN s.label || ' / ' || substr(s.id,9) ELSE s.label END AS source,d.state,d.attempts,d.retry_at + FROM documents d JOIN document_sources s ON s.id=d.source_id + WHERE d.state='failed' OR d.attempts>0 ORDER BY d.retry_at LIMIT 50''')] + return result + + +def fingerprint(info): + return ':'.join(str(value) for value in (info.st_dev, info.st_ino, info.st_size, info.st_mtime_ns, info.st_ctime_ns)) + + +def valid_parts(path): + parts = path.split('/') + if not path or any(part in {'', '.', '..', '.trash'} or '\x00' in part for part in parts): + raise FileNotFoundError('Invalid document path') + return parts + + +@contextlib.contextmanager +def open_file(root, path): + """Pin every path component; reject links, special files and nested mounts.""" + parts = valid_parts(path) + fd = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + handle = None + try: + device = os.fstat(fd).st_dev + for part in parts[:-1]: + child = os.open(part, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=fd) + os.close(fd) + fd = child + if os.fstat(fd).st_dev != device: + raise FileNotFoundError('Mount boundary') + child = os.open(parts[-1], os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK, dir_fd=fd) + handle = os.fdopen(child, 'rb') + info = os.fstat(handle.fileno()) + if not stat.S_ISREG(info.st_mode) or info.st_dev != device: + raise FileNotFoundError('Not a regular document') + yield handle, info + finally: + if handle is not None: + handle.close() + os.close(fd) + + +def source_inventory(): + conn = directory.open_db() + try: + # Reconciliation and login populate these identities. Never import groups here. + sources = [] + for row in conn.execute('SELECT * FROM shares WHERE isActive=1'): + root = access_control.safe_folder_path(row) + if os.path.isdir(root): + sources.append({'id': 'data:' + row['objectGUID'], 'kind': 'data', 'label': row['shareName'], 'root': root, 'owner_sid': ''}) + users = {row['sam'].casefold(): row['sid'] for row in conn.execute('SELECT sid,sam FROM access_users') + if not is_excluded_user(row['sam'])} + if os.path.isdir(directory.PRIVATE_ROOT): + for entry in os.scandir(directory.PRIVATE_ROOT): + sid = users.get(entry.name.casefold()) + if sid and entry.is_dir(follow_symlinks=False): + sources.append({'id': 'private:' + entry.name.casefold(), 'kind': 'private', + 'label': 'Private', 'root': entry.path, 'owner_sid': sid}) + return sources + finally: + conn.close() + + +def allowed_sources(conn, identity): + """Use current policy on every request, including previews and downloads.""" + if is_excluded_user(str(identity.get('sub', ''))): + return {} + policy = directory.open_db() + try: + excluded = access_control.excluded_user_sids(policy) + sid = identity.get('sid') + if not sid or sid in excluded: + return {} + folders = {row['objectGUID']: row for row in policy.execute('SELECT * FROM shares WHERE isActive=1')} + levels = {row['folderId']: row['level'] for row in policy.execute( + "SELECT folderId,level FROM folder_permissions WHERE kind='user' AND principalId=?", (sid,))} + result = {} + for source in conn.execute('SELECT * FROM document_sources'): + if source['kind'] == 'data': + folder_id = source['id'][5:] + folder = folders.get(folder_id) + if folder is None or (identity.get('role') != 'domain-admin' and levels.get(folder_id, 0) < 1): + continue + root = access_control.safe_folder_path(folder) + if source['root'] != root: + continue + elif source['kind'] == 'private': + # The user portal always shows only one's own home, even for admins. + if source['owner_sid'] != sid or os.path.dirname(source['root']) != os.path.abspath(directory.PRIVATE_ROOT): + continue + try: + if os.stat(source['root'], follow_symlinks=False).st_uid != identity.get('uid'): + continue + except OSError: + continue + else: + continue + if not os.path.islink(source['root']): + result[source['id']] = dict(source) + return result + finally: + policy.close() + + +def discover_file(conn, source, relative, seen=0): + try: + with open_file(source['root'], relative) as (_, info): + stamp = fingerprint(info) + except (OSError, ValueError): + conn.execute('DELETE FROM documents WHERE source_id=? AND path=?', (source['id'], relative)) + return + old = conn.execute('SELECT fingerprint FROM documents WHERE source_id=? AND path=?', (source['id'], relative)).fetchone() + if old and old['fingerprint'] == stamp: + if seen: + conn.execute('UPDATE documents SET seen=? WHERE source_id=? AND path=?', (seen, source['id'], relative)) + return + name = relative.rsplit('/', 1)[-1] + conn.execute('''INSERT INTO documents(id,source_id,path,name,extension,size,modified,fingerprint,seen) + VALUES(?,?,?,?,?,?,?,?,?) ON CONFLICT(source_id,path) DO UPDATE SET + name=excluded.name,extension=excluded.extension,size=excluded.size, + modified=excluded.modified,fingerprint=excluded.fingerprint,seen=excluded.seen, + body='',state='pending',preview='',pages=0,attempts=0,retry_at=0,indexed_at=0''', + (uuid.uuid4().hex, source['id'], relative, name, os.path.splitext(name)[1].lower(), + info.st_size, info.st_mtime, stamp, seen)) + + +def walk_source(conn, source, watch=None, should_stop=None): + """Deletion affects index records only. Original files are never removed.""" + seen = time.time_ns() + completed = True + visited = 0 + def scan(relative=''): + nonlocal completed, visited + if should_stop and should_stop(): + completed = False + return + path = os.path.join(source['root'], relative) + try: + with open_directory(source['root'], relative) as fd: + if watch: + watch(path, source['id'], relative) + device = os.fstat(fd).st_dev + for entry in os.scandir(fd): + visited += 1 + if visited % 100 == 0: + conn.commit() + if should_stop and should_stop(): + completed = False + return + if entry.name == '.trash' or entry.is_symlink(): + continue + child = '/'.join(filter(None, (relative, entry.name))) + info = entry.stat(follow_symlinks=False) + if info.st_dev != device: + continue + if stat.S_ISDIR(info.st_mode): + scan(child) + elif stat.S_ISREG(info.st_mode): + discover_file(conn, source, child, seen) + conn.commit() + except OSError: + completed = False + scan() + if completed: + conn.execute('DELETE FROM documents WHERE source_id=? AND seen<>?', (source['id'], seen)) + conn.commit() + return completed + + +@contextlib.contextmanager +def open_directory(root, relative): + fd = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + device = os.fstat(fd).st_dev + for part in valid_parts(relative) if relative else []: + child = os.open(part, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=fd) + os.close(fd) + fd = child + if os.fstat(fd).st_dev != device: + raise FileNotFoundError('Mount boundary') + yield fd + finally: + os.close(fd) + + +def private_readable(source, relative, uid): + if source['kind'] != 'private': + return True + try: + with open_directory(source['root'], '') as fd: + info = os.fstat(fd) + if info.st_uid != uid or info.st_mode & 0o500 != 0o500: + return False + parts = valid_parts(relative) + for index in range(1, len(parts)): + with open_directory(source['root'], '/'.join(parts[:index])) as fd: + info = os.fstat(fd) + if info.st_uid != uid or info.st_mode & 0o500 != 0o500: + return False + with open_file(source['root'], relative) as (_, info): + return info.st_uid == uid and bool(info.st_mode & 0o400) + except OSError: + return False + + +def current_document(conn, document_id, identity): + sources = allowed_sources(conn, identity) + row = conn.execute('SELECT * FROM documents WHERE id=?', (document_id,)).fetchone() + if row is None or row['source_id'] not in sources: + raise FileNotFoundError('Document not found') + source = sources[row['source_id']] + if not private_readable(source, row['path'], identity.get('uid')): + raise FileNotFoundError('Document not readable') + try: + with open_file(source['root'], row['path']) as (_, info): + if fingerprint(info) != row['fingerprint']: + raise FileNotFoundError('Document changed') + except OSError as exc: + raise FileNotFoundError('Document unavailable') from exc + return dict(row), source + + +def public_document(row, source, snippet=''): + version = hashlib.sha256(row['fingerprint'].encode()).hexdigest()[:16] + return {'id': row['id'], 'name': row['name'], 'path': row['path'], 'sourceId': row['source_id'], + 'source': source['label'], 'kind': source['kind'], 'extension': row['extension'].lstrip('.'), + 'size': row['size'], 'modified': row['modified'], 'state': row['state'], 'pages': row['pages'], + 'hasPreview': bool(row['preview']), 'version': version, 'snippet': snippet, 'attempts': row['attempts']} + + +def search(conn, identity, params): + sources = allowed_sources(conn, identity) + query = str(params.get('q', [''])[0]).strip()[:256] + scope = params.get('scope', ['all'])[0] + source_filter = params.get('source', [''])[0] + extension = params.get('type', [''])[0].lower().lstrip('.') + sort = params.get('sort', ['recent'])[0] + try: + limit = min(100, max(1, int(params.get('limit', ['40'])[0]))) + offset = max(0, int(params.get('offset', ['0'])[0])) + except ValueError as exc: + raise ValueError('Ungültige Seitennummer') from exc + selected = {key: value for key, value in sources.items() if not source_filter or key == source_filter} + base = {'items': [], 'total': 0, 'hasMore': False, 'sources': [{'id': key, 'label': value['label'], 'kind': value['kind']} for key, value in sources.items()], 'types': [], 'pending': 0} + if not sources: + return base + # Filter before snippets, pagination, counts and facets. Hidden files + # must never influence any user-visible result. + ids = list(selected) + if not ids: + return base + placeholders = ','.join('?' for _ in ids) + conn.execute('CREATE TEMP TABLE IF NOT EXISTS readable_private(id TEXT PRIMARY KEY)') + conn.execute('DELETE FROM readable_private') + for row in conn.execute(f"SELECT id,source_id,path FROM documents WHERE source_id IN ({placeholders}) AND source_id LIKE 'private:%'", ids).fetchall(): + if private_readable(selected[row['source_id']], row['path'], identity.get('uid')): + conn.execute('INSERT INTO readable_private VALUES(?)', (row['id'],)) + private_condition = "(d.source_id NOT LIKE 'private:%' OR EXISTS (SELECT 1 FROM readable_private r WHERE r.id=d.id))" + conditions = [f'd.source_id IN ({placeholders})', private_condition] + values = ids[:] + from_sql = 'documents d' + snippet_sql = "substr(d.body,1,220)" + if query: + if scope == 'name': + conditions.append("d.name LIKE ? ESCAPE '\\'") + values.append('%' + query.replace('\\', '\\\\').replace('%', '\\%').replace('_', '\\_') + '%') + else: + words = re.findall(r'\w+', query, flags=re.UNICODE)[:24] + if not words: + return base + match = ' AND '.join('"' + word + '"*' for word in words) + if scope == 'content': + match = 'body : (' + match + ')' + from_sql += ' JOIN document_fts ON document_fts.rowid=d.rowid' + conditions.append('document_fts MATCH ?') + values.append(match) + snippet_sql = "snippet(document_fts,2,char(1),char(2),' … ',32)" + if extension: + conditions.append('d.extension=?') + values.append('.' + extension) + where = ' AND '.join(conditions) + order = {'name': 'd.name COLLATE NOCASE', 'oldest': 'd.modified', 'recent': 'd.modified DESC'}.get(sort, 'd.modified DESC') + if query and scope != 'name' and sort == 'relevance': + order = 'bm25(document_fts,6,2,1)' + base['total'] = conn.execute(f'SELECT count(*) FROM {from_sql} WHERE {where}', values).fetchone()[0] + base['hasMore'] = offset + limit < base['total'] + columns = 'd.id,d.source_id,d.path,d.name,d.extension,d.size,d.modified,d.fingerprint,d.state,d.preview,d.pages,d.attempts' + rows = conn.execute(f'SELECT {columns}, {snippet_sql} AS excerpt FROM {from_sql} WHERE {where} ORDER BY {order},d.id LIMIT ? OFFSET ?', [*values,limit,offset]) + visible = [] + for row in rows: + try: + with open_file(selected[row['source_id']]['root'], row['path']) as (_, info): + if fingerprint(info) == row['fingerprint'] and private_readable(selected[row['source_id']], row['path'], identity.get('uid')): + visible.append(row) + except OSError: + continue + base['items'] = [public_document(row, selected[row['source_id']], row['excerpt']) for row in visible] + base['types'] = sorted({row['extension'].lstrip('.') for row in conn.execute(f'SELECT DISTINCT extension FROM documents d WHERE source_id IN ({placeholders}) AND {private_condition}', ids) if row['extension']}) + base['pending'] = conn.execute(f"SELECT count(*) FROM documents d WHERE source_id IN ({placeholders}) AND {private_condition} AND state IN ('pending','ocr')", ids).fetchone()[0] + return base + + +def detail(conn, identity, document_id): + row, source = current_document(conn, document_id, identity) + result = public_document(row, source) + result['text'] = row['body'] + return result + + +@contextlib.contextmanager +def download(conn, identity, document_id): + row, source = current_document(conn, document_id, identity) + with open_file(source['root'], row['path']) as (handle, info): + if fingerprint(info) != row['fingerprint']: + raise FileNotFoundError('Document changed') + yield handle, row + + +@contextlib.contextmanager +def preview(conn, identity, document_id): + row, _ = current_document(conn, document_id, identity) + if not row['preview'] or not re.fullmatch(r'[a-f0-9]{32}\.jpg', row['preview']): + raise FileNotFoundError('Preview unavailable') + with open_file(os.path.join(SEARCH_ROOT, 'previews'), row['preview']) as (handle, info): + yield handle, info.st_size diff --git a/app/extract_document.py b/app/extract_document.py new file mode 100644 index 0000000..a6eee78 --- /dev/null +++ b/app/extract_document.py @@ -0,0 +1,169 @@ +#!/usr/bin/env python3 +"""Run document parsers in a restricted, disposable workspace, never on originals.""" + +import ctypes +import errno +import json +import os +from pathlib import Path +import platform +import re +import resource +import subprocess +import sys +import warnings +import zipfile +from xml.etree import ElementTree + +MAX_TEXT = 2 * 1024 * 1024 + + +def sandbox(workspace): + """Landlock filesystem rules and seccomp prevent parser access to live data/network.""" + libc = ctypes.CDLL(None, use_errno=True) + abi = libc.syscall(444, 0, 0, 1) + if abi < 1: + raise RuntimeError('Document isolation requires Linux Landlock') + class Ruleset(ctypes.Structure): + _fields_ = [('fs', ctypes.c_uint64)] + class PathRule(ctypes.Structure): + _pack_ = 1 + _fields_ = [('access', ctypes.c_uint64), ('fd', ctypes.c_int32)] + handled = (1 << 13) - 1 + if abi >= 2: + handled |= 1 << 13 + if abi >= 3: + handled |= 1 << 14 + rules = Ruleset(handled) + rules_fd = libc.syscall(444, ctypes.byref(rules), ctypes.sizeof(rules), 0) + if rules_fd < 0: + raise OSError(ctypes.get_errno(), 'Cannot create document isolation') + read = (1 << 0) | (1 << 2) | (1 << 3) + try: + for path, access in [(workspace, handled), ('/usr', read), ('/lib', read), ('/lib64', read), ('/bin', read), + ('/etc/fonts', read), ('/etc/ghostscript', read), ('/etc/ld.so.cache', 1 << 2), + ('/etc/papersize', 1 << 2), ('/dev/null', (1 << 1) | (1 << 2)), + ('/dev/urandom', 1 << 2), ('/proc/cpuinfo', 1 << 2), ('/proc/meminfo', 1 << 2)]: + if not os.path.exists(path): + continue + fd = os.open(path, os.O_PATH | os.O_CLOEXEC) + try: + rule = PathRule(access, fd) + if libc.syscall(445, rules_fd, 1, ctypes.byref(rule), 0) != 0: + raise OSError(ctypes.get_errno(), 'Cannot restrict document path') + finally: + os.close(fd) + if libc.prctl(38, 1, 0, 0, 0) != 0 or libc.syscall(446, rules_fd, 0) != 0: + raise OSError(ctypes.get_errno(), 'Cannot enter document isolation') + finally: + os.close(rules_fd) + # Permit local IPC used by OCR workers, deny creation of internet sockets. + socket_number = {'aarch64': 198, 'x86_64': 41}.get(platform.machine()) + if socket_number is None: + raise RuntimeError('Unsupported document sandbox architecture') + class Filter(ctypes.Structure): + _fields_ = [('code', ctypes.c_ushort), ('jt', ctypes.c_ubyte), ('jf', ctypes.c_ubyte), ('k', ctypes.c_uint32)] + class Program(ctypes.Structure): + _fields_ = [('length', ctypes.c_ushort), ('filter', ctypes.POINTER(Filter))] + filters = (Filter * 6)(Filter(0x20,0,0,0), Filter(0x15,0,3,socket_number), + Filter(0x20,0,0,16), Filter(0x15,1,0,1), + Filter(0x06,0,0,0x50000 | errno.EPERM), Filter(0x06,0,0,0x7fff0000)) + program = Program(len(filters), filters) + if libc.prctl(22, 2, ctypes.byref(program), 0, 0) != 0: + raise OSError(ctypes.get_errno(), 'Cannot isolate document networking') + + +def command(args, timeout=120): + result = subprocess.run(args, stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=timeout, check=False) + if result.returncode: + raise RuntimeError(f'{args[0]} failed ({result.returncode})') + return result.stdout.decode('utf-8', errors='replace') + + +def read_text(path): + with open(path, 'rb') as handle: + raw = handle.read(MAX_TEXT) + if raw.startswith((b'\xff\xfe', b'\xfe\xff')): + return raw.decode('utf-16', errors='replace') + return raw.decode('utf-8-sig', errors='replace') + + +def office_text(path): + text = [] + total = 0 + with zipfile.ZipFile(path) as archive: + entries = archive.infolist() + if len(entries) > 10000 or sum(entry.file_size for entry in entries) > 64 * 1024 * 1024: + raise RuntimeError('Office document exceeds extraction limits') + for entry in sorted(entries, key=lambda item: item.filename): + name = entry.filename + wanted = (name.startswith(('word/', 'ppt/slides/', 'xl/worksheets/')) or name in {'xl/sharedStrings.xml', 'content.xml'}) and name.endswith('.xml') + if not wanted or entry.file_size > 16 * 1024 * 1024: + continue + node = ElementTree.fromstring(archive.read(entry)) + fragment = ' '.join(node.itertext()) + text.append(fragment) + total += len(fragment) + if total > MAX_TEXT: + break + return '\n'.join(text)[:MAX_TEXT] + + +def extract(config): + suffix = config['extension'] + source = Path('input') + body, pages, needs_ocr = '', 0, False + if suffix == '.pdf': + info = command(['pdfinfo', str(source)]) + found = re.search(r'^Pages:\s*(\d+)', info, flags=re.MULTILINE) + pages = int(found.group(1)) if found else 0 + if pages > config['max_pages']: + raise RuntimeError('PDF exceeds page limit') + if config['phase'] == 'ocr': + command(['ocrmypdf', '--redo-ocr', '--output-type', 'pdf', '--optimize', '0', '--jobs', '1', + '--language', config['language'], '--tesseract-timeout', '120', '--skip-big', '50', + str(source), 'ocr.pdf'], timeout=config['timeout'] - 15) + source = Path('ocr.pdf') + command(['pdftotext', '-enc', 'UTF-8', '-layout', str(source), 'text.txt']) + body = read_text('text.txt') + if config['phase'] != 'ocr': + image_list = command(['pdfimages', '-list', str(source)]) + has_images = bool(re.search(r'^\s*\d+\s+\d+\s+image\s', image_list, re.MULTILINE)) + page_text = body.split('\f')[:pages] + needs_ocr = has_images and any(len(re.sub(r'\s+', '', page)) < 80 for page in page_text) + command(['pdftoppm', '-f', '1', '-l', '1', '-singlefile', '-scale-to', '1400', '-jpeg', str(source), 'preview']) + elif suffix in {'.docx', '.xlsx', '.pptx', '.odt', '.ods', '.odp'}: + body = office_text(source) + elif suffix in {'.jpg', '.jpeg', '.png', '.tif', '.tiff', '.webp', '.bmp'}: + from PIL import Image + Image.MAX_IMAGE_PIXELS = 25_000_000 + warnings.simplefilter('error', Image.DecompressionBombWarning) + with Image.open(source) as image: + image.thumbnail((1400,1400)) + image.convert('RGB').save('preview.jpg', 'JPEG', quality=85) + else: + body = read_text(source) + if suffix in {'.html', '.htm', '.xml'}: + body = re.sub(r'<[^>]+>', ' ', body) + return {'body': body[:MAX_TEXT], 'pages': pages, 'needsOcr': needs_ocr} + + +def main(): + workspace = os.path.abspath(sys.argv[1]) + os.chdir(workspace) + config = json.loads(Path('config.json').read_text()) + resource.setrlimit(resource.RLIMIT_CORE, (0,0)) + resource.setrlimit(resource.RLIMIT_AS, (2 * 1024**3, 2 * 1024**3)) + resource.setrlimit(resource.RLIMIT_FSIZE, (512 * 1024**2, 512 * 1024**2)) + resource.setrlimit(resource.RLIMIT_CPU, (config['timeout'], config['timeout'])) + try: + sandbox(workspace) + result = extract(config) + except Exception as exc: + result = {'error': str(exc)[:200]} + Path('result.json').write_text(json.dumps(result)) + return 0 if 'error' not in result else 1 + + +if __name__ == '__main__': + sys.exit(main()) diff --git a/app/init.sh b/app/init.sh index 0565271..8025729 100755 --- a/app/init.sh +++ b/app/init.sh @@ -559,7 +559,15 @@ start_observability_services() { ;; esac fi - log "Starting administration web UI for https://${WEB_HOSTNAME}" + log 'Starting document search and OCR worker' + ( + while true; do + python3 /app/document_index.py || log 'Document worker failed' + log 'Document worker stopped; restarting in two seconds' + sleep 2 + done + ) & + log "Starting web UI for https://${WEB_HOSTNAME}" python3 /app/web_ui.py & else log 'WEB_ENABLED is false; web UI disabled' diff --git a/app/web/app.js b/app/web/app.js index d5a9bc2..7cf535f 100644 --- a/app/web/app.js +++ b/app/web/app.js @@ -11,7 +11,7 @@ const esc = value => String(value ?? "").replace(/[&<>'"]/g, char => ({"&":"& function actionIcon(action) { const paths = { download: '', - restore: '', + restore: '', }; return ``; } @@ -60,6 +60,7 @@ const isoDay = date => date.toISOString().slice(0, 10); const sum = (rows, field) => (rows || []).reduce((total, row) => total + Number(row[field] || 0), 0); async function api(path, options = {}) { + if (path.startsWith("/api/") && !["/api/login", "/api/logout", "/api/session"].includes(path)) path = "/admin" + path; const response = await fetch(path, { ...options, headers: {"Content-Type": "application/json", ...(options.headers || {})}, @@ -91,11 +92,12 @@ function showLogin() { } function showApp(session) { + if (session.role !== "domain-admin") { location.replace("/"); return; } state.session = {user: session.user, expiresAt: session.expiresAt}; document.querySelector("#session-user").textContent = session.user; loginView.hidden = true; appView.hidden = false; - navigate(location.pathname === "/" ? "/overview" : location.pathname, true); + navigate(["/", "/admin", "/admin/"].includes(location.pathname) ? "/admin/overview" : location.pathname, true); } function setLoading() { content.innerHTML = '
Wird geladen…
'; } @@ -106,10 +108,12 @@ function empty(message) { return `
${esc(message)}
`; } function badge(text, kind = "") { return `${esc(text)}`; } function routeFor(path) { + path = path.replace(/^\/admin(?=\/|$)/, "") || "/overview"; if (path.startsWith("/storage/data")) return "storage-data"; if (path.startsWith("/storage/users")) return "storage-users"; if (path.startsWith("/access")) return "access"; if (path.startsWith("/reconciliation")) return "reconciliation"; + if (path.startsWith("/documents")) return "documents"; if (path.startsWith("/activity/fslogix")) return "activity-fslogix"; if (path.startsWith("/activity")) return "activity"; if (path.startsWith("/trash")) return "trash"; @@ -124,7 +128,8 @@ async function navigate(path, replace = false) { if (replace && state.currentPath) history.replaceState({}, "", state.currentPath); return; } - if (path === "/shares" || path.startsWith("/shares/")) { path = "/access"; replace = true; } + if (path === "/shares" || path.startsWith("/shares/") || path === "/admin/shares") { path = "/admin/access"; replace = true; } + if (!path.startsWith("/admin/")) path = "/admin" + path; state.accessDirty = false; state.currentPath = path; clearInterval(state.timer); @@ -138,6 +143,7 @@ async function navigate(path, replace = false) { if (route === "overview") await renderOverview(); if (route === "access") await renderAccess(); if (route === "reconciliation") await renderReconciliation(); + if (route === "documents") await renderDocuments(); if (route === "storage-data") await renderStorage("data"); if (route === "storage-users") await renderStorage("users"); if (route === "activity") await renderActivity(); @@ -489,7 +495,7 @@ async function renderTrash() { ${esc(item.path)} ${bytes(item.size)} ${esc(utcTime(item.expiresAt))} -
${actionIcon("download")}
+
${actionIcon("download")}
`).join(""); } const shown = result.items.length.toLocaleString("de-DE"); @@ -731,6 +737,55 @@ async function renderReconciliation() { state.timer = window.setInterval(() => { void poll(); }, 1500); } +async function renderDocuments() { + content.innerHTML = pageHead("Dokumentindex", "Dateierfassung, Volltextsuche und PDF-Texterkennung.") + ` +

Status wird geladen…

+

+

+
Letzte Dateierfassung—
Letztes Lebenszeichen—
+
+
+

Ausstehende Wiederholungen

Ordner / DateiAuftragVersucheNächster Versuch (UTC)
+

Aktivität

Zeit (UTC)VorgangOrdner / DateiBenutzer / Ergebnis
`; + const button = document.querySelector("#document-worker-control"); + let status, polling = false, controlling = false; + const actionLabels = {started:"Worker gestartet",text:"Text erfassen",ocr:"PDF-Texterkennung",complete:"Abgeschlossen","queued-ocr":"Für OCR vorgemerkt",changed:"Datei verändert",failed:"Wiederholung geplant",interrupted:"Auftrag unterbrochen",pause:"Anhalten",resume:"Fortsetzen"}; + const refresh = async () => { + if (polling || !button.isConnected) return; + polling = true; + try { + const data = await api("/api/documents/status"); + if (!button.isConnected) return; + status = data; + const paused = data.paused; + const caption = !data.online ? "Worker nicht erreichbar" : paused ? (data.current ? "Wird angehalten…" : "Angehalten") : data.state === "ocr" ? "PDF-Texterkennung läuft" : data.state === "text" ? "Text wird erfasst" : data.scanning ? "Dateikatalog wird aktualisiert" : "Bereit"; + document.querySelector("#document-worker-state").textContent = caption; + button.textContent = paused ? "Fortsetzen" : "Anhalten"; + button.disabled = controlling; + document.querySelector("#document-worker-current").textContent = data.current ? `${data.current.source} / ${data.current.path}` : "Kein aktiver Auftrag"; + document.querySelector("#document-worker-progress").value = data.progress; + document.querySelector("#document-worker-progress-label").textContent = `${data.counts.complete.toLocaleString("de-DE")} von ${data.counts.total.toLocaleString("de-DE")} Dateien erfasst · ${data.progress} %`; + document.querySelector("#document-worker-scan").textContent = data.last_scan ? utcTime(new Date(data.last_scan * 1000)) : "Ausstehend"; + document.querySelector("#document-worker-heartbeat").textContent = data.heartbeat ? utcTime(new Date(data.heartbeat * 1000)) : "Ausstehend"; + document.querySelector("#document-worker-counts").innerHTML = [["Text-Warteschlange",data.counts.pending],["OCR-Warteschlange",data.counts.ocr],["Wiederholungen",data.counts.errors],["Nur Dateiname",data.counts.nameOnly]].map(([label,value]) => `
${label}${value.toLocaleString("de-DE")}
`).join(""); + document.querySelector("#document-worker-failures").innerHTML = data.failures.length ? data.failures.map(row => `${esc(row.source)} / ${esc(row.path)}${row.state === "ocr" ? "OCR" : "Text"}${Number(row.attempts)}${esc(utcTime(new Date(row.retry_at * 1000)))}`).join("") : 'Keine fehlgeschlagenen Aufträge'; + document.querySelector("#document-worker-activity").innerHTML = data.activity.length ? data.activity.map(row => `${esc(utcTime(new Date(row.timestamp * 1000)))}${esc(actionLabels[row.action] || row.action)}${esc(row.source)}${row.path ? " / " + esc(row.path) : "—"}${esc(row.actor || row.message || "—")}`).join("") : 'Noch keine Aktivität'; + } catch (error) { notice(error.message); } + finally { polling = false; } + }; + button.addEventListener("click", async () => { + if (!status || controlling) return; + controlling = true; button.disabled = true; + try { + await api("/api/documents/control", {method:"POST",body:JSON.stringify({action:status.paused ? "resume" : "pause"})}); + notice(status.paused ? "Dokumentindex wird fortgesetzt." : "Dokumentindex wird angehalten."); + } catch (error) { notice(error.message); } + finally { controlling = false; await refresh(); } + }); + await refresh(); + if (button.isConnected) state.timer = setInterval(() => { void refresh(); }, 1500); +} + async function renderReport() { content.innerHTML = pageHead("PDF-Bericht", "Vollständiger, druckbarer Stand der Dateiserver-Konfiguration und Belegung.") + `
diff --git a/app/web/index.html b/app/web/index.html index aa7647d..3cdf6ab 100644 --- a/app/web/index.html +++ b/app/web/index.html @@ -11,7 +11,7 @@