From 14874e504e4e1f6ac31e8f9bfa6e21675c1a4076 Mon Sep 17 00:00:00 2001 From: Ludwig Lehnert Date: Tue, 11 Aug 2026 15:20:53 +0000 Subject: [PATCH] way more efficient logging --- README.md | 12 +- app/audit_collector.py | 1 + app/audit_store.py | 767 ++++++++++++++++++++++++++++++++++++----- app/web/app.js | 19 +- app/web_ui.py | 2 +- dev/e2e.py | 20 +- tests/test_web_ui.py | 165 ++++++++- 7 files changed, 886 insertions(+), 100 deletions(-) diff --git a/README.md b/README.md index 6c99c0c..d9508b1 100644 --- a/README.md +++ b/README.md @@ -54,9 +54,12 @@ 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_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; - `web_cache`: cached storage scan results. -Activity lookup uses UTC epoch seconds, keyset pagination, and multi-column indexes for time, user, share, action, and result. WAL mode allows the collector, reconciler, scanner, and read-only web requests to operate concurrently. +Activity pages use UTC epoch seconds, keyset pagination, and compact indexes for time, scalar filters, and selective substring path searches. Exact scalar-filter counts and facets read the daily rollups instead of scanning raw events. WAL mode allows the collector, reconciler, scanner, and read-only web requests to operate concurrently. Standard SQLite does not provide transparent general-purpose compression, so this project deliberately does not depend on a non-core compression VFS. Structured columns avoid repeated JSON field names and make indexed queries much cheaper; the state volume still needs capacity for retained activity. @@ -396,10 +399,13 @@ Samba emits only the selected high-level `full_audit` operations, and the collec - 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; - 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 queries run directly against indexed SQLite columns and use a stable time/id cursor; +- 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; +- 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. On the first collector start after this migration, recognized legacy 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, 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. ## Backups diff --git a/app/audit_collector.py b/app/audit_collector.py index a04a555..7c9f521 100644 --- a/app/audit_collector.py +++ b/app/audit_collector.py @@ -165,6 +165,7 @@ def main() -> int: count = collect_once(store) if count: log(f"Stored {count} event(s)") + store.backfill_rollups() except Exception as exc: # pylint: disable=broad-except log(f"Collector cycle failed: {exc}") time.sleep(POLL_SECONDS) diff --git a/app/audit_store.py b/app/audit_store.py index 31d0cb7..ff085fe 100644 --- a/app/audit_store.py +++ b/app/audit_store.py @@ -5,6 +5,7 @@ import datetime as dt import glob import os import sqlite3 +from collections import Counter from typing import Dict, Iterable, List, Optional, Set, Tuple try: @@ -34,6 +35,7 @@ CREATE TABLE IF NOT EXISTS audit_events ( result TEXT NOT NULL, success INTEGER NOT NULL CHECK (success IN (0, 1)), path TEXT NOT NULL, + path_id INTEGER, source TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS audit_sources ( @@ -55,11 +57,102 @@ 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 TABLE IF NOT EXISTS audit_paths ( + id INTEGER PRIMARY KEY, + path TEXT NOT NULL UNIQUE +); +CREATE VIRTUAL TABLE IF NOT EXISTS audit_paths_fts USING fts5( + path, + content = 'audit_paths', + content_rowid = 'id', + tokenize = 'trigram' +); +CREATE TRIGGER IF NOT EXISTS audit_paths_insert +AFTER INSERT ON audit_paths BEGIN + INSERT INTO audit_paths_fts (rowid, path) VALUES (new.id, new.path); +END; +CREATE TABLE IF NOT EXISTS audit_path_events ( + path_id INTEGER NOT NULL, + event_second INTEGER NOT NULL, + event_id INTEGER NOT NULL, + PRIMARY KEY (path_id, event_second DESC, event_id DESC) +) WITHOUT ROWID; + + +CREATE TABLE IF NOT EXISTS audit_daily_totals ( + day INTEGER NOT NULL, + account TEXT NOT NULL, + event_count INTEGER NOT NULL, + PRIMARY KEY (day, account) +) WITHOUT ROWID; +CREATE TABLE IF NOT EXISTS audit_daily_counts ( + day INTEGER NOT NULL, + account TEXT NOT NULL, + user TEXT COLLATE NOCASE NOT NULL, + share TEXT COLLATE NOCASE NOT NULL, + action TEXT NOT NULL, + result TEXT COLLATE NOCASE NOT NULL, + success INTEGER NOT NULL CHECK (success IN (0, 1)), + event_count INTEGER NOT NULL, + PRIMARY KEY (day, account, user, share, action, result, success) +) WITHOUT ROWID; +CREATE TABLE IF NOT EXISTS audit_daily_facets ( + kind TEXT NOT NULL, + day INTEGER NOT NULL, + account TEXT NOT NULL, + value TEXT COLLATE NOCASE NOT NULL, + PRIMARY KEY (kind, day, account, value) +) WITHOUT ROWID; +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)) +); + +""" + +ROLLUP_TOTAL_SQL = """ +INSERT INTO audit_daily_totals (day, account, event_count) +VALUES (?, ?, ?) +ON CONFLICT (day, account) DO UPDATE SET + event_count = event_count + excluded.event_count +""" + +ROLLUP_COUNT_SQL = """ +INSERT INTO audit_daily_counts ( + day, account, user, share, action, result, success, event_count +) VALUES (?, ?, ?, ?, ?, ?, ?, ?) +ON CONFLICT (day, account, user, share, action, result, success) DO UPDATE SET + event_count = event_count + excluded.event_count +""" + +ROLLUP_FACET_SQL = """ +INSERT OR IGNORE INTO audit_daily_facets (kind, day, account, value) +VALUES (?, ?, ?, ?) """ def ensure_audit_schema(conn: sqlite3.Connection) -> None: conn.executescript(AUDIT_SCHEMA) + columns = { + str(row["name"]) + for row in conn.execute("PRAGMA table_info(audit_events)") + } + if "path_id" not in columns: + conn.execute("ALTER TABLE audit_events ADD COLUMN path_id INTEGER") + + max_id = int( + conn.execute("SELECT coalesce(max(id), 0) FROM audit_events").fetchone()[0] + ) + conn.execute( + """ + INSERT OR IGNORE INTO audit_rollup_state ( + singleton, backfill_next_id, backfill_max_id, ready + ) VALUES (1, 1, ?, ?) + """, + (max_id, int(max_id == 0)), + ) conn.commit() @@ -127,43 +220,106 @@ class AuditStore: source_updates: Dict[str, Dict[str, object]], seen_paths: Set[str], ) -> int: - inserted = 0 + event_rows = [] + totals = Counter() + counts = Counter() + facets = set() + previous_key = read_deduplication_key(last_event(self.conn) or {}) + for event in events: + current_key = read_deduplication_key(event) + if current_key is not None and current_key == previous_key: + continue + occurred_second = parse_event_second(event["timestamp"]) + user = str(event["user"]) + account = account_name(user).casefold() + share = str(event["share"]) + action = str(event["action"]) + result = str(event["result"]) + success = int(bool(event["success"])) + day = occurred_second // 86400 + event_rows.append( + ( + str(event["timestamp"]), + occurred_second, + str(event["ingestedAt"]), + user, + account, + str(event["clientIp"]), + str(event["client"]), + share, + action, + result, + success, + str(event["path"]), + str(event["source"]), + ) + ) + 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 + self.conn.execute("BEGIN IMMEDIATE") try: - previous_key = read_deduplication_key(last_event(self.conn) or {}) - for event in events: - current_key = read_deduplication_key(event) - if current_key is not None and current_key == previous_key: - continue - self.conn.execute( + if event_rows: + paths = sorted({str(row[11]) for row in event_rows}) + self.conn.executemany( + "INSERT OR IGNORE INTO audit_paths (path) VALUES (?)", + ((path,) for path in paths), + ) + path_ids = {} + for offset in range(0, len(paths), 500): + path_chunk = paths[offset : offset + 500] + placeholders = ",".join("?" for _ in path_chunk) + for row in self.conn.execute( + f""" + SELECT id, path + FROM audit_paths + WHERE path IN ({placeholders}) + """, + path_chunk, + ): + path_ids[str(row["path"])] = int(row["id"]) + indexed_rows = [ + (*row[:12], path_ids[str(row[11])], row[12]) + for row in event_rows + ] + self.conn.executemany( """ INSERT INTO audit_events ( occurred_at, occurred_second, ingested_at, user, account, client_ip, client, share, action, result, success, path, - source - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + path_id, source + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, - ( - str(event["timestamp"]), - parse_event_second(event["timestamp"]), - str(event["ingestedAt"]), - str(event["user"]), - account_name(str(event["user"])).casefold(), - str(event["clientIp"]), - str(event["client"]), - str(event["share"]), - str(event["action"]), - str(event["result"]), - int(bool(event["success"])), - str(event["path"]), - str(event["source"]), - ), + indexed_rows, ) - inserted += 1 - previous_key = current_key - - for path, entry in source_updates.items(): self.conn.execute( + """ + INSERT OR IGNORE INTO audit_path_events ( + path_id, event_second, event_id + ) + SELECT path_id, occurred_second, id + FROM audit_events + ORDER BY id DESC + LIMIT ? + """, + (len(indexed_rows),), + ) + self.conn.executemany( + ROLLUP_TOTAL_SQL, + ((*key, count) for key, count in totals.items()), + ) + self.conn.executemany( + ROLLUP_COUNT_SQL, + ((*key, count) for key, count in counts.items()), + ) + self.conn.executemany(ROLLUP_FACET_SQL, facets) + + if source_updates: + self.conn.executemany( """ INSERT INTO audit_sources (path, inode, offset) VALUES (?, ?, ?) @@ -171,17 +327,158 @@ class AuditStore: inode = excluded.inode, offset = excluded.offset """, - (path, int(entry["inode"]), int(entry["offset"])), + ( + (path, int(entry["inode"]), int(entry["offset"])) + for path, entry in source_updates.items() + ), + ) + missing_paths = [ + (str(row["path"]),) + for row in self.conn.execute("SELECT path FROM audit_sources") + if str(row["path"]) not in seen_paths + ] + if missing_paths: + self.conn.executemany( + "DELETE FROM audit_sources WHERE path = ?", missing_paths ) - for row in self.conn.execute("SELECT path FROM audit_sources").fetchall(): - path = str(row["path"]) - if path not in seen_paths: - self.conn.execute("DELETE FROM audit_sources WHERE path = ?", (path,)) self.conn.commit() except Exception: self.conn.rollback() raise - return inserted + return len(event_rows) + + def backfill_rollups(self, limit: int = 10000) -> int: + """Backfill one bounded legacy-event chunk without delaying live appends.""" + state = self.conn.execute( + """ + SELECT backfill_next_id, backfill_max_id, ready + FROM audit_rollup_state + WHERE singleton = 1 + """ + ).fetchone() + if state is None or bool(state["ready"]): + return 0 + start_id = int(state["backfill_next_id"]) + max_id = int(state["backfill_max_id"]) + chunk = self.conn.execute( + """ + SELECT count(*), max(id) + FROM ( + SELECT id + FROM audit_events + WHERE id >= ? AND id <= ? + ORDER BY id + LIMIT ? + ) + """, + (start_id, max_id, max(1, limit)), + ).fetchone() + row_count = int(chunk[0]) + if row_count == 0: + self.conn.execute( + "UPDATE audit_rollup_state SET ready = 1 WHERE singleton = 1" + ) + self.conn.commit() + return 0 + + end_id = int(chunk[1]) + self.conn.execute("BEGIN IMMEDIATE") + try: + self.conn.execute( + """ + INSERT OR IGNORE INTO audit_paths (path) + SELECT DISTINCT path + FROM audit_events + WHERE id BETWEEN ? AND ? + """, + (start_id, end_id), + ) + self.conn.execute( + """ + UPDATE audit_events + SET path_id = ( + SELECT id + FROM audit_paths + WHERE audit_paths.path = audit_events.path + ) + WHERE id BETWEEN ? AND ? + """, + (start_id, end_id), + ) + self.conn.execute( + """ + INSERT OR IGNORE INTO audit_path_events ( + path_id, event_second, event_id + ) + SELECT path_id, occurred_second, id + FROM audit_events + WHERE id BETWEEN ? AND ? + """, + (start_id, end_id), + ) + + self.conn.execute( + """ + INSERT INTO audit_daily_totals (day, account, event_count) + SELECT occurred_second / 86400, account, count(*) + FROM audit_events + WHERE id BETWEEN ? AND ? + GROUP BY occurred_second / 86400, account + ON CONFLICT (day, account) DO UPDATE SET + event_count = event_count + excluded.event_count + """, + (start_id, end_id), + ) + self.conn.execute( + """ + INSERT INTO audit_daily_counts ( + day, account, user, share, action, result, success, event_count + ) + SELECT occurred_second / 86400, account, user, share, action, + result, success, count(*) + FROM audit_events + WHERE id BETWEEN ? AND ? + GROUP BY occurred_second / 86400, account, user, share, action, + result, success + ON CONFLICT ( + day, account, user, share, action, result, success + ) DO UPDATE SET + event_count = event_count + excluded.event_count + """, + (start_id, end_id), + ) + for kind, column in ( + ("user", "user"), + ("share", "share"), + ("action", "action"), + ): + self.conn.execute( + f""" + INSERT OR IGNORE INTO audit_daily_facets ( + kind, day, account, value + ) + SELECT ?, occurred_second / 86400, account, {column} + FROM audit_events + WHERE id BETWEEN ? AND ? AND {column} <> '' + GROUP BY occurred_second / 86400, account, {column} + """, + (kind, start_id, end_id), + ) + finished = end_id >= max_id + self.conn.execute( + """ + UPDATE audit_rollup_state + SET backfill_next_id = ?, ready = ? + WHERE singleton = 1 + """, + (end_id + 1, int(finished)), + ) + self.conn.commit() + except Exception: + self.conn.rollback() + raise + return row_count + def escape_like(value: str) -> str: @@ -194,29 +491,55 @@ def date_seconds(day: dt.date) -> int: ) +def activity_filters(params: Dict[str, List[str]]) -> Dict[str, str]: + return { + key: params.get(key, [""])[0].casefold().strip() + for key in ("user", "share", "operation", "action", "path", "result") + } + + +def query_flag( + params: Dict[str, List[str]], name: str, *, default: bool = True +) -> bool: + fallback = "1" if default else "0" + return params.get(name, [fallback])[0].strip().casefold() not in { + "0", + "false", + "no", + "off", + } + + +def excluded_account_conditions() -> Tuple[List[str], List[object]]: + conditions: List[str] = [] + values: List[object] = [] + for suffix in skipped_user_suffixes(): + conditions.append("account NOT LIKE ? ESCAPE '\\'") + values.append(f"%{escape_like(suffix)}") + return conditions, values + + def activity_conditions( start: dt.date, end: dt.date, params: Dict[str, List[str]], *, include_filters: bool, + include_path: bool = True, ) -> Tuple[List[str], List[object]]: conditions = ["occurred_second >= ?", "occurred_second < ?"] values: List[object] = [ date_seconds(start), date_seconds(end + dt.timedelta(days=1)), ] - for suffix in skipped_user_suffixes(): - conditions.append("lower(account) NOT LIKE ? ESCAPE '\\'") - values.append(f"%{escape_like(suffix)}") + excluded_conditions, excluded_values = excluded_account_conditions() + conditions.extend(excluded_conditions) + values.extend(excluded_values) if not include_filters: return conditions, values - filters = { - key: params.get(key, [""])[0].casefold().strip() - for key in ("user", "share", "operation", "action", "path", "result") - } + filters = activity_filters(params) user = filters["user"] if user: if "\\" in user or "@" in user: @@ -229,7 +552,7 @@ def activity_conditions( conditions.append("share = ? COLLATE NOCASE") values.append(share) path = filters["path"] - if path: + if path and include_path: conditions.append("lower(path) LIKE ? ESCAPE '\\'") values.append(f"%{escape_like(path)}%") action = filters["action"] or filters["operation"] @@ -245,6 +568,197 @@ def activity_conditions( return conditions, values +def rollups_ready(conn: sqlite3.Connection) -> bool: + try: + row = conn.execute( + "SELECT ready FROM audit_rollup_state WHERE singleton = 1" + ).fetchone() + except sqlite3.Error: + return False + return row is not None and bool(row[0]) + + +def rollup_range( + start: dt.date, end: dt.date +) -> Tuple[List[str], List[object]]: + conditions = ["day >= ?", "day <= ?"] + values: List[object] = [ + date_seconds(start) // 86400, + date_seconds(end) // 86400, + ] + excluded_conditions, excluded_values = excluded_account_conditions() + conditions.extend(excluded_conditions) + values.extend(excluded_values) + return conditions, values + + +def rollup_matched( + conn: sqlite3.Connection, + start: dt.date, + end: dt.date, + params: Dict[str, List[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) + + row = conn.execute( + f""" + SELECT coalesce(sum(event_count), 0) + FROM {table} + WHERE {" AND ".join(conditions)} + """, + values, + ).fetchone() + return int(row[0]) + + +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]], +) -> Dict[str, List[str]]: + conditions, values = activity_conditions( + start, end, params, include_filters=False + ) + where_sql = " AND ".join(conditions) + + def distinct(column: str) -> List[str]: + return [ + str(row[0]) + for row in conn.execute( + f""" + SELECT DISTINCT {column} + FROM audit_events + WHERE {where_sql} AND {column} <> '' + ORDER BY {column} COLLATE NOCASE + """, + values, + ) + ] + + actions = distinct("action") + return { + "users": distinct("user"), + "shares": distinct("share"), + "operations": actions, + "actions": actions, + } + + +def indexed_path_ids( + conn: sqlite3.Connection, path: str, *, maximum: int = 500 +) -> Optional[List[int]]: + if len(path) < 3: + return None + phrase = '"' + path.replace('"', '""') + '"' + try: + rows = conn.execute( + """ + SELECT rowid + FROM audit_paths_fts + WHERE audit_paths_fts MATCH ? + LIMIT ? + """, + (phrase, maximum + 1), + ).fetchall() + except sqlite3.Error: + return None + if len(rows) > maximum: + return None + return [int(row[0]) for row in rows] + + + +def path_event_set_is_selective( + conn: sqlite3.Connection, + path_ids: List[int], + start: dt.date, + end: dt.date, + *, + maximum_events: int = 5000, +) -> bool: + placeholders = ",".join("?" for _ in path_ids) + row = conn.execute( + f""" + SELECT count(*) + FROM ( + SELECT 1 + FROM audit_path_events + WHERE path_id IN ({placeholders}) + AND event_second >= ? AND event_second < ? + LIMIT ? + ) + """, + ( + *path_ids, + date_seconds(start), + date_seconds(end + dt.timedelta(days=1)), + maximum_events + 1, + ), + ).fetchone() + return int(row[0]) <= maximum_events + + def query_activity( conn: sqlite3.Connection, start: dt.date, @@ -252,35 +766,81 @@ def query_activity( params: Dict[str, List[str]], ) -> Dict[str, object]: limit = min(500, max(1, int(params.get("limit", ["100"])[0]))) - conditions, values = activity_conditions(start, end, params, include_filters=True) - where_sql = " AND ".join(conditions) - matched = int( - conn.execute( - f"SELECT count(*) FROM audit_events WHERE {where_sql}", - values, - ).fetchone()[0] + ready = rollups_ready(conn) + path = activity_filters(params)["path"] + path_ids = indexed_path_ids(conn, path) if ready and path else None + if ( + path_ids + and len(path_ids) > 1 + and not path_event_set_is_selective(conn, path_ids, start, end) + ): + path_ids = None + conditions, values = activity_conditions( + start, + end, + params, + include_filters=True, + include_path=path_ids is None, ) + single_path_id = None + if path_ids is not None: + if len(path_ids) == 1: + single_path_id = path_ids[0] + elif path_ids: + path_placeholders = ",".join("?" for _ in path_ids) + conditions.append( + f"id IN (SELECT event_id FROM audit_path_events " + f"WHERE path_id IN ({path_placeholders}) " + "AND event_second >= ? AND event_second < ?)" + ) + values.extend(path_ids) + values.extend( + (date_seconds(start), date_seconds(end + dt.timedelta(days=1))) + ) + else: + conditions.append("0") cursor = params.get("cursor", [""])[0].strip() page_conditions = list(conditions) page_values = list(values) + from_sql = "audit_events" + order_sql = "occurred_second DESC, id DESC" + if single_path_id is not None: + from_sql = ( + "audit_path_events AS pe " + "CROSS JOIN audit_events ON id = pe.event_id" + ) + order_sql = "pe.event_second DESC, pe.event_id DESC" + page_conditions.extend( + ["pe.path_id = ?", "pe.event_second >= ?", "pe.event_second < ?"] + ) + page_values.extend( + ( + single_path_id, + date_seconds(start), + date_seconds(end + dt.timedelta(days=1)), + ) + ) if cursor: try: cursor_second, cursor_id = (int(part) for part in cursor.split(":", 1)) except (TypeError, ValueError) as exc: raise ValueError("Ungültiger Seitenzeiger") from exc - page_conditions.append( - "(occurred_second < ? OR (occurred_second = ? AND id < ?))" + cursor_condition = ( + "(pe.event_second < ? OR " + "(pe.event_second = ? AND pe.event_id < ?))" + if single_path_id is not None + else "(occurred_second < ? OR (occurred_second = ? AND id < ?))" ) + page_conditions.append(cursor_condition) page_values.extend((cursor_second, cursor_second, cursor_id)) - rows = conn.execute( f""" SELECT id, occurred_at, occurred_second, ingested_at, user, client_ip, client, share, action, result, success, path, source - FROM audit_events + FROM {from_sql} WHERE {" AND ".join(page_conditions)} - ORDER BY occurred_second DESC, id DESC + ORDER BY {order_sql} LIMIT ? """, (*page_values, limit + 1), @@ -292,48 +852,77 @@ def query_activity( if has_more and rows: next_cursor = f"{int(rows[-1]['occurred_second'])}:{int(rows[-1]['id'])}" - facet_conditions, facet_values = activity_conditions( - start, end, params, include_filters=False - ) - facet_where = " AND ".join(facet_conditions) - - def distinct(column: str) -> List[str]: - return [ - str(row[0]) - for row in conn.execute( - f""" - SELECT DISTINCT {column} - FROM audit_events - WHERE {facet_where} AND {column} <> '' - ORDER BY {column} COLLATE NOCASE - """, - facet_values, + 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) + + 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": []} - actions = distinct("action") return { "events": events, + "hasMore": has_more, "nextCursor": next_cursor, "matched": matched, - "facets": { - "users": distinct("user"), - "shares": distinct("share"), - "operations": actions, - "actions": actions, - }, + "matchedExact": matched is not None, + "facets": facets, } def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, object]: - 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 - """ - ).fetchone() + if rollups_ready(conn): + row = conn.execute( + """ + SELECT count(DISTINCT day), min(day), max(day) + FROM audit_daily_totals + """ + ).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]) + else: + 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 + """ + ).fetchone() + days = int(row[0] or 0) + oldest = row[1] + newest = row[2] + database_bytes = 0 for suffix in ("", "-wal"): try: @@ -341,10 +930,10 @@ def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, obj except OSError: pass return { - "days": int(row[0] or 0), + "days": days, "bytes": database_bytes, - "oldest": row[1], - "newest": row[2], + "oldest": oldest, + "newest": newest, } diff --git a/app/web/app.js b/app/web/app.js index bd87218..11e0586 100644 --- a/app/web/app.js +++ b/app/web/app.js @@ -260,20 +260,31 @@ async function renderActivity() { `; const form = document.querySelector("#activity-filter"); let cursor = 0; + let loaded = 0; + let matched = null; const load = async append => { const params = new URLSearchParams(new FormData(form)); params.set("limit", "100"); if (cursor) params.set("cursor", String(cursor)); + if (append) { + params.set("count", "0"); + params.set("facets", "0"); + } const result = await api(`/api/activity?${params}`); state.activity = result; const rows = document.querySelector("#activity-rows"); if (append) rows.insertAdjacentHTML("beforeend", eventRows(result.events)); else rows.innerHTML = eventRows(result.events); - document.querySelector("#activity-summary").textContent = `${result.matched.toLocaleString("de-DE")} passende Ereignisse`; + 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 more = document.querySelector("#load-more"); cursor = result.nextCursor || 0; - more.hidden = !result.nextCursor; - for (const [id, values] of [["users-list", result.facets.users], ["shares-list", result.facets.shares]]) { - document.querySelector(`#${id}`).innerHTML = values.filter(Boolean).map(value => `