From 0ea17c6401488f7ff680d17fb781c623605a0742 Mon Sep 17 00:00:00 2001 From: Ludwig Lehnert Date: Wed, 12 Aug 2026 10:20:32 +0000 Subject: [PATCH] fslogix separation --- README.md | 7 +- app/audit_collector.py | 5 +- app/audit_store.py | 365 ++++++++++++++++++---------- app/backup_to_destination.py | 3 + app/web/app.js | 19 +- app/web/index.html | 1 + app/web/report.mjs | 1 + app/web_ui.py | 87 ++++++- tests/test_backup_to_destination.py | 16 ++ tests/test_web_ui.py | 282 ++++++++++++++++++++- 10 files changed, 633 insertions(+), 153 deletions(-) diff --git a/README.md b/README.md index d9508b1..d1e45e7 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ This repository provides a production-oriented Samba file server container that - Setup prompts for well-known authorization groups by SID (`DOMAIN_USERS_SID`, `DOMAIN_ADMINS_SID`) to avoid localized group names. - `FSLOGIX_GROUP_SID` controls who can access the default FSLogix share (defaults to `DOMAIN_USERS_SID`). - Startup resolves those SIDs to NSS group names via winbind, then uses those resolved groups in Samba `valid users` rules. -- Samba `full_audit` is restricted to successful and failed reads, writes, renames, and deletions. +- 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. - A plain HTTPS administration console provides read-only statistics and logs plus narrowly scoped actions for 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. @@ -393,15 +393,17 @@ If `WEB_ENABLED` is absent on an upgraded installation and no TLS settings/certi ## Activity Database -Samba emits only the selected high-level `full_audit` operations, and the collector stores them durably: +Samba emits selected high-level `full_audit` operations for Private, Data, and FSLogix, and the collector stores them durably: - it tails every `/var/log/samba/log.*` source and stores inode/byte offsets transactionally in `audit_sources`, so source progress and inserted events commit together; - each event records a UTC timestamp, user, client address/name, share, result, path, and one of `read`, `write`, `move`, or `delete`; - directory listings, sessions, metadata access, file-open/create noise, and all other VFS operations are discarded; users ending in a configured `AUDIT_SKIP_USER_SUFFIXES` value are also discarded; +- FSLogix profile-container events are retained and queried through their own partial indexes, API route, and UI log instead of appearing in the main activity stream; - immediately consecutive identical reads within the same UTC second are collapsed into one event, including across collector polling cycles; a different event interrupts the sequence and preserves later reads; - activity pages read only `limit + 1` indexed rows and use a stable time/id cursor; - exact counts for date, user, share, action, and result filters and all facet lists come from daily rollups; - selective substring path searches use the deduplicated trigram index and skip an expensive exact count while more pages exist; +- while a legacy backfill is running, pages return immediately without raw-table count or facet scans and expose indexing progress; - collector inserts and rollup updates are batched in one transaction; - no activity retention deletion is performed. @@ -434,6 +436,7 @@ Collection starts even when the web UI is disabled. Existing databases remain qu - Upload progress is logged per file every `BACKUP_PROGRESS_INTERVAL_SECONDS` seconds and again when a file reaches 100%, including percentage and transferred/remaining bytes with auto-scaled units. - `BACKUP_PROGRESS=auto` shows an interactive multi-line progress view only for TTY/manual runs. Current file uploads are shown as separate rows, capped at 12 rows, with the total progress row at the bottom. Use `always` to force it or `never` to suppress the bar. File and total progress are still logged. - Every run atomically updates `BACKUP_STATUS_FILE` (default `/state/backup-status.json`) with its state, current source, active files, byte totals, percentage, snapshot, and final result for the live web view. +- Active status is reconciled against the kernel-held `/state/backup.lock`; if a container restart or forced termination releases the lock while status still says `starting` or `running`, the web API atomically records the run as interrupted instead of reporting a phantom backup. - Rclone-backed destinations cap concurrent file transfers at 12. Rsync remains single-streamed by rsync itself. - Retention logic: - daily: newest N snapshots diff --git a/app/audit_collector.py b/app/audit_collector.py index 7c9f521..d92c5e2 100644 --- a/app/audit_collector.py +++ b/app/audit_collector.py @@ -161,13 +161,16 @@ def main() -> int: log(f"Watching {SAMBA_LOG_GLOB}; storing events in SQLite") try: while not STOP: + backfilled = 0 try: count = collect_once(store) if count: log(f"Stored {count} event(s)") - store.backfill_rollups() + backfilled = store.backfill_rollups() except Exception as exc: # pylint: disable=broad-except log(f"Collector cycle failed: {exc}") + if backfilled: + continue time.sleep(POLL_SECONDS) finally: store.close() diff --git a/app/audit_store.py b/app/audit_store.py index ff085fe..19cd42d 100644 --- a/app/audit_store.py +++ b/app/audit_store.py @@ -16,7 +16,11 @@ try: ) from .state_db import connect_state_db except ImportError: - from audit_policy import account_name, read_deduplication_key, skipped_user_suffixes + from audit_policy import ( + account_name, + read_deduplication_key, + skipped_user_suffixes, + ) from state_db import connect_state_db @@ -43,20 +47,46 @@ CREATE TABLE IF NOT EXISTS audit_sources ( inode INTEGER NOT NULL, offset INTEGER NOT NULL ) WITHOUT ROWID; -CREATE INDEX IF NOT EXISTS audit_events_time - ON audit_events (occurred_second DESC, id DESC); -CREATE INDEX IF NOT EXISTS audit_events_action_time - ON audit_events (action, occurred_second DESC, id DESC); -CREATE INDEX IF NOT EXISTS audit_events_success_time - ON audit_events (success, occurred_second DESC, id DESC); -CREATE INDEX IF NOT EXISTS audit_events_user_time - ON audit_events (user COLLATE NOCASE, occurred_second DESC, id DESC); -CREATE INDEX IF NOT EXISTS audit_events_account_time - ON audit_events (account, occurred_second DESC, id DESC); -CREATE INDEX IF NOT EXISTS audit_events_share_time - ON audit_events (share COLLATE NOCASE, occurred_second DESC, id DESC); -CREATE INDEX IF NOT EXISTS audit_events_result_time - ON audit_events (result COLLATE NOCASE, occurred_second DESC, id DESC); +CREATE INDEX IF NOT EXISTS audit_events_main_time + ON audit_events (occurred_second DESC, id DESC) + WHERE share <> 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_main_action_time + ON audit_events (action, occurred_second DESC, id DESC) + WHERE share <> 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_main_success_time + ON audit_events (success, occurred_second DESC, id DESC) + WHERE share <> 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_main_user_time + ON audit_events (user COLLATE NOCASE, occurred_second DESC, id DESC) + WHERE share <> 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_main_account_time + ON audit_events (account, occurred_second DESC, id DESC) + WHERE share <> 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_main_share_time + ON audit_events (share COLLATE NOCASE, occurred_second DESC, id DESC) + WHERE share <> 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_main_result_time + ON audit_events (result COLLATE NOCASE, occurred_second DESC, id DESC) + WHERE share <> 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_fslogix_time + ON audit_events (occurred_second DESC, id DESC) + WHERE share = 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_fslogix_action_time + ON audit_events (action, occurred_second DESC, id DESC) + WHERE share = 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_fslogix_success_time + ON audit_events (success, occurred_second DESC, id DESC) + WHERE share = 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_fslogix_user_time + ON audit_events (user COLLATE NOCASE, occurred_second DESC, id DESC) + WHERE share = 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_fslogix_account_time + ON audit_events (account, occurred_second DESC, id DESC) + WHERE share = 'FSLogix' COLLATE NOCASE; +CREATE INDEX IF NOT EXISTS audit_events_fslogix_result_time + ON audit_events (result COLLATE NOCASE, occurred_second DESC, id DESC) + WHERE share = 'FSLogix' COLLATE NOCASE; + CREATE TABLE IF NOT EXISTS audit_paths ( id INTEGER PRIMARY KEY, path TEXT NOT NULL UNIQUE @@ -107,7 +137,8 @@ CREATE TABLE IF NOT EXISTS audit_rollup_state ( singleton INTEGER PRIMARY KEY CHECK (singleton = 1), backfill_next_id INTEGER NOT NULL, backfill_max_id INTEGER NOT NULL, - ready INTEGER NOT NULL CHECK (ready IN (0, 1)) + ready INTEGER NOT NULL CHECK (ready IN (0, 1)), + policy_version INTEGER NOT NULL DEFAULT 3 ); """ @@ -135,6 +166,17 @@ VALUES (?, ?, ?, ?) def ensure_audit_schema(conn: sqlite3.Connection) -> None: conn.executescript(AUDIT_SCHEMA) + for legacy_index in ( + "audit_events_time", + "audit_events_action_time", + "audit_events_success_time", + "audit_events_user_time", + "audit_events_account_time", + "audit_events_share_time", + "audit_events_result_time", + ): + conn.execute(f"DROP INDEX IF EXISTS {legacy_index}") + columns = { str(row["name"]) for row in conn.execute("PRAGMA table_info(audit_events)") @@ -142,6 +184,16 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None: if "path_id" not in columns: conn.execute("ALTER TABLE audit_events ADD COLUMN path_id INTEGER") + state_columns = { + str(row["name"]) + for row in conn.execute("PRAGMA table_info(audit_rollup_state)") + } + if "policy_version" not in state_columns: + conn.execute( + "ALTER TABLE audit_rollup_state " + "ADD COLUMN policy_version INTEGER NOT NULL DEFAULT 1" + ) + max_id = int( conn.execute("SELECT coalesce(max(id), 0) FROM audit_events").fetchone()[0] ) @@ -153,6 +205,29 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None: """, (max_id, int(max_id == 0)), ) + state = conn.execute( + """ + SELECT ready, policy_version + FROM audit_rollup_state + WHERE singleton = 1 + """ + ).fetchone() + if state is not None and int(state["policy_version"]) < 3: + for table in ( + "audit_daily_totals", + "audit_daily_counts", + "audit_daily_facets", + ): + conn.execute(f"DELETE FROM {table}") + conn.execute( + """ + UPDATE audit_rollup_state + SET backfill_next_id = 1, backfill_max_id = ?, ready = ?, + policy_version = 3 + WHERE singleton = 1 + """, + (max_id, int(max_id == 0)), + ) conn.commit() @@ -347,7 +422,7 @@ class AuditStore: raise return len(event_rows) - def backfill_rollups(self, limit: int = 10000) -> int: + def backfill_rollups(self, limit: int = 50000) -> int: """Backfill one bounded legacy-event chunk without delaying live appends.""" state = self.conn.execute( """ @@ -459,7 +534,8 @@ class AuditStore: ) SELECT ?, occurred_second / 86400, account, {column} FROM audit_events - WHERE id BETWEEN ? AND ? AND {column} <> '' + WHERE id BETWEEN ? AND ? + AND {column} <> '' GROUP BY occurred_second / 86400, account, {column} """, (kind, start_id, end_id), @@ -519,15 +595,28 @@ def excluded_account_conditions() -> Tuple[List[str], List[object]]: return conditions, values +def activity_stream_condition(stream: str) -> str: + if stream == "main": + return "share <> 'FSLogix' COLLATE NOCASE" + if stream == "fslogix": + return "share = 'FSLogix' COLLATE NOCASE" + raise ValueError("Ungültiger Aktivitätsstrom") + + def activity_conditions( start: dt.date, end: dt.date, params: Dict[str, List[str]], *, + stream: str, include_filters: bool, include_path: bool = True, ) -> Tuple[List[str], List[object]]: - conditions = ["occurred_second >= ?", "occurred_second < ?"] + conditions = [ + "occurred_second >= ?", + "occurred_second < ?", + activity_stream_condition(stream), + ] values: List[object] = [ date_seconds(start), date_seconds(end + dt.timedelta(days=1)), @@ -568,14 +657,44 @@ def activity_conditions( return conditions, values -def rollups_ready(conn: sqlite3.Connection) -> bool: +def rollup_status(conn: sqlite3.Connection) -> Dict[str, object]: try: row = conn.execute( - "SELECT ready FROM audit_rollup_state WHERE singleton = 1" + """ + SELECT backfill_next_id, backfill_max_id, ready + FROM audit_rollup_state + WHERE singleton = 1 + """ ).fetchone() except sqlite3.Error: - return False - return row is not None and bool(row[0]) + row = None + if row is None: + return { + "ready": False, + "processedEvents": 0, + "totalEvents": 0, + "percent": 0.0, + } + + total = max(0, int(row["backfill_max_id"])) + ready = bool(row["ready"]) + processed = ( + total + if ready + else min(total, max(0, int(row["backfill_next_id"]) - 1)) + ) + percent = ( + 100.0 + if ready or total == 0 + else round(processed * 100 / total, 1) + ) + return { + "ready": ready, + "processedEvents": processed, + "totalEvents": total, + "percent": percent, + } + def rollup_range( @@ -597,44 +716,40 @@ def rollup_matched( start: dt.date, end: dt.date, params: Dict[str, List[str]], + *, + stream: str, ) -> Optional[int]: filters = activity_filters(params) if filters["path"]: return None conditions, values = rollup_range(start, end) - has_filters = any( - filters[name] for name in ("user", "share", "operation", "action", "result") - ) - if not has_filters: - table = "audit_daily_totals" - else: - table = "audit_daily_counts" - user = filters["user"] - if user: - if "\\" in user or "@" in user: - conditions.append("user = ? COLLATE NOCASE") - else: - conditions.append("account = ?") - values.append(user) - if filters["share"]: - conditions.append("share = ? COLLATE NOCASE") - values.append(filters["share"]) - action = filters["action"] or filters["operation"] - if action: - conditions.append("action = ?") - values.append(action) - result = filters["result"] - if result == "fail": - conditions.append("success = 0") - elif result: - conditions.append("result = ? COLLATE NOCASE") - values.append(result) + conditions.append(activity_stream_condition(stream)) + user = filters["user"] + if user: + if "\\" in user or "@" in user: + conditions.append("user = ? COLLATE NOCASE") + else: + conditions.append("account = ?") + values.append(user) + if filters["share"]: + conditions.append("share = ? COLLATE NOCASE") + values.append(filters["share"]) + action = filters["action"] or filters["operation"] + if action: + conditions.append("action = ?") + values.append(action) + result = filters["result"] + if result == "fail": + conditions.append("success = 0") + elif result: + conditions.append("result = ? COLLATE NOCASE") + values.append(result) row = conn.execute( f""" SELECT coalesce(sum(event_count), 0) - FROM {table} + FROM audit_daily_counts WHERE {" AND ".join(conditions)} """, values, @@ -643,56 +758,30 @@ def rollup_matched( def rollup_facets( - conn: sqlite3.Connection, start: dt.date, end: dt.date -) -> Dict[str, List[str]]: - range_conditions, range_values = rollup_range(start, end) - - def distinct(kind: str) -> List[str]: - conditions = ["kind = ?", *range_conditions, "value <> ''"] - return [ - str(row[0]) - for row in conn.execute( - f""" - SELECT DISTINCT value - FROM audit_daily_facets - WHERE {" AND ".join(conditions)} - ORDER BY value COLLATE NOCASE - """, - (kind, *range_values), - ) - ] - - actions = distinct("action") - return { - "users": distinct("user"), - "shares": distinct("share"), - "operations": actions, - "actions": actions, - } - - -def raw_facets( conn: sqlite3.Connection, start: dt.date, end: dt.date, - params: Dict[str, List[str]], + *, + stream: str, ) -> Dict[str, List[str]]: - conditions, values = activity_conditions( - start, end, params, include_filters=False - ) - where_sql = " AND ".join(conditions) + range_conditions, range_values = rollup_range(start, end) def distinct(column: str) -> List[str]: + conditions = [ + *range_conditions, + activity_stream_condition(stream), + f"{column} <> ?", + ] return [ str(row[0]) for row in conn.execute( f""" SELECT DISTINCT {column} - FROM audit_events - WHERE {where_sql} AND {column} <> '' + FROM audit_daily_counts + WHERE {" AND ".join(conditions)} ORDER BY {column} COLLATE NOCASE """, - values, + (*range_values, ""), ) ] @@ -764,9 +853,12 @@ def query_activity( start: dt.date, end: dt.date, params: Dict[str, List[str]], + *, + stream: str = "main", ) -> Dict[str, object]: limit = min(500, max(1, int(params.get("limit", ["100"])[0]))) - ready = rollups_ready(conn) + indexing = rollup_status(conn) + ready = bool(indexing["ready"]) path = activity_filters(params)["path"] path_ids = indexed_path_ids(conn, path) if ready and path else None if ( @@ -779,6 +871,7 @@ def query_activity( start, end, params, + stream=stream, include_filters=True, include_path=path_ids is None, ) @@ -854,32 +947,23 @@ def query_activity( include_count = query_flag(params, "count") matched: Optional[int] = None - if include_count: - if ready: - matched = rollup_matched(conn, start, end, params) - else: - matched = int( - conn.execute( - f""" - SELECT count(*) - FROM audit_events - WHERE {" AND ".join(conditions)} - """, - values, - ).fetchone()[0] - ) - if matched is None and not cursor and not has_more: - matched = len(events) + if include_count and ready: + matched = rollup_matched(conn, start, end, params, stream=stream) + if include_count and matched is None and not cursor and not has_more: + matched = len(events) + empty_facets = { + "users": [], + "shares": [], + "operations": [], + "actions": [], + } include_facets = query_flag(params, "facets") - if include_facets: - facets = ( - rollup_facets(conn, start, end) - if ready - else raw_facets(conn, start, end, params) - ) - else: - facets = {"users": [], "shares": [], "operations": [], "actions": []} + facets = ( + rollup_facets(conn, start, end, stream=stream) + if include_facets and ready + else empty_facets + ) return { "events": events, @@ -888,40 +972,54 @@ def query_activity( "matched": matched, "matchedExact": matched is not None, "facets": facets, + "indexing": indexing, + "stream": stream, } def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, object]: - if rollups_ready(conn): + indexing = rollup_status(conn) + + def second_day(value: object) -> Optional[str]: + if value is None: + return None + return dt.datetime.fromtimestamp( + int(value), tz=dt.timezone.utc + ).date().isoformat() + + if indexing["ready"]: row = conn.execute( """ SELECT count(DISTINCT day), min(day), max(day) - FROM audit_daily_totals + FROM audit_daily_counts + WHERE share <> 'FSLogix' COLLATE NOCASE """ ).fetchone() - - def day_text(value: object) -> Optional[str]: - if value is None: - return None - return dt.datetime.fromtimestamp( - int(value) * 86400, tz=dt.timezone.utc - ).date().isoformat() - days = int(row[0] or 0) - oldest = day_text(row[1]) - newest = day_text(row[2]) + oldest = second_day(None if row[1] is None else int(row[1]) * 86400) + newest = second_day(None if row[2] is None else int(row[2]) * 86400) else: - row = conn.execute( + oldest_row = conn.execute( """ - SELECT count(DISTINCT substr(occurred_at, 1, 10)), - min(substr(occurred_at, 1, 10)), - max(substr(occurred_at, 1, 10)) - FROM audit_events + SELECT occurred_second FROM audit_events + WHERE share <> 'FSLogix' COLLATE NOCASE + ORDER BY occurred_second, id LIMIT 1 """ ).fetchone() - days = int(row[0] or 0) - oldest = row[1] - newest = row[2] + newest_row = conn.execute( + """ + SELECT occurred_second FROM audit_events + WHERE share <> 'FSLogix' COLLATE NOCASE + ORDER BY occurred_second DESC, id DESC LIMIT 1 + """ + ).fetchone() + oldest = second_day(oldest_row[0]) if oldest_row else None + newest = second_day(newest_row[0]) if newest_row else None + days = ( + (dt.date.fromisoformat(newest) - dt.date.fromisoformat(oldest)).days + 1 + if oldest and newest + else 0 + ) database_bytes = 0 for suffix in ("", "-wal"): @@ -934,6 +1032,7 @@ def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, obj "bytes": database_bytes, "oldest": oldest, "newest": newest, + "indexing": indexing, } diff --git a/app/backup_to_destination.py b/app/backup_to_destination.py index c64a33a..2277997 100644 --- a/app/backup_to_destination.py +++ b/app/backup_to_destination.py @@ -204,6 +204,7 @@ class BackupStatus: self.write( state="starting", trigger=trigger, + workerPid=os.getpid(), startedAt=dt.datetime.now(dt.timezone.utc).isoformat(timespec="seconds"), finishedAt=None, destination=destination, @@ -241,6 +242,7 @@ class BackupStatus: def complete(self, message: str) -> None: self.write( state="completed", + workerPid=None, finishedAt=dt.datetime.now(dt.timezone.utc).isoformat(timespec="seconds"), percent=100.0, transferredBytes=self.value.get("totalBytes", 0), @@ -252,6 +254,7 @@ class BackupStatus: def fail(self, message: str) -> None: self.write( state="failed", + workerPid=None, finishedAt=dt.datetime.now(dt.timezone.utc).isoformat(timespec="seconds"), activeFiles=[], message=message, diff --git a/app/web/app.js b/app/web/app.js index 11e0586..49dd6e6 100644 --- a/app/web/app.js +++ b/app/web/app.js @@ -43,6 +43,7 @@ function backupMessage(data) { if (message.startsWith("Syncing ")) return `${message.slice(8)} wird synchronisiert.`; const completed = message.match(/^Backup completed; (\d+) snapshot\(s\) retained$/); if (completed) return `Sicherung abgeschlossen; ${completed[1]} Sicherungsstände werden aufbewahrt.`; + if (data.interrupted) return "Sicherung wurde durch einen Neustart oder Prozessabbruch unterbrochen."; if (data.state === "failed") return "Sicherung fehlgeschlagen. Einzelheiten stehen im Sicherungsprotokoll."; if (data.state === "completed") return "Sicherung abgeschlossen."; if (data.state === "running") return "Sicherung läuft."; @@ -102,6 +103,7 @@ function routeFor(path) { if (path.startsWith("/storage/users")) return "storage-users"; if (path.startsWith("/shares")) return "shares"; if (path.startsWith("/reconciliation")) return "reconciliation"; + if (path.startsWith("/activity/fslogix")) return "activity-fslogix"; if (path.startsWith("/activity")) return "activity"; if (path.startsWith("/backup")) return "backup"; if (path.startsWith("/report")) return "report"; @@ -124,6 +126,7 @@ async function navigate(path, replace = false) { if (route === "storage-data") await renderStorage("data"); if (route === "storage-users") await renderStorage("users"); if (route === "activity") await renderActivity(); + if (route === "activity-fslogix") await renderActivity("fslogix"); if (route === "backup") await renderBackup(); if (route === "report") await renderReport(); if (route === "system") await renderSystem(); @@ -238,10 +241,17 @@ async function renderStorage(type) { draw(); } -async function renderActivity() { +async function renderActivity(stream = "main") { const today = new Date(); const yesterday = new Date(Date.now() - 86400000); - content.innerHTML = pageHead("Aktivitätsprotokoll", "Lese-, Schreib-, Umbenennungs- und Löschvorgänge nach Datum, Identität, Freigabe, Ergebnis oder Pfad durchsuchen.") + ` + const fslogix = stream === "fslogix"; + const endpoint = fslogix ? "/api/fslogix-activity" : "/api/activity"; + content.innerHTML = pageHead( + fslogix ? "FSLogix-Protokoll" : "Aktivitätsprotokoll", + fslogix + ? "FSLogix-Profilcontainer-Aktivitäten getrennt vom regulären Dateiprotokoll durchsuchen." + : "Lese-, Schreib-, Umbenennungs- und Löschvorgänge nach Datum, Identität, Freigabe, Ergebnis oder Pfad durchsuchen.", + ) + `
@@ -270,14 +280,15 @@ async function renderActivity() { params.set("count", "0"); params.set("facets", "0"); } - const result = await api(`/api/activity?${params}`); + const result = await api(`${endpoint}?${params}`); state.activity = result; const rows = document.querySelector("#activity-rows"); if (append) rows.insertAdjacentHTML("beforeend", eventRows(result.events)); else rows.innerHTML = eventRows(result.events); loaded = append ? loaded + result.events.length : result.events.length; if (!append) matched = result.matchedExact ? result.matched : null; const shown = matched === null ? `${loaded.toLocaleString("de-DE")}${result.hasMore ? "+" : ""}` : matched.toLocaleString("de-DE"); - document.querySelector("#activity-summary").textContent = `${shown} passende Ereignisse`; + const indexNote = result.indexing?.ready === false ? ` · Indexaufbau ${Number(result.indexing.percent || 0).toLocaleString("de-DE")} %` : ""; + document.querySelector("#activity-summary").textContent = `${shown} passende Ereignisse${indexNote}`; const more = document.querySelector("#load-more"); cursor = result.nextCursor || 0; more.hidden = !result.hasMore; diff --git a/app/web/index.html b/app/web/index.html index 44d78da..aa244b6 100644 --- a/app/web/index.html +++ b/app/web/index.html @@ -32,6 +32,7 @@ Datenbelegung Benutzerbelegung Aktivitätsprotokoll + FSLogix-Protokoll Sicherungen PDF-Bericht System diff --git a/app/web/report.mjs b/app/web/report.mjs index fec2126..4e11b5e 100644 --- a/app/web/report.mjs +++ b/app/web/report.mjs @@ -102,6 +102,7 @@ function backupMessage(data) { if (message.startsWith("Syncing ")) return message.slice(8) + " wird synchronisiert."; const completed = message.match(/^Backup completed; (\d+) snapshot\(s\) retained$/); if (completed) return "Sicherung abgeschlossen; " + completed[1] + " Sicherungsstände werden aufbewahrt."; + if (data.interrupted) return "Sicherung wurde durch einen Neustart oder Prozessabbruch unterbrochen."; if (data.state === "failed") return "Sicherung fehlgeschlagen."; if (data.state === "completed") return "Sicherung abgeschlossen."; if (data.state === "running") return "Sicherung läuft."; diff --git a/app/web_ui.py b/app/web_ui.py index a78ec93..053deeb 100644 --- a/app/web_ui.py +++ b/app/web_ui.py @@ -60,7 +60,7 @@ STATIC_ROOT = os.getenv("WEB_STATIC_DIR", "/app/web") STATE_DB = STATE_DB_PATH BACKUP_STATUS_FILE = os.getenv("BACKUP_STATUS_FILE", "/state/backup-status.json") BACKUP_LOG_FILE = os.getenv("BACKUP_LOG_FILE", "/var/log/backup.log") -BACKUP_LOCK_FILE = "/state/backup.lock" +BACKUP_LOCK_FILE = os.path.join(os.getenv("STATE_ROOT", "/state"), "backup.lock") RECONCILE_STATUS_FILE = os.getenv( "RECONCILE_STATUS_FILE", "/state/reconcile-status.json" ) @@ -76,6 +76,7 @@ SID_RE = re.compile(r"S-\d+(?:-\d+)+", re.IGNORECASE) LOGIN_LIMIT: Dict[str, deque] = {} LOGIN_LIMIT_LOCK = threading.Lock() ACTION_LAUNCH_LOCK = threading.Lock() +ACTIVE_PROCESS_STATES = frozenset({"starting", "running"}) def log(message: str) -> None: @@ -176,6 +177,69 @@ def read_json(path: str, default): return default +def write_json_atomic(path: str, value: Dict[str, object]) -> None: + directory = os.path.dirname(path) + if directory: + os.makedirs(directory, exist_ok=True) + temp_path = f"{path}.recovery.tmp" + try: + with open(temp_path, "w", encoding="utf-8") as handle: + json.dump(value, handle, separators=(",", ":"), sort_keys=True) + handle.flush() + os.fsync(handle.fileno()) + os.replace(temp_path, path) + except (OSError, TypeError, ValueError): + try: + os.remove(temp_path) + except OSError: + pass + raise + + +def reconcile_backup_process_status(value: Dict[str, object]) -> Dict[str, object]: + """Replace stale active state when no worker owns the backup lock.""" + if str(value.get("state", "")) not in ACTIVE_PROCESS_STATES: + value["processRunning"] = False + return value + + try: + lock_dir = os.path.dirname(BACKUP_LOCK_FILE) + if lock_dir: + os.makedirs(lock_dir, exist_ok=True) + with open(BACKUP_LOCK_FILE, "a+", encoding="utf-8") as lock_file: + try: + fcntl.flock(lock_file, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + value["processRunning"] = True + return value + + latest = read_json(BACKUP_STATUS_FILE, value) + if str(latest.get("state", "")) in ACTIVE_PROCESS_STATES: + latest.update( + { + "state": "failed", + "finishedAt": now_utc().isoformat(timespec="seconds"), + "activeFiles": [], + "currentSource": None, + "workerPid": None, + "interrupted": True, + "message": ( + "Backup interrupted because the worker process " + "is no longer running" + ), + } + ) + write_json_atomic(BACKUP_STATUS_FILE, latest) + log("Recovered stale backup status after worker interruption") + latest["processRunning"] = False + fcntl.flock(lock_file, fcntl.LOCK_UN) + return latest + except OSError as exc: + log(f"Unable to verify backup worker state: {exc}") + value["processRunning"] = None + return value + + def base64url(value: bytes) -> str: return base64.urlsafe_b64encode(value).rstrip(b"=").decode("ascii") @@ -711,7 +775,9 @@ class DirectoryCache: -def query_audit(params: Dict[str, List[str]]) -> Dict[str, object]: +def query_audit( + params: Dict[str, List[str]], *, stream: str = "main" +) -> Dict[str, object]: today = now_utc().date() default_start = today - dt.timedelta(days=1) try: @@ -724,7 +790,7 @@ def query_audit(params: Dict[str, List[str]]) -> Dict[str, object]: raise ValueError(f"Der Datumsbereich darf höchstens {max_days} Tage umfassen") conn = connect_state_db(STATE_DB, read_only=True) try: - return query_activity(conn, start, end, params) + return query_activity(conn, start, end, params, stream=stream) finally: conn.close() @@ -747,6 +813,7 @@ def tail_lines(path: str, count: int) -> List[str]: def backup_payload(include_log: bool = True) -> Dict[str, object]: value = read_json(BACKUP_STATUS_FILE, {}) + value = reconcile_backup_process_status(value) value.pop("log", None) configured = bool(os.getenv("BACKUP_DESTINATION", "").strip()) automatic = configured and env_bool("BACKUP_AUTO_ENABLED", True) @@ -837,7 +904,17 @@ class App: def overview(self) -> Dict[str, object]: usage = self.usage.snapshot() recent = query_audit({"limit": ["12"], "facets": ["0"]}) - return {"usage": usage, "activeGroups": share_count(), "recentEvents": recent["events"], "eventCount": recent["matched"], "backup": backup_payload(), "audit": audit_archive_summary()} + event_count = recent["matched"] + if event_count is None: + event_count = f"{len(recent['events'])}+" + return { + "usage": usage, + "activeGroups": share_count(), + "recentEvents": recent["events"], + "eventCount": event_count, + "backup": backup_payload(), + "audit": audit_archive_summary(), + } def system_summary(self) -> Dict[str, object]: checks = {} @@ -1026,6 +1103,8 @@ class Handler(BaseHTTPRequestHandler): self.send_json(APP.usage.snapshot()) elif path == "/api/activity": self.send_json(query_audit(params)) + elif path == "/api/fslogix-activity": + self.send_json(query_audit(params, stream="fslogix")) elif path == "/api/backup": self.send_json( backup_payload(include_log=query_includes_log(params)) diff --git a/tests/test_backup_to_destination.py b/tests/test_backup_to_destination.py index b7f0819..61dd12d 100644 --- a/tests/test_backup_to_destination.py +++ b/tests/test_backup_to_destination.py @@ -32,6 +32,22 @@ class BackupLoggerTests(unittest.TestCase): ) +class BackupStatusTests(unittest.TestCase): + def test_status_records_and_clears_worker_pid(self): + with tempfile.TemporaryDirectory() as tmpdir: + status = backup.BackupStatus( + os.path.join(tmpdir, "backup-status.json") + ) + status.begin("rsync://backup.test/target", "manual") + self.assertEqual(status.value["workerPid"], os.getpid()) + + status.complete("done") + + self.assertIsNone(status.value["workerPid"]) + self.assertEqual(status.value["state"], "completed") + + + class ProgressConfigTests(unittest.TestCase): def test_parse_progress_mode_accepts_known_values(self): with mock.patch.dict(os.environ, {"BACKUP_PROGRESS": "ALWAYS"}): diff --git a/tests/test_web_ui.py b/tests/test_web_ui.py index 48c0946..94f8d67 100644 --- a/tests/test_web_ui.py +++ b/tests/test_web_ui.py @@ -102,6 +102,77 @@ class AdministrationActionTests(unittest.TestCase): self.assertEqual(payload["state"], "waiting") self.assertEqual(payload["log"], ["backup log"]) + @mock.patch("app.web_ui.tail_lines", return_value=[]) + def test_backup_payload_marks_unowned_active_status_interrupted( + self, _tail_lines + ): + with tempfile.TemporaryDirectory() as tmpdir: + status_path = os.path.join(tmpdir, "backup-status.json") + lock_path = os.path.join(tmpdir, "backup.lock") + web_ui.write_json_atomic( + status_path, + { + "state": "running", + "workerPid": 4242, + "currentSource": "data/fslogix", + "activeFiles": [{"path": "profile.vhdx"}], + }, + ) + with ( + mock.patch.object(web_ui, "BACKUP_STATUS_FILE", status_path), + mock.patch.object(web_ui, "BACKUP_LOCK_FILE", lock_path), + mock.patch.dict( + os.environ, + {"BACKUP_DESTINATION": "rsync://backup.test/target"}, + clear=True, + ), + ): + payload = web_ui.backup_payload() + persisted = web_ui.read_json(status_path, {}) + + self.assertEqual(payload["state"], "failed") + self.assertFalse(payload["processRunning"]) + self.assertTrue(payload["interrupted"]) + self.assertIsNone(payload["workerPid"]) + self.assertIsNone(payload["currentSource"]) + self.assertEqual(payload["activeFiles"], []) + self.assertIsNotNone(payload["finishedAt"]) + self.assertEqual(persisted["state"], "failed") + self.assertTrue(persisted["interrupted"]) + + @mock.patch("app.web_ui.tail_lines", return_value=[]) + def test_backup_payload_keeps_status_active_while_lock_is_owned( + self, _tail_lines + ): + with tempfile.TemporaryDirectory() as tmpdir: + status_path = os.path.join(tmpdir, "backup-status.json") + lock_path = os.path.join(tmpdir, "backup.lock") + web_ui.write_json_atomic( + status_path, {"state": "running", "workerPid": 4242} + ) + with open(lock_path, "w", encoding="utf-8") as lock_file: + web_ui.fcntl.flock(lock_file, web_ui.fcntl.LOCK_EX) + with ( + mock.patch.object( + web_ui, "BACKUP_STATUS_FILE", status_path + ), + mock.patch.object(web_ui, "BACKUP_LOCK_FILE", lock_path), + mock.patch.dict( + os.environ, + {"BACKUP_DESTINATION": "rsync://backup.test/target"}, + clear=True, + ), + ): + payload = web_ui.backup_payload() + web_ui.fcntl.flock(lock_file, web_ui.fcntl.LOCK_UN) + + persisted = web_ui.read_json(status_path, {}) + + self.assertEqual(payload["state"], "running") + self.assertTrue(payload["processRunning"]) + self.assertEqual(persisted["state"], "running") + + @mock.patch("app.web_ui.tail_lines", return_value=["reconcile log"]) @mock.patch("app.web_ui.read_json", return_value={}) def test_reconciliation_payload_has_progress_schedule_and_log( @@ -343,6 +414,20 @@ class AuditParsingTests(unittest.TestCase): audit_collector.parse_audit_line(line, "/var/log/samba/log.pc01") ) + def test_parses_fslogix_activity_for_separate_storage(self): + line = ( + "smbd_audit: x|alice|192.0.2.5|PC01|" + "fSlOgIx|pwrite|OK|profile.vhdx\n" + ) + + event = audit_collector.parse_audit_line( + line, "/var/log/samba/log.pc01" + ) + + self.assertIsNotNone(event) + self.assertEqual(event["share"], "fSlOgIx") + self.assertEqual(event["path"], "profile.vhdx") + def test_tracks_rotated_file_by_inode_without_reingesting_it(self): with tempfile.TemporaryDirectory() as tmpdir: active = os.path.join(tmpdir, "log.pc01") @@ -551,6 +636,57 @@ class AuditQueryTests(unittest.TestCase): self.assertNotIn("robot_svc", visible["facets"]["users"]) self.assertEqual(set(visible["facets"]["actions"]), {"read", "delete"}) + def test_store_separates_main_and_fslogix_streams(self): + with tempfile.TemporaryDirectory() as tmpdir: + today = dt.datetime.now(dt.timezone.utc).date() + store = audit_store.AuditStore(os.path.join(tmpdir, "state.db")) + main_event = self.make_event(f"{today}T12:00:00+00:00", "alice") + fslogix_event = dict(main_event, share="FSLogix", path="profile.vhdx") + try: + inserted = store.append_batch( + [fslogix_event, main_event], {}, set() + ) + main = audit_store.query_activity( + store.conn, today, today, {"limit": ["100"]} + ) + fslogix = audit_store.query_activity( + store.conn, + today, + today, + {"limit": ["100"]}, + stream="fslogix", + ) + plans = {} + for stream, operator, index in ( + ("main", "<>", "audit_events_main_time"), + ("fslogix", "=", "audit_events_fslogix_time"), + ): + plans[stream] = " ".join( + row["detail"] + for row in store.conn.execute( + f""" + EXPLAIN QUERY PLAN + SELECT id FROM audit_events + WHERE occurred_second >= ? AND occurred_second < ? + AND share {operator} 'FSLogix' COLLATE NOCASE + ORDER BY occurred_second DESC, id DESC + LIMIT 101 + """, + (0, 9999999999), + ) + ) + self.assertIn(index, plans[stream]) + finally: + store.close() + + self.assertEqual(inserted, 2) + self.assertEqual(main["matched"], 1) + self.assertEqual(main["events"][0]["share"], "Data") + self.assertEqual(main["stream"], "main") + self.assertEqual(fslogix["matched"], 1) + self.assertEqual(fslogix["events"][0]["share"], "FSLogix") + self.assertEqual(fslogix["stream"], "fslogix") + def test_query_metadata_uses_rollups_instead_of_raw_event_scans(self): with tempfile.TemporaryDirectory() as tmpdir: today = dt.datetime.now(dt.timezone.utc).date() @@ -575,10 +711,10 @@ class AuditQueryTests(unittest.TestCase): self.assertTrue(result["matchedExact"]) self.assertTrue(result["hasMore"]) self.assertTrue( - any("from audit_daily_totals" in statement for statement in normalized) + any("from audit_daily_counts" in statement for statement in normalized) ) self.assertTrue( - any("from audit_daily_facets" in statement for statement in normalized) + any("from audit_daily_counts" in statement for statement in normalized) ) self.assertFalse( any( @@ -593,6 +729,62 @@ class AuditQueryTests(unittest.TestCase): ) ) + def test_pending_rollups_never_scan_raw_metadata(self): + with tempfile.TemporaryDirectory() as tmpdir: + today = dt.datetime.now(dt.timezone.utc).date() + database = os.path.join(tmpdir, "state.db") + store = audit_store.AuditStore(database) + events = [ + self.make_event(f"{today}T12:00:0{index}+00:00", "alice") + for index in range(4) + ] + store.append_batch(events, {}, set()) + for table in ( + "audit_daily_totals", + "audit_daily_counts", + "audit_daily_facets", + "audit_rollup_state", + ): + store.conn.execute(f"DELETE FROM {table}") + store.conn.commit() + store.close() + + store = audit_store.AuditStore(database) + statements = [] + store.conn.set_trace_callback(statements.append) + try: + result = audit_store.query_activity( + store.conn, today, today, {"limit": ["1"]} + ) + summary = audit_store.audit_summary(store.conn, database) + finally: + store.conn.set_trace_callback(None) + store.close() + + normalized = [ + " ".join(statement.casefold().split()) + for statement in statements + ] + self.assertTrue(result["hasMore"]) + self.assertIsNone(result["matched"]) + self.assertFalse(result["matchedExact"]) + self.assertEqual(result["facets"]["users"], []) + self.assertFalse(result["indexing"]["ready"]) + self.assertFalse(summary["indexing"]["ready"]) + self.assertFalse( + any( + "select count(*) from audit_events" in statement + for statement in normalized + ) + ) + self.assertFalse( + any( + "select distinct" in statement and "from audit_events" in statement + for statement in normalized + ) + ) + + def test_path_query_does_not_block_on_an_exact_count(self): with tempfile.TemporaryDirectory() as tmpdir: today = dt.datetime.now(dt.timezone.utc).date() @@ -697,6 +889,65 @@ class AuditQueryTests(unittest.TestCase): self.assertEqual(path_event_count, 4) + def test_upgrade_rebuilds_main_and_fslogix_rollups(self): + with tempfile.TemporaryDirectory() as tmpdir: + today = dt.datetime.now(dt.timezone.utc).date() + database = os.path.join(tmpdir, "state.db") + store = audit_store.AuditStore(database) + main_event = self.make_event(f"{today}T12:00:00+00:00", "alice") + fslogix_event = dict( + main_event, share="FSLogix", path="profile.vhdx" + ) + store.append_batch([main_event, fslogix_event], {}, set()) + store.conn.execute("DROP TABLE audit_rollup_state") + store.conn.execute( + """ + CREATE TABLE audit_rollup_state ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + backfill_next_id INTEGER NOT NULL, + backfill_max_id INTEGER NOT NULL, + ready INTEGER NOT NULL CHECK (ready IN (0, 1)) + ) + """ + ) + store.conn.execute( + "INSERT INTO audit_rollup_state VALUES (1, 3, 2, 1)" + ) + store.conn.commit() + store.close() + + store = audit_store.AuditStore(database) + try: + self.assertEqual(store.backfill_rollups(limit=1), 1) + self.assertEqual(store.backfill_rollups(limit=1), 1) + self.assertEqual(store.backfill_rollups(limit=1), 0) + main = audit_store.query_activity( + store.conn, today, today, {"limit": ["1"]} + ) + fslogix = audit_store.query_activity( + store.conn, + today, + today, + {"limit": ["1"]}, + stream="fslogix", + ) + total = store.conn.execute( + "SELECT sum(event_count) FROM audit_daily_totals" + ).fetchone()[0] + policy_version = store.conn.execute( + """ + SELECT policy_version FROM audit_rollup_state + WHERE singleton = 1 + """ + ).fetchone()[0] + finally: + store.close() + + self.assertEqual(main["matched"], 1) + self.assertEqual(fslogix["matched"], 1) + self.assertEqual(total, 2) + self.assertEqual(policy_version, 3) + def test_schema_has_filter_indexes_and_rollup_tables(self): with tempfile.TemporaryDirectory() as tmpdir: store = audit_store.AuditStore(os.path.join(tmpdir, "state.db")) @@ -717,13 +968,19 @@ class AuditQueryTests(unittest.TestCase): store.close() self.assertTrue( { - "audit_events_time", - "audit_events_action_time", - "audit_events_success_time", - "audit_events_user_time", - "audit_events_account_time", - "audit_events_share_time", - "audit_events_result_time", + "audit_events_main_time", + "audit_events_main_action_time", + "audit_events_main_success_time", + "audit_events_main_user_time", + "audit_events_main_account_time", + "audit_events_main_share_time", + "audit_events_main_result_time", + "audit_events_fslogix_time", + "audit_events_fslogix_action_time", + "audit_events_fslogix_success_time", + "audit_events_fslogix_user_time", + "audit_events_fslogix_account_time", + "audit_events_fslogix_result_time", }.issubset(indexes) ) self.assertTrue( @@ -800,6 +1057,10 @@ class WebPresentationTests(unittest.TestCase): self.assertIn("Umbenennen", script) self.assertNotIn("Umbenennen/Verschieben", script) self.assertIn("Löschen", script) + self.assertIn("Indexaufbau", script) + self.assertIn("/activity/fslogix", html) + self.assertIn("/api/fslogix-activity", script) + self.assertIn('renderActivity("fslogix")', script) self.assertNotIn(">Auflisten<", script) self.assertNotIn(">Metadaten<", script) self.assertNotIn(">Sitzung<", script) @@ -822,6 +1083,7 @@ class WebPresentationTests(unittest.TestCase): self.assertIn('"/api/reconciliation?log=0"', script) self.assertEqual(script.count("if (active || previousActive)"), 2) self.assertNotIn("BACKUP_AUTO_ENABLED", script) + self.assertIn("data.interrupted", script) def test_pdf_report_is_client_side_and_uses_only_log_free_endpoint(self): script = self.asset("app.js") @@ -861,6 +1123,8 @@ class WebPresentationTests(unittest.TestCase): ] self.assertEqual(len(success_lines), 3) self.assertEqual(len(failure_lines), 3) + fslogix_config = config.split("[FSLogix]", 1)[1] + self.assertIn("full_audit", fslogix_config) for line in success_lines + failure_lines: self.assertEqual(set(line.split("=", 1)[1].split()), expected)