even more performance benefits

This commit is contained in:
Ludwig Lehnert
2026-08-12 13:26:31 +00:00
parent 0ea17c6401
commit 29340d778d
5 changed files with 353 additions and 92 deletions
+163 -29
View File
@@ -11,14 +11,16 @@ from typing import Dict, Iterable, List, Optional, Set, Tuple
try:
from .audit_policy import (
account_name,
read_deduplication_key,
read_deduplication_fingerprint,
read_deduplication_fingerprint_values,
skipped_user_suffixes,
)
from .state_db import connect_state_db
except ImportError:
from audit_policy import (
account_name,
read_deduplication_key,
read_deduplication_fingerprint,
read_deduplication_fingerprint_values,
skipped_user_suffixes,
)
from state_db import connect_state_db
@@ -47,6 +49,13 @@ CREATE TABLE IF NOT EXISTS audit_sources (
inode INTEGER NOT NULL,
offset INTEGER NOT NULL
) WITHOUT ROWID;
CREATE TABLE IF NOT EXISTS audit_read_dedup (
fingerprint BLOB PRIMARY KEY,
occurred_second INTEGER NOT NULL,
event_id INTEGER
) WITHOUT ROWID;
CREATE INDEX IF NOT EXISTS audit_read_dedup_time
ON audit_read_dedup (occurred_second);
CREATE INDEX IF NOT EXISTS audit_events_main_time
ON audit_events (occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
@@ -138,7 +147,7 @@ CREATE TABLE IF NOT EXISTS audit_rollup_state (
backfill_next_id INTEGER NOT NULL,
backfill_max_id INTEGER NOT NULL,
ready INTEGER NOT NULL CHECK (ready IN (0, 1)),
policy_version INTEGER NOT NULL DEFAULT 3
policy_version INTEGER NOT NULL DEFAULT 4
);
"""
@@ -212,7 +221,7 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
WHERE singleton = 1
"""
).fetchone()
if state is not None and int(state["policy_version"]) < 3:
if state is not None and int(state["policy_version"]) < 4:
for table in (
"audit_daily_totals",
"audit_daily_counts",
@@ -223,7 +232,7 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
"""
UPDATE audit_rollup_state
SET backfill_next_id = 1, backfill_max_id = ?, ready = ?,
policy_version = 3
policy_version = 4
WHERE singleton = 1
""",
(max_id, int(max_id == 0)),
@@ -231,6 +240,21 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
conn.commit()
READ_DEDUP_DEFAULT_WINDOW_SECONDS = 2 * 86400
def read_deduplication_cutoff() -> int:
try:
window = int(os.getenv(
"AUDIT_READ_DEDUP_WINDOW_SECONDS",
str(READ_DEDUP_DEFAULT_WINDOW_SECONDS),
))
except ValueError:
window = READ_DEDUP_DEFAULT_WINDOW_SECONDS
now = int(dt.datetime.now(dt.timezone.utc).timestamp())
return now - max(1, window)
def parse_event_second(timestamp: object) -> int:
parsed = dt.datetime.fromisoformat(str(timestamp).replace("Z", "+00:00"))
if parsed.tzinfo is None:
@@ -257,22 +281,13 @@ def row_to_event(row: sqlite3.Row) -> Dict[str, object]:
}
def last_event(conn: sqlite3.Connection) -> Optional[Dict[str, object]]:
row = conn.execute(
"""
SELECT id, occurred_at, ingested_at, user, client_ip, client, share,
action, result, success, path, source
FROM audit_events
ORDER BY id DESC
LIMIT 1
"""
).fetchone()
return row_to_event(row) if row is not None else None
class AuditStore:
def __init__(self, database_path: Optional[str] = None):
self.conn = connect_state_db(database_path)
self.conn.create_function(
"audit_read_fingerprint", 7,
read_deduplication_fingerprint_values, deterministic=True,
)
ensure_audit_schema(self.conn)
def close(self) -> None:
@@ -299,10 +314,13 @@ class AuditStore:
totals = Counter()
counts = Counter()
facets = set()
previous_key = read_deduplication_key(last_event(self.conn) or {})
seen_read_fingerprints = set()
for event in events:
current_key = read_deduplication_key(event)
if current_key is not None and current_key == previous_key:
read_fingerprint = read_deduplication_fingerprint(event)
if (
read_fingerprint is not None
and read_fingerprint in seen_read_fingerprints
):
continue
occurred_second = parse_event_second(event["timestamp"])
user = str(event["user"])
@@ -311,7 +329,6 @@ class AuditStore:
action = str(event["action"])
result = str(event["result"])
success = int(bool(event["success"]))
day = occurred_second // 86400
event_rows.append(
(
str(event["timestamp"]),
@@ -327,17 +344,62 @@ class AuditStore:
success,
str(event["path"]),
str(event["source"]),
read_fingerprint,
)
)
totals[(day, account)] += 1
counts[(day, account, user, share, action, result, success)] += 1
for kind, value in (("user", user), ("share", share), ("action", action)):
if value:
facets.add((kind, day, account, value))
previous_key = current_key
if read_fingerprint is not None:
seen_read_fingerprints.add(read_fingerprint)
self.conn.execute("BEGIN IMMEDIATE")
try:
fingerprints = [
row[13] for row in event_rows if row[13] is not None
]
existing_fingerprints = set()
for offset in range(0, len(fingerprints), 500):
chunk = fingerprints[offset : offset + 500]
placeholders = ",".join("?" for _ in chunk)
existing_fingerprints.update(
bytes(row[0])
for row in self.conn.execute(
f"""
SELECT fingerprint
FROM audit_read_dedup
WHERE fingerprint IN ({placeholders})
""",
chunk,
)
)
if existing_fingerprints:
event_rows = [
row
for row in event_rows
if (
row[13] is None
or row[13] not in existing_fingerprints
)
]
for row in event_rows:
day = int(row[1]) // 86400
account = str(row[4])
user = str(row[3])
share = str(row[7])
action = str(row[8])
result = str(row[9])
success = int(row[10])
totals[(day, account)] += 1
counts[(
day, account, user, share, action, result, success
)] += 1
for kind, value in (
("user", user),
("share", share),
("action", action),
):
if value:
facets.add((kind, day, account, value))
if event_rows:
paths = sorted({str(row[11]) for row in event_rows})
self.conn.executemany(
@@ -383,6 +445,24 @@ class AuditStore:
""",
(len(indexed_rows),),
)
last_inserted_id = int(
self.conn.execute(
"SELECT max(id) FROM audit_events"
).fetchone()[0]
)
first_inserted_id = last_inserted_id - len(event_rows) + 1
self.conn.executemany(
"""
INSERT INTO audit_read_dedup (
fingerprint, occurred_second, event_id
) VALUES (?, ?, ?)
""",
(
(row[13], int(row[1]), first_inserted_id + offset)
for offset, row in enumerate(event_rows)
if row[13] is not None
),
)
self.conn.executemany(
ROLLUP_TOTAL_SQL,
((*key, count) for key, count in totals.items()),
@@ -393,6 +473,11 @@ class AuditStore:
)
self.conn.executemany(ROLLUP_FACET_SQL, facets)
self.conn.execute(
"DELETE FROM audit_read_dedup WHERE occurred_second < ?",
(read_deduplication_cutoff(),),
)
if source_updates:
self.conn.executemany(
"""
@@ -422,7 +507,7 @@ class AuditStore:
raise
return len(event_rows)
def backfill_rollups(self, limit: int = 50000) -> int:
def backfill_rollups(self, limit: int = 5000) -> int:
"""Backfill one bounded legacy-event chunk without delaying live appends."""
state = self.conn.execute(
"""
@@ -459,6 +544,55 @@ class AuditStore:
end_id = int(chunk[1])
self.conn.execute("BEGIN IMMEDIATE")
try:
dedup_cutoff = read_deduplication_cutoff()
self.conn.execute(
"DELETE FROM audit_read_dedup WHERE occurred_second < ?",
(dedup_cutoff,),
)
self.conn.execute(
"""
INSERT INTO audit_read_dedup (
fingerprint, occurred_second, event_id
)
SELECT audit_read_fingerprint(
e.occurred_at, e.user, e.client_ip, e.share, e.path,
e.success, e.result
), e.occurred_second, e.id
FROM audit_events AS e
WHERE e.id BETWEEN ? AND ?
AND e.action = 'read'
AND e.occurred_second >= ?
ORDER BY e.id
ON CONFLICT (fingerprint) DO UPDATE SET
event_id = coalesce(
audit_read_dedup.event_id, excluded.event_id
)
""",
(start_id, end_id, dedup_cutoff),
)
duplicate_read_sql = """
SELECT e.id
FROM audit_events AS e
JOIN audit_read_dedup AS d
ON d.fingerprint = audit_read_fingerprint(
e.occurred_at, e.user, e.client_ip, e.share, e.path,
e.success, e.result
)
WHERE e.id BETWEEN ? AND ?
AND e.action = 'read'
AND e.occurred_second >= ?
AND (d.event_id IS NULL OR d.event_id <> e.id)
"""
self.conn.execute(
f"DELETE FROM audit_path_events "
f"WHERE event_id IN ({duplicate_read_sql})",
(start_id, end_id, dedup_cutoff),
)
self.conn.execute(
f"DELETE FROM audit_events "
f"WHERE id IN ({duplicate_read_sql})",
(start_id, end_id, dedup_cutoff),
)
self.conn.execute(
"""
INSERT OR IGNORE INTO audit_paths (path)