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 => `