diff --git a/README.md b/README.md index d1e45e7..d3090f9 100644 --- a/README.md +++ b/README.md @@ -54,6 +54,7 @@ The database contains: - `shares`: AD group-to-folder lifecycle and ACL reconciliation state; - `audit_events`: normalized read, write, move, and delete events; - `audit_sources`: Samba log inode/offset checkpoints; +- `audit_read_dedup`: bounded, persistent fingerprints for restart-safe read deduplication; - `audit_daily_totals`, `audit_daily_counts`, and `audit_daily_facets`: compact materialized metadata for fast activity counts and filters; - `audit_paths`, `audit_paths_fts`, and `audit_path_events`: deduplicated trigram path search with an incrementally maintained event mapping; - `audit_rollup_state`: bounded legacy-event backfill progress; @@ -180,7 +181,7 @@ The E2E suite verifies: - SMB allow/deny behavior and real file operations; - Data, Private, and FSLogix usage aggregation; - high-level `full_audit` ingestion for all four actions, service-account exclusion, filters, facets, and pagination; -- shared SQLite schema, integrity, indexes, legacy-log removal, and ordered read deduplication; +- shared SQLite schema, integrity, indexes, legacy-log removal, and read deduplication across interleaved events and collector polls; - real rsync transfer progress, completed backup status, log output, remote snapshot marker, and per-group encrypted non-solid 7z archives; - anonymous action rejection plus authenticated manual backup and reconciliation actions, terminal progress, and live reconciliation output; - overview and system-health aggregation; @@ -399,7 +400,7 @@ Samba emits selected high-level `full_audit` operations for Private, Data, and F - 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; +- identical reads within the same UTC second are collapsed into one event regardless of intervening events, log source, or collector polling cycle; the persistent fingerprint cache covers the latest 48 hours and can be tuned with `AUDIT_READ_DEDUP_WINDOW_SECONDS`; - 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; @@ -407,7 +408,7 @@ Samba emits selected high-level `full_audit` operations for Private, Data, and F - collector inserts and rollup updates are batched in one transaction; - no activity retention deletion is performed. -Collection starts even when the web UI is disabled. Existing databases remain queryable while the collector backfills rollups, path IDs, and path-event mappings in bounded chunks after prioritizing each live append. On the first collector start after the legacy archive migration, recognized daily `.jsonl`/`.jsonl.gz` files and the old collector state file under `/state/audit` are deleted without import. Existing raw Samba log content is then processed using the current action, suffix, and deduplication policy. +Collection starts even when the web UI is disabled. Existing databases remain queryable while the collector backfills rollups, path IDs, path-event mappings, and recent read deduplication in bounded chunks after prioritizing each live append. On the first collector start after the legacy archive migration, recognized daily `.jsonl`/`.jsonl.gz` files and the old collector state file under `/state/audit` are deleted without import. Existing raw Samba log content is then processed using the current action, suffix, and deduplication policy. ## Backups diff --git a/app/audit_policy.py b/app/audit_policy.py index 18c1d42..3d373ba 100644 --- a/app/audit_policy.py +++ b/app/audit_policy.py @@ -2,10 +2,10 @@ """Shared policy for the small set of user-facing audit events.""" import datetime as dt +import hashlib import os from typing import Mapping, Optional, Tuple - AUDIT_ACTIONS = frozenset({"read", "write", "move", "delete"}) OPERATION_ACTIONS = { @@ -30,11 +30,9 @@ OPERATION_ACTIONS = { "rmdir": "delete", } - def action_for(operation: str) -> Optional[str]: return OPERATION_ACTIONS.get(operation.strip().casefold()) - def skipped_user_suffixes() -> Tuple[str, ...]: raw = os.getenv("AUDIT_SKIP_USER_SUFFIXES", "_svc,_ServiceAcc") return tuple( @@ -43,20 +41,16 @@ def skipped_user_suffixes() -> Tuple[str, ...]: if suffix.strip() ) - def account_name(user: str) -> str: account = user.strip().rsplit("\\", 1)[-1] return account.split("@", 1)[0] - def skip_user(user: str) -> bool: account = account_name(user).casefold() return any(account.endswith(suffix) for suffix in skipped_user_suffixes()) - ReadEventKey = Tuple[str, ...] - def read_deduplication_key(event: Mapping[str, object]) -> Optional[ReadEventKey]: if str(event.get("action", "")).casefold() != "read": return None @@ -84,3 +78,39 @@ def read_deduplication_key(event: Mapping[str, object]) -> Optional[ReadEventKey "success" if success else "failure", "" if success else str(event.get("result", "")), ) + +def read_deduplication_fingerprint( + event: Mapping[str, object] +) -> Optional[bytes]: + key = read_deduplication_key(event) + if key is None: + return None + digest = hashlib.blake2b(digest_size=16) + for value in key: + encoded = value.encode("utf-8", errors="surrogatepass") + digest.update(len(encoded).to_bytes(4, "big")) + digest.update(encoded) + return digest.digest() + +def read_deduplication_fingerprint_values( + timestamp: object, + user: object, + client_ip: object, + share: object, + path: object, + success: object, + result: object, +) -> bytes: + """Fingerprint legacy SQLite columns with the live-ingest policy.""" + fingerprint = read_deduplication_fingerprint({ + "action": "read", + "timestamp": timestamp, + "user": user, + "clientIp": client_ip, + "share": share, + "path": path, + "success": bool(success), + "result": result, + }) + assert fingerprint is not None + return fingerprint diff --git a/app/audit_store.py b/app/audit_store.py index 19cd42d..a736e54 100644 --- a/app/audit_store.py +++ b/app/audit_store.py @@ -11,14 +11,16 @@ from typing import Dict, Iterable, List, Optional, Set, Tuple try: from .audit_policy import ( account_name, - read_deduplication_key, + read_deduplication_fingerprint, + read_deduplication_fingerprint_values, skipped_user_suffixes, ) from .state_db import connect_state_db except ImportError: from audit_policy import ( account_name, - read_deduplication_key, + read_deduplication_fingerprint, + read_deduplication_fingerprint_values, skipped_user_suffixes, ) from state_db import connect_state_db @@ -47,6 +49,13 @@ CREATE TABLE IF NOT EXISTS audit_sources ( inode INTEGER NOT NULL, offset INTEGER NOT NULL ) WITHOUT ROWID; +CREATE TABLE IF NOT EXISTS audit_read_dedup ( + fingerprint BLOB PRIMARY KEY, + occurred_second INTEGER NOT NULL, + event_id INTEGER +) WITHOUT ROWID; +CREATE INDEX IF NOT EXISTS audit_read_dedup_time + ON audit_read_dedup (occurred_second); CREATE INDEX IF NOT EXISTS audit_events_main_time ON audit_events (occurred_second DESC, id DESC) WHERE share <> 'FSLogix' COLLATE NOCASE; @@ -138,7 +147,7 @@ CREATE TABLE IF NOT EXISTS audit_rollup_state ( backfill_next_id INTEGER NOT NULL, backfill_max_id INTEGER NOT NULL, ready INTEGER NOT NULL CHECK (ready IN (0, 1)), - policy_version INTEGER NOT NULL DEFAULT 3 + policy_version INTEGER NOT NULL DEFAULT 4 ); """ @@ -212,7 +221,7 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None: WHERE singleton = 1 """ ).fetchone() - if state is not None and int(state["policy_version"]) < 3: + if state is not None and int(state["policy_version"]) < 4: for table in ( "audit_daily_totals", "audit_daily_counts", @@ -223,7 +232,7 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None: """ UPDATE audit_rollup_state SET backfill_next_id = 1, backfill_max_id = ?, ready = ?, - policy_version = 3 + policy_version = 4 WHERE singleton = 1 """, (max_id, int(max_id == 0)), @@ -231,6 +240,21 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None: conn.commit() +READ_DEDUP_DEFAULT_WINDOW_SECONDS = 2 * 86400 + + +def read_deduplication_cutoff() -> int: + try: + window = int(os.getenv( + "AUDIT_READ_DEDUP_WINDOW_SECONDS", + str(READ_DEDUP_DEFAULT_WINDOW_SECONDS), + )) + except ValueError: + window = READ_DEDUP_DEFAULT_WINDOW_SECONDS + now = int(dt.datetime.now(dt.timezone.utc).timestamp()) + return now - max(1, window) + + def parse_event_second(timestamp: object) -> int: parsed = dt.datetime.fromisoformat(str(timestamp).replace("Z", "+00:00")) if parsed.tzinfo is None: @@ -257,22 +281,13 @@ def row_to_event(row: sqlite3.Row) -> Dict[str, object]: } -def last_event(conn: sqlite3.Connection) -> Optional[Dict[str, object]]: - row = conn.execute( - """ - SELECT id, occurred_at, ingested_at, user, client_ip, client, share, - action, result, success, path, source - FROM audit_events - ORDER BY id DESC - LIMIT 1 - """ - ).fetchone() - return row_to_event(row) if row is not None else None - - class AuditStore: def __init__(self, database_path: Optional[str] = None): self.conn = connect_state_db(database_path) + self.conn.create_function( + "audit_read_fingerprint", 7, + read_deduplication_fingerprint_values, deterministic=True, + ) ensure_audit_schema(self.conn) def close(self) -> None: @@ -299,10 +314,13 @@ class AuditStore: totals = Counter() counts = Counter() facets = set() - previous_key = read_deduplication_key(last_event(self.conn) or {}) + seen_read_fingerprints = set() for event in events: - current_key = read_deduplication_key(event) - if current_key is not None and current_key == previous_key: + read_fingerprint = read_deduplication_fingerprint(event) + if ( + read_fingerprint is not None + and read_fingerprint in seen_read_fingerprints + ): continue occurred_second = parse_event_second(event["timestamp"]) user = str(event["user"]) @@ -311,7 +329,6 @@ class AuditStore: action = str(event["action"]) result = str(event["result"]) success = int(bool(event["success"])) - day = occurred_second // 86400 event_rows.append( ( str(event["timestamp"]), @@ -327,17 +344,62 @@ class AuditStore: success, str(event["path"]), str(event["source"]), + read_fingerprint, ) ) - totals[(day, account)] += 1 - counts[(day, account, user, share, action, result, success)] += 1 - for kind, value in (("user", user), ("share", share), ("action", action)): - if value: - facets.add((kind, day, account, value)) - previous_key = current_key + if read_fingerprint is not None: + seen_read_fingerprints.add(read_fingerprint) self.conn.execute("BEGIN IMMEDIATE") try: + fingerprints = [ + row[13] for row in event_rows if row[13] is not None + ] + existing_fingerprints = set() + for offset in range(0, len(fingerprints), 500): + chunk = fingerprints[offset : offset + 500] + placeholders = ",".join("?" for _ in chunk) + existing_fingerprints.update( + bytes(row[0]) + for row in self.conn.execute( + f""" + SELECT fingerprint + FROM audit_read_dedup + WHERE fingerprint IN ({placeholders}) + """, + chunk, + ) + ) + if existing_fingerprints: + event_rows = [ + row + for row in event_rows + if ( + row[13] is None + or row[13] not in existing_fingerprints + ) + ] + + for row in event_rows: + day = int(row[1]) // 86400 + account = str(row[4]) + user = str(row[3]) + share = str(row[7]) + action = str(row[8]) + result = str(row[9]) + success = int(row[10]) + totals[(day, account)] += 1 + counts[( + day, account, user, share, action, result, success + )] += 1 + for kind, value in ( + ("user", user), + ("share", share), + ("action", action), + ): + if value: + facets.add((kind, day, account, value)) + if event_rows: paths = sorted({str(row[11]) for row in event_rows}) self.conn.executemany( @@ -383,6 +445,24 @@ class AuditStore: """, (len(indexed_rows),), ) + last_inserted_id = int( + self.conn.execute( + "SELECT max(id) FROM audit_events" + ).fetchone()[0] + ) + first_inserted_id = last_inserted_id - len(event_rows) + 1 + self.conn.executemany( + """ + INSERT INTO audit_read_dedup ( + fingerprint, occurred_second, event_id + ) VALUES (?, ?, ?) + """, + ( + (row[13], int(row[1]), first_inserted_id + offset) + for offset, row in enumerate(event_rows) + if row[13] is not None + ), + ) self.conn.executemany( ROLLUP_TOTAL_SQL, ((*key, count) for key, count in totals.items()), @@ -393,6 +473,11 @@ class AuditStore: ) self.conn.executemany(ROLLUP_FACET_SQL, facets) + self.conn.execute( + "DELETE FROM audit_read_dedup WHERE occurred_second < ?", + (read_deduplication_cutoff(),), + ) + if source_updates: self.conn.executemany( """ @@ -422,7 +507,7 @@ class AuditStore: raise return len(event_rows) - def backfill_rollups(self, limit: int = 50000) -> int: + def backfill_rollups(self, limit: int = 5000) -> int: """Backfill one bounded legacy-event chunk without delaying live appends.""" state = self.conn.execute( """ @@ -459,6 +544,55 @@ class AuditStore: end_id = int(chunk[1]) self.conn.execute("BEGIN IMMEDIATE") try: + dedup_cutoff = read_deduplication_cutoff() + self.conn.execute( + "DELETE FROM audit_read_dedup WHERE occurred_second < ?", + (dedup_cutoff,), + ) + self.conn.execute( + """ + INSERT INTO audit_read_dedup ( + fingerprint, occurred_second, event_id + ) + SELECT audit_read_fingerprint( + e.occurred_at, e.user, e.client_ip, e.share, e.path, + e.success, e.result + ), e.occurred_second, e.id + FROM audit_events AS e + WHERE e.id BETWEEN ? AND ? + AND e.action = 'read' + AND e.occurred_second >= ? + ORDER BY e.id + ON CONFLICT (fingerprint) DO UPDATE SET + event_id = coalesce( + audit_read_dedup.event_id, excluded.event_id + ) + """, + (start_id, end_id, dedup_cutoff), + ) + duplicate_read_sql = """ + SELECT e.id + FROM audit_events AS e + JOIN audit_read_dedup AS d + ON d.fingerprint = audit_read_fingerprint( + e.occurred_at, e.user, e.client_ip, e.share, e.path, + e.success, e.result + ) + WHERE e.id BETWEEN ? AND ? + AND e.action = 'read' + AND e.occurred_second >= ? + AND (d.event_id IS NULL OR d.event_id <> e.id) + """ + self.conn.execute( + f"DELETE FROM audit_path_events " + f"WHERE event_id IN ({duplicate_read_sql})", + (start_id, end_id, dedup_cutoff), + ) + self.conn.execute( + f"DELETE FROM audit_events " + f"WHERE id IN ({duplicate_read_sql})", + (start_id, end_id, dedup_cutoff), + ) self.conn.execute( """ INSERT OR IGNORE INTO audit_paths (path) diff --git a/app/web/app.js b/app/web/app.js index 49dd6e6..c0595cb 100644 --- a/app/web/app.js +++ b/app/web/app.js @@ -341,6 +341,7 @@ async function renderBackup() { '' ) + '
'; const button = document.querySelector("#start-backup"); + const body = document.querySelector("#backup-body"); let logLines = []; let previousActive = false; let polling = false; @@ -349,9 +350,10 @@ async function renderBackup() { const currentLog = document.querySelector("#backup-log"); const previousScrollTop = currentLog?.scrollTop || 0; const data = await api(includeLog ? "/api/backup" : "/api/backup?log=0"); + if (!body.isConnected || !button.isConnected) return null; if (Array.isArray(data.log)) logLines = data.log; data.log = logLines; - document.querySelector("#backup-body").innerHTML = backupMarkup(data); + body.innerHTML = backupMarkup(data); button.disabled = !data.manualEnabled || processIsActive(data); button.title = data.manualEnabled ? "" : "Es ist kein Sicherungsziel eingerichtet."; const logView = document.querySelector("#backup-log"); @@ -394,6 +396,7 @@ async function renderBackup() { } }); const initial = await refresh(true); + if (!body.isConnected || !button.isConnected) return; previousActive = processIsActive(initial); state.timer = window.setInterval(() => { void poll(); }, 2000); } @@ -449,6 +452,7 @@ async function renderReconciliation() { '' ) + '
'; const button = document.querySelector("#start-reconciliation"); + const body = document.querySelector("#reconciliation-body"); let logLines = []; let previousActive = false; let polling = false; @@ -459,9 +463,10 @@ async function renderReconciliation() { const data = await api( includeLog ? "/api/reconciliation" : "/api/reconciliation?log=0" ); + if (!body.isConnected || !button.isConnected) return null; if (Array.isArray(data.log)) logLines = data.log; data.log = logLines; - document.querySelector("#reconciliation-body").innerHTML = reconciliationMarkup(data); + body.innerHTML = reconciliationMarkup(data); button.disabled = processIsActive(data); const logView = document.querySelector("#reconciliation-log"); if (logView) { @@ -503,6 +508,7 @@ async function renderReconciliation() { } }); const initial = await refresh(true); + if (!body.isConnected || !button.isConnected) return; previousActive = processIsActive(initial); state.timer = window.setInterval(() => { void poll(); }, 1500); } diff --git a/tests/test_web_ui.py b/tests/test_web_ui.py index 94f8d67..fb9d696 100644 --- a/tests/test_web_ui.py +++ b/tests/test_web_ui.py @@ -466,8 +466,9 @@ class AuditParsingTests(unittest.TestCase): with tempfile.TemporaryDirectory() as tmpdir: active = os.path.join(tmpdir, "log.pc01") database = os.path.join(tmpdir, "state.db") + today = dt.datetime.now(dt.timezone.utc).strftime("%Y/%m/%d") read = ( - "[2026/07/31 12:34:56.000000, 1] smbd_audit: " + f"[{today} 12:34:56.000000, 1] smbd_audit: " "x|alice|192.0.2.5|PC01|Data|pread|OK|a.txt\n" ) write = read.replace("pread|OK|a.txt", "pwrite|OK|changed.txt") @@ -499,65 +500,87 @@ class AuditParsingTests(unittest.TestCase): finally: store.close() - def test_deduplicates_only_uninterrupted_identical_reads_in_one_second(self): + + def test_deduplicates_reads_per_second_across_interleaved_events(self): with tempfile.TemporaryDirectory() as tmpdir: - active = os.path.join(tmpdir, "log.pc01") database = os.path.join(tmpdir, "state.db") - - def line(operation, path, second="12:34:56"): - return ( - f"[2026/07/31 {second}.000000, 1] smbd_audit: " - f"x|alice|192.0.2.5|PC01|Data|{operation}|OK|{path}\n" - ) - - with open(active, "w", encoding="utf-8") as handle: - handle.writelines( - [ - line("pread", "a.txt"), - line("pread", "a.txt"), - line("pread", "b.txt"), - line("pread", "a.txt"), - line("pread", "a.txt"), - ] - ) - store = audit_store.AuditStore(database) - try: - with mock.patch.object( - audit_collector, - "SAMBA_LOG_GLOB", - os.path.join(tmpdir, "log.*"), - ): - self.assertEqual(audit_collector.collect_once(store), 3) - with open(active, "a", encoding="utf-8") as handle: - handle.write(line("pread", "a.txt")) - self.assertEqual(audit_collector.collect_once(store), 0) - with open(active, "a", encoding="utf-8") as handle: - handle.write(line("pwrite", "changed.txt")) - handle.write(line("pread", "a.txt")) - self.assertEqual(audit_collector.collect_once(store), 2) - with open(active, "a", encoding="utf-8") as handle: - handle.write(line("pread", "a.txt", "12:34:57")) - self.assertEqual(audit_collector.collect_once(store), 1) - rows = store.conn.execute( - "SELECT action, path, occurred_at FROM audit_events ORDER BY id" - ).fetchall() - self.assertEqual( - [(row["action"], row["path"]) for row in rows], + def event( + action="read", + path="a.txt", + second="12:34:56", + share="Data", + source="log.pc01", + ): + timestamp = f"{dt.datetime.now(dt.timezone.utc).date()}T{second}+00:00" + return { + "timestamp": timestamp, + "ingestedAt": timestamp, + "user": "LOCAL\\silvia.mueller", + "clientIp": "10.100.0.22", + "client": "PC01", + "share": share, + "operation": "pread" if action == "read" else "pwrite", + "action": action, + "path": path, + "result": "OK", + "success": True, + "source": source, + } + + try: + inserted = store.append_batch( [ - ("read", "a.txt"), - ("read", "b.txt"), - ("read", "a.txt"), - ("write", "changed.txt"), - ("read", "a.txt"), - ("read", "a.txt"), + event(source="log.pc01"), + event(path="profile.vhdx", share="FSLogix"), + event(action="write", path="changed.txt"), + event(action="write", path="changed.txt", source="log.pc02"), + event(source="log.pc02"), + event(path="b.txt"), + event(source="log.pc03"), ], + {}, + set(), ) - self.assertTrue(rows[-1]["occurred_at"].endswith("12:34:57+00:00")) + repeated_poll = store.append_batch( + [event(source="log.pc04")], {}, set() + ) + next_second = store.append_batch( + [event(second="12:34:57")], {}, set() + ) + rows = store.conn.execute( + """ + SELECT action, share, path, occurred_second + FROM audit_events + ORDER BY id + """ + ).fetchall() + fingerprints = store.conn.execute( + """ + SELECT count(*) FROM audit_read_dedup + """ + ).fetchone()[0] finally: store.close() + self.assertEqual(inserted, 5) + self.assertEqual(repeated_poll, 0) + self.assertEqual(next_second, 1) + self.assertEqual(len(rows), 6) + self.assertEqual(fingerprints, 4) + self.assertEqual( + [(row["action"], row["share"], row["path"]) for row in rows], + [ + ("read", "Data", "a.txt"), + ("read", "FSLogix", "profile.vhdx"), + ("write", "Data", "changed.txt"), + ("write", "Data", "changed.txt"), + ("read", "Data", "b.txt"), + ("read", "Data", "a.txt"), + ], + ) + class AuditQueryTests(unittest.TestCase): def make_event(self, timestamp, user, success=True): @@ -946,7 +969,69 @@ class AuditQueryTests(unittest.TestCase): self.assertEqual(main["matched"], 1) self.assertEqual(fslogix["matched"], 1) self.assertEqual(total, 2) - self.assertEqual(policy_version, 3) + self.assertEqual(policy_version, 4) + + def test_upgrade_collapses_recent_legacy_read_duplicates(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) + event = self.make_event(f"{today}T12:00:00+00:00", "alice") + store.append_batch([event], {}, set()) + for source in ("log.pc02", "log.pc03"): + store.conn.execute( + """ + INSERT INTO audit_events ( + occurred_at, occurred_second, ingested_at, user, account, + client_ip, client, share, action, result, success, path, + path_id, source + ) + SELECT occurred_at, occurred_second, ingested_at, user, + account, client_ip, client, share, action, result, + success, path, path_id, ? + FROM audit_events + WHERE id = 1 + """, + (source,), + ) + store.conn.execute("DELETE FROM audit_read_dedup") + store.conn.execute( + "UPDATE audit_rollup_state SET policy_version = 3" + ) + store.conn.commit() + store.close() + + store = audit_store.AuditStore(database) + try: + live = dict(event, source="log.live") + live_inserted = store.append_batch([live], {}, set()) + while store.backfill_rollups(limit=1): + pass + event_count = store.conn.execute( + "SELECT count(*) FROM audit_events" + ).fetchone()[0] + rolled_up = store.conn.execute( + "SELECT sum(event_count) FROM audit_daily_totals" + ).fetchone()[0] + dedup_keys = store.conn.execute( + "SELECT count(*) FROM audit_read_dedup" + ).fetchone()[0] + policy_version = store.conn.execute( + "SELECT policy_version FROM audit_rollup_state" + ).fetchone()[0] + retained_source = store.conn.execute( + "SELECT source FROM audit_events" + ).fetchone()[0] + finally: + store.close() + + self.assertEqual(event_count, 1) + self.assertEqual(rolled_up, 1) + self.assertEqual(live_inserted, 1) + self.assertEqual(retained_source, "log.live") + self.assertEqual(dedup_keys, 1) + self.assertEqual(policy_version, 4) + def test_schema_has_filter_indexes_and_rollup_tables(self): with tempfile.TemporaryDirectory() as tmpdir: @@ -981,6 +1066,7 @@ class AuditQueryTests(unittest.TestCase): "audit_events_fslogix_user_time", "audit_events_fslogix_account_time", "audit_events_fslogix_result_time", + "audit_read_dedup_time", }.issubset(indexes) ) self.assertTrue( @@ -992,6 +1078,7 @@ class AuditQueryTests(unittest.TestCase): "audit_paths", "audit_paths_fts", "audit_path_events", + "audit_read_dedup", }.issubset(tables) ) @@ -1082,6 +1169,9 @@ class WebPresentationTests(unittest.TestCase): self.assertIn('"/api/backup?log=0"', script) self.assertIn('"/api/reconciliation?log=0"', script) self.assertEqual(script.count("if (active || previousActive)"), 2) + self.assertIn('const body = document.querySelector("#reconciliation-body")', script) + self.assertIn("if (!body.isConnected || !button.isConnected)", script) + self.assertNotIn('document.querySelector("#reconciliation-body").innerHTML', script) self.assertNotIn("BACKUP_AUTO_ENABLED", script) self.assertIn("data.interrupted", script)