way more efficient logging

This commit is contained in:
Ludwig Lehnert
2026-08-11 15:20:53 +00:00
parent 79cd02695a
commit 14874e504e
7 changed files with 886 additions and 100 deletions
+9 -3
View File
@@ -54,9 +54,12 @@ The database contains:
- `shares`: AD group-to-folder lifecycle and ACL reconciliation state; - `shares`: AD group-to-folder lifecycle and ACL reconciliation state;
- `audit_events`: normalized read, write, move, and delete events; - `audit_events`: normalized read, write, move, and delete events;
- `audit_sources`: Samba log inode/offset checkpoints; - `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. - `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. 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`; - 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; - 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; - 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. - 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 ## Backups
+1
View File
@@ -165,6 +165,7 @@ def main() -> int:
count = collect_once(store) count = collect_once(store)
if count: if count:
log(f"Stored {count} event(s)") log(f"Stored {count} event(s)")
store.backfill_rollups()
except Exception as exc: # pylint: disable=broad-except except Exception as exc: # pylint: disable=broad-except
log(f"Collector cycle failed: {exc}") log(f"Collector cycle failed: {exc}")
time.sleep(POLL_SECONDS) time.sleep(POLL_SECONDS)
+664 -75
View File
@@ -5,6 +5,7 @@ import datetime as dt
import glob import glob
import os import os
import sqlite3 import sqlite3
from collections import Counter
from typing import Dict, Iterable, List, Optional, Set, Tuple from typing import Dict, Iterable, List, Optional, Set, Tuple
try: try:
@@ -34,6 +35,7 @@ CREATE TABLE IF NOT EXISTS audit_events (
result TEXT NOT NULL, result TEXT NOT NULL,
success INTEGER NOT NULL CHECK (success IN (0, 1)), success INTEGER NOT NULL CHECK (success IN (0, 1)),
path TEXT NOT NULL, path TEXT NOT NULL,
path_id INTEGER,
source TEXT NOT NULL source TEXT NOT NULL
); );
CREATE TABLE IF NOT EXISTS audit_sources ( 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); ON audit_events (share COLLATE NOCASE, occurred_second DESC, id DESC);
CREATE INDEX IF NOT EXISTS audit_events_result_time CREATE INDEX IF NOT EXISTS audit_events_result_time
ON audit_events (result COLLATE NOCASE, occurred_second DESC, id DESC); 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: def ensure_audit_schema(conn: sqlite3.Connection) -> None:
conn.executescript(AUDIT_SCHEMA) 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() conn.commit()
@@ -127,43 +220,106 @@ class AuditStore:
source_updates: Dict[str, Dict[str, object]], source_updates: Dict[str, Dict[str, object]],
seen_paths: Set[str], seen_paths: Set[str],
) -> int: ) -> int:
inserted = 0 event_rows = []
self.conn.execute("BEGIN IMMEDIATE") totals = Counter()
try: counts = Counter()
facets = set()
previous_key = read_deduplication_key(last_event(self.conn) or {}) previous_key = read_deduplication_key(last_event(self.conn) or {})
for event in events: for event in events:
current_key = read_deduplication_key(event) current_key = read_deduplication_key(event)
if current_key is not None and current_key == previous_key: if current_key is not None and current_key == previous_key:
continue continue
self.conn.execute( 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:
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 ( INSERT INTO audit_events (
occurred_at, occurred_second, ingested_at, user, account, occurred_at, occurred_second, ingested_at, user, account,
client_ip, client, share, action, result, success, path, client_ip, client, share, action, result, success, path,
source path_id, source
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""", """,
( indexed_rows,
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"]),
),
) )
inserted += 1
previous_key = current_key
for path, entry in source_updates.items():
self.conn.execute( 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) INSERT INTO audit_sources (path, inode, offset)
VALUES (?, ?, ?) VALUES (?, ?, ?)
@@ -171,17 +327,158 @@ class AuditStore:
inode = excluded.inode, inode = excluded.inode,
offset = excluded.offset 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() self.conn.commit()
except Exception: except Exception:
self.conn.rollback() self.conn.rollback()
raise 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: 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( def activity_conditions(
start: dt.date, start: dt.date,
end: dt.date, end: dt.date,
params: Dict[str, List[str]], params: Dict[str, List[str]],
*, *,
include_filters: bool, include_filters: bool,
include_path: bool = True,
) -> Tuple[List[str], List[object]]: ) -> Tuple[List[str], List[object]]:
conditions = ["occurred_second >= ?", "occurred_second < ?"] conditions = ["occurred_second >= ?", "occurred_second < ?"]
values: List[object] = [ values: List[object] = [
date_seconds(start), date_seconds(start),
date_seconds(end + dt.timedelta(days=1)), date_seconds(end + dt.timedelta(days=1)),
] ]
for suffix in skipped_user_suffixes(): excluded_conditions, excluded_values = excluded_account_conditions()
conditions.append("lower(account) NOT LIKE ? ESCAPE '\\'") conditions.extend(excluded_conditions)
values.append(f"%{escape_like(suffix)}") values.extend(excluded_values)
if not include_filters: if not include_filters:
return conditions, values return conditions, values
filters = { filters = activity_filters(params)
key: params.get(key, [""])[0].casefold().strip()
for key in ("user", "share", "operation", "action", "path", "result")
}
user = filters["user"] user = filters["user"]
if user: if user:
if "\\" in user or "@" in user: if "\\" in user or "@" in user:
@@ -229,7 +552,7 @@ def activity_conditions(
conditions.append("share = ? COLLATE NOCASE") conditions.append("share = ? COLLATE NOCASE")
values.append(share) values.append(share)
path = filters["path"] path = filters["path"]
if path: if path and include_path:
conditions.append("lower(path) LIKE ? ESCAPE '\\'") conditions.append("lower(path) LIKE ? ESCAPE '\\'")
values.append(f"%{escape_like(path)}%") values.append(f"%{escape_like(path)}%")
action = filters["action"] or filters["operation"] action = filters["action"] or filters["operation"]
@@ -245,6 +568,197 @@ def activity_conditions(
return conditions, values 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( def query_activity(
conn: sqlite3.Connection, conn: sqlite3.Connection,
start: dt.date, start: dt.date,
@@ -252,35 +766,81 @@ def query_activity(
params: Dict[str, List[str]], params: Dict[str, List[str]],
) -> Dict[str, object]: ) -> Dict[str, object]:
limit = min(500, max(1, int(params.get("limit", ["100"])[0]))) limit = min(500, max(1, int(params.get("limit", ["100"])[0])))
conditions, values = activity_conditions(start, end, params, include_filters=True) ready = rollups_ready(conn)
where_sql = " AND ".join(conditions) path = activity_filters(params)["path"]
matched = int( path_ids = indexed_path_ids(conn, path) if ready and path else None
conn.execute( if (
f"SELECT count(*) FROM audit_events WHERE {where_sql}", path_ids
values, and len(path_ids) > 1
).fetchone()[0] 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() cursor = params.get("cursor", [""])[0].strip()
page_conditions = list(conditions) page_conditions = list(conditions)
page_values = list(values) 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: if cursor:
try: try:
cursor_second, cursor_id = (int(part) for part in cursor.split(":", 1)) cursor_second, cursor_id = (int(part) for part in cursor.split(":", 1))
except (TypeError, ValueError) as exc: except (TypeError, ValueError) as exc:
raise ValueError("Ungültiger Seitenzeiger") from exc raise ValueError("Ungültiger Seitenzeiger") from exc
page_conditions.append( cursor_condition = (
"(occurred_second < ? OR (occurred_second = ? AND id < ?))" "(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)) page_values.extend((cursor_second, cursor_second, cursor_id))
rows = conn.execute( rows = conn.execute(
f""" f"""
SELECT id, occurred_at, occurred_second, ingested_at, user, client_ip, SELECT id, occurred_at, occurred_second, ingested_at, user, client_ip,
client, share, action, result, success, path, source client, share, action, result, success, path, source
FROM audit_events FROM {from_sql}
WHERE {" AND ".join(page_conditions)} WHERE {" AND ".join(page_conditions)}
ORDER BY occurred_second DESC, id DESC ORDER BY {order_sql}
LIMIT ? LIMIT ?
""", """,
(*page_values, limit + 1), (*page_values, limit + 1),
@@ -292,40 +852,65 @@ def query_activity(
if has_more and rows: if has_more and rows:
next_cursor = f"{int(rows[-1]['occurred_second'])}:{int(rows[-1]['id'])}" next_cursor = f"{int(rows[-1]['occurred_second'])}:{int(rows[-1]['id'])}"
facet_conditions, facet_values = activity_conditions( include_count = query_flag(params, "count")
start, end, params, include_filters=False matched: Optional[int] = None
) if include_count:
facet_where = " AND ".join(facet_conditions) if ready:
matched = rollup_matched(conn, start, end, params)
def distinct(column: str) -> List[str]: else:
return [ matched = int(
str(row[0]) conn.execute(
for row in conn.execute(
f""" f"""
SELECT DISTINCT {column} SELECT count(*)
FROM audit_events FROM audit_events
WHERE {facet_where} AND {column} <> '' WHERE {" AND ".join(conditions)}
ORDER BY {column} COLLATE NOCASE
""", """,
facet_values, 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 { return {
"events": events, "events": events,
"hasMore": has_more,
"nextCursor": next_cursor, "nextCursor": next_cursor,
"matched": matched, "matched": matched,
"facets": { "matchedExact": matched is not None,
"users": distinct("user"), "facets": facets,
"shares": distinct("share"),
"operations": actions,
"actions": actions,
},
} }
def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, object]: def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, object]:
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( row = conn.execute(
""" """
SELECT count(DISTINCT substr(occurred_at, 1, 10)), SELECT count(DISTINCT substr(occurred_at, 1, 10)),
@@ -334,6 +919,10 @@ def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, obj
FROM audit_events FROM audit_events
""" """
).fetchone() ).fetchone()
days = int(row[0] or 0)
oldest = row[1]
newest = row[2]
database_bytes = 0 database_bytes = 0
for suffix in ("", "-wal"): for suffix in ("", "-wal"):
try: try:
@@ -341,10 +930,10 @@ def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, obj
except OSError: except OSError:
pass pass
return { return {
"days": int(row[0] or 0), "days": days,
"bytes": database_bytes, "bytes": database_bytes,
"oldest": row[1], "oldest": oldest,
"newest": row[2], "newest": newest,
} }
+13 -2
View File
@@ -260,21 +260,32 @@ async function renderActivity() {
</section>`; </section>`;
const form = document.querySelector("#activity-filter"); const form = document.querySelector("#activity-filter");
let cursor = 0; let cursor = 0;
let loaded = 0;
let matched = null;
const load = async append => { const load = async append => {
const params = new URLSearchParams(new FormData(form)); const params = new URLSearchParams(new FormData(form));
params.set("limit", "100"); params.set("limit", "100");
if (cursor) params.set("cursor", String(cursor)); if (cursor) params.set("cursor", String(cursor));
if (append) {
params.set("count", "0");
params.set("facets", "0");
}
const result = await api(`/api/activity?${params}`); const result = await api(`/api/activity?${params}`);
state.activity = result; state.activity = result;
const rows = document.querySelector("#activity-rows"); const rows = document.querySelector("#activity-rows");
if (append) rows.insertAdjacentHTML("beforeend", eventRows(result.events)); else rows.innerHTML = eventRows(result.events); 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"); const more = document.querySelector("#load-more");
cursor = result.nextCursor || 0; cursor = result.nextCursor || 0;
more.hidden = !result.nextCursor; more.hidden = !result.hasMore;
if (!append) {
for (const [id, values] of [["users-list", result.facets.users], ["shares-list", result.facets.shares]]) { for (const [id, values] of [["users-list", result.facets.users], ["shares-list", result.facets.shares]]) {
document.querySelector(`#${id}`).innerHTML = values.filter(Boolean).map(value => `<option value="${esc(value)}">`).join(""); document.querySelector(`#${id}`).innerHTML = values.filter(Boolean).map(value => `<option value="${esc(value)}">`).join("");
} }
}
}; };
form.addEventListener("submit", async event => { event.preventDefault(); cursor = 0; try { await load(false); } catch (error) { notice(error.message); } }); form.addEventListener("submit", async event => { event.preventDefault(); cursor = 0; try { await load(false); } catch (error) { notice(error.message); } });
document.querySelector("#load-more").addEventListener("click", async () => { try { await load(true); } catch (error) { notice(error.message); } }); document.querySelector("#load-more").addEventListener("click", async () => { try { await load(true); } catch (error) { notice(error.message); } });
+1 -1
View File
@@ -836,7 +836,7 @@ class App:
def overview(self) -> Dict[str, object]: def overview(self) -> Dict[str, object]:
usage = self.usage.snapshot() usage = self.usage.snapshot()
recent = query_audit({"limit": ["12"]}) 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()} return {"usage": usage, "activeGroups": share_count(), "recentEvents": recent["events"], "eventCount": recent["matched"], "backup": backup_payload(), "audit": audit_archive_summary()}
def system_summary(self) -> Dict[str, object]: def system_summary(self) -> Dict[str, object]:
+18 -2
View File
@@ -372,7 +372,19 @@ def main() -> int:
).stdout.splitlines() ).stdout.splitlines()
) )
check( check(
{"shares", "audit_events", "audit_sources", "web_cache"}.issubset(tables), {
"shares",
"audit_events",
"audit_sources",
"audit_daily_totals",
"audit_daily_counts",
"audit_daily_facets",
"audit_rollup_state",
"audit_paths",
"audit_paths_fts",
"audit_path_events",
"web_cache",
}.issubset(tables),
f"shared SQLite tables are incomplete: {sorted(tables)}", f"shared SQLite tables are incomplete: {sorted(tables)}",
) )
indexes = set( indexes = set(
@@ -385,7 +397,11 @@ def main() -> int:
).stdout.splitlines() ).stdout.splitlines()
) )
check( check(
{"audit_events_time", "audit_events_user_time", "audit_events_action_time"}.issubset(indexes), {
"audit_events_time",
"audit_events_user_time",
"audit_events_action_time",
}.issubset(indexes),
f"audit indexes are incomplete: {sorted(indexes)}", f"audit indexes are incomplete: {sorted(indexes)}",
) )
duplicate_reads = engine_run( duplicate_reads = engine_run(
+164 -1
View File
@@ -551,7 +551,153 @@ class AuditQueryTests(unittest.TestCase):
self.assertNotIn("robot_svc", visible["facets"]["users"]) self.assertNotIn("robot_svc", visible["facets"]["users"])
self.assertEqual(set(visible["facets"]["actions"]), {"read", "delete"}) self.assertEqual(set(visible["facets"]["actions"]), {"read", "delete"})
def test_schema_has_filter_and_time_indexes(self): 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()
store = audit_store.AuditStore(os.path.join(tmpdir, "state.db"))
events = [
self.make_event(f"{today}T12:00:0{index}+00:00", "alice")
for index in range(4)
]
store.append_batch(events, {}, set())
statements = []
store.conn.set_trace_callback(statements.append)
try:
result = audit_store.query_activity(
store.conn, today, today, {"limit": ["1"]}
)
finally:
store.conn.set_trace_callback(None)
store.close()
normalized = [" ".join(statement.casefold().split()) for statement in statements]
self.assertEqual(result["matched"], 4)
self.assertTrue(result["matchedExact"])
self.assertTrue(result["hasMore"])
self.assertTrue(
any("from audit_daily_totals" in statement for statement in normalized)
)
self.assertTrue(
any("from audit_daily_facets" in statement for statement in normalized)
)
self.assertFalse(
any(
"select count(*) from audit_events" in statement
for statement in normalized
)
)
self.assertFalse(
any(
"select distinct user 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()
store = audit_store.AuditStore(os.path.join(tmpdir, "state.db"))
events = [
self.make_event(f"{today}T12:00:0{index}+00:00", "alice")
for index in range(3)
]
store.append_batch(events, {}, set())
statements = []
store.conn.set_trace_callback(statements.append)
try:
result = audit_store.query_activity(
store.conn,
today,
today,
{"path": ["file"], "limit": ["1"], "facets": ["0"]},
)
finally:
store.conn.set_trace_callback(None)
store.close()
normalized = [" ".join(statement.casefold().split()) for statement in statements]
self.assertIsNone(result["matched"])
self.assertFalse(result["matchedExact"])
self.assertTrue(result["hasMore"])
self.assertTrue(
any("from audit_paths_fts" in statement for statement in normalized)
)
self.assertTrue(
any("from audit_path_events" in statement for statement in normalized)
)
self.assertTrue(
any("pe.path_id =" in statement for statement in normalized)
)
self.assertFalse(
any("lower(path) like" in statement for statement in normalized)
)
self.assertFalse(
any(
"select count(*) from audit_events" in statement
for statement in normalized
)
)
def test_incrementally_backfills_legacy_events_without_double_counting_live_rows(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)
legacy = [
self.make_event(f"{today}T12:00:0{index}+00:00", "alice")
for index in range(3)
]
store.append_batch(legacy, {}, set())
store.conn.execute("UPDATE audit_events SET path_id = NULL")
store.conn.execute("DELETE FROM audit_path_events")
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)
live = self.make_event(f"{today}T12:00:09+00:00", "bob")
store.append_batch([live], {}, set())
before = audit_store.query_activity(
store.conn, today, today, {"limit": ["10"]}
)
self.assertEqual(store.backfill_rollups(limit=2), 2)
self.assertEqual(store.backfill_rollups(limit=2), 1)
self.assertEqual(store.backfill_rollups(limit=2), 0)
after = audit_store.query_activity(
store.conn, today, today, {"limit": ["10"]}
)
total = store.conn.execute(
"SELECT sum(event_count) FROM audit_daily_totals"
).fetchone()[0]
ready = store.conn.execute(
"SELECT ready FROM audit_rollup_state WHERE singleton = 1"
).fetchone()[0]
missing_path_ids = store.conn.execute(
"SELECT count(*) FROM audit_events WHERE path_id IS NULL"
).fetchone()[0]
path_event_count = store.conn.execute(
"SELECT count(*) FROM audit_path_events"
).fetchone()[0]
store.close()
self.assertEqual(before["matched"], 4)
self.assertEqual(after["matched"], 4)
self.assertEqual(total, 4)
self.assertEqual(ready, 1)
self.assertEqual(missing_path_ids, 0)
self.assertEqual(path_event_count, 4)
def test_schema_has_filter_indexes_and_rollup_tables(self):
with tempfile.TemporaryDirectory() as tmpdir: with tempfile.TemporaryDirectory() as tmpdir:
store = audit_store.AuditStore(os.path.join(tmpdir, "state.db")) store = audit_store.AuditStore(os.path.join(tmpdir, "state.db"))
try: try:
@@ -561,6 +707,12 @@ class AuditQueryTests(unittest.TestCase):
"SELECT name FROM sqlite_schema WHERE type = 'index'" "SELECT name FROM sqlite_schema WHERE type = 'index'"
) )
} }
tables = {
row[0]
for row in store.conn.execute(
"SELECT name FROM sqlite_schema WHERE type = 'table'"
)
}
finally: finally:
store.close() store.close()
self.assertTrue( self.assertTrue(
@@ -574,6 +726,17 @@ class AuditQueryTests(unittest.TestCase):
"audit_events_result_time", "audit_events_result_time",
}.issubset(indexes) }.issubset(indexes)
) )
self.assertTrue(
{
"audit_daily_totals",
"audit_daily_counts",
"audit_daily_facets",
"audit_rollup_state",
"audit_paths",
"audit_paths_fts",
"audit_path_events",
}.issubset(tables)
)
def test_drops_only_known_legacy_audit_files(self): def test_drops_only_known_legacy_audit_files(self):
with tempfile.TemporaryDirectory() as tmpdir: with tempfile.TemporaryDirectory() as tmpdir: