fslogix separation

This commit is contained in:
Ludwig Lehnert
2026-08-12 10:20:32 +00:00
parent 14874e504e
commit 0ea17c6401
10 changed files with 633 additions and 153 deletions
+4 -1
View File
@@ -161,13 +161,16 @@ def main() -> int:
log(f"Watching {SAMBA_LOG_GLOB}; storing events in SQLite")
try:
while not STOP:
backfilled = 0
try:
count = collect_once(store)
if count:
log(f"Stored {count} event(s)")
store.backfill_rollups()
backfilled = store.backfill_rollups()
except Exception as exc: # pylint: disable=broad-except
log(f"Collector cycle failed: {exc}")
if backfilled:
continue
time.sleep(POLL_SECONDS)
finally:
store.close()
+232 -133
View File
@@ -16,7 +16,11 @@ try:
)
from .state_db import connect_state_db
except ImportError:
from audit_policy import account_name, read_deduplication_key, skipped_user_suffixes
from audit_policy import (
account_name,
read_deduplication_key,
skipped_user_suffixes,
)
from state_db import connect_state_db
@@ -43,20 +47,46 @@ CREATE TABLE IF NOT EXISTS audit_sources (
inode INTEGER NOT NULL,
offset INTEGER NOT NULL
) WITHOUT ROWID;
CREATE INDEX IF NOT EXISTS audit_events_time
ON audit_events (occurred_second DESC, id DESC);
CREATE INDEX IF NOT EXISTS audit_events_action_time
ON audit_events (action, occurred_second DESC, id DESC);
CREATE INDEX IF NOT EXISTS audit_events_success_time
ON audit_events (success, occurred_second DESC, id DESC);
CREATE INDEX IF NOT EXISTS audit_events_user_time
ON audit_events (user COLLATE NOCASE, occurred_second DESC, id DESC);
CREATE INDEX IF NOT EXISTS audit_events_account_time
ON audit_events (account, occurred_second DESC, id DESC);
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 INDEX IF NOT EXISTS audit_events_main_time
ON audit_events (occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_main_action_time
ON audit_events (action, occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_main_success_time
ON audit_events (success, occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_main_user_time
ON audit_events (user COLLATE NOCASE, occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_main_account_time
ON audit_events (account, occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_main_share_time
ON audit_events (share COLLATE NOCASE, occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_main_result_time
ON audit_events (result COLLATE NOCASE, occurred_second DESC, id DESC)
WHERE share <> 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_fslogix_time
ON audit_events (occurred_second DESC, id DESC)
WHERE share = 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_fslogix_action_time
ON audit_events (action, occurred_second DESC, id DESC)
WHERE share = 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_fslogix_success_time
ON audit_events (success, occurred_second DESC, id DESC)
WHERE share = 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_fslogix_user_time
ON audit_events (user COLLATE NOCASE, occurred_second DESC, id DESC)
WHERE share = 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_fslogix_account_time
ON audit_events (account, occurred_second DESC, id DESC)
WHERE share = 'FSLogix' COLLATE NOCASE;
CREATE INDEX IF NOT EXISTS audit_events_fslogix_result_time
ON audit_events (result COLLATE NOCASE, occurred_second DESC, id DESC)
WHERE share = 'FSLogix' COLLATE NOCASE;
CREATE TABLE IF NOT EXISTS audit_paths (
id INTEGER PRIMARY KEY,
path TEXT NOT NULL UNIQUE
@@ -107,7 +137,8 @@ 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))
ready INTEGER NOT NULL CHECK (ready IN (0, 1)),
policy_version INTEGER NOT NULL DEFAULT 3
);
"""
@@ -135,6 +166,17 @@ VALUES (?, ?, ?, ?)
def ensure_audit_schema(conn: sqlite3.Connection) -> None:
conn.executescript(AUDIT_SCHEMA)
for legacy_index in (
"audit_events_time",
"audit_events_action_time",
"audit_events_success_time",
"audit_events_user_time",
"audit_events_account_time",
"audit_events_share_time",
"audit_events_result_time",
):
conn.execute(f"DROP INDEX IF EXISTS {legacy_index}")
columns = {
str(row["name"])
for row in conn.execute("PRAGMA table_info(audit_events)")
@@ -142,6 +184,16 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
if "path_id" not in columns:
conn.execute("ALTER TABLE audit_events ADD COLUMN path_id INTEGER")
state_columns = {
str(row["name"])
for row in conn.execute("PRAGMA table_info(audit_rollup_state)")
}
if "policy_version" not in state_columns:
conn.execute(
"ALTER TABLE audit_rollup_state "
"ADD COLUMN policy_version INTEGER NOT NULL DEFAULT 1"
)
max_id = int(
conn.execute("SELECT coalesce(max(id), 0) FROM audit_events").fetchone()[0]
)
@@ -153,6 +205,29 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
""",
(max_id, int(max_id == 0)),
)
state = conn.execute(
"""
SELECT ready, policy_version
FROM audit_rollup_state
WHERE singleton = 1
"""
).fetchone()
if state is not None and int(state["policy_version"]) < 3:
for table in (
"audit_daily_totals",
"audit_daily_counts",
"audit_daily_facets",
):
conn.execute(f"DELETE FROM {table}")
conn.execute(
"""
UPDATE audit_rollup_state
SET backfill_next_id = 1, backfill_max_id = ?, ready = ?,
policy_version = 3
WHERE singleton = 1
""",
(max_id, int(max_id == 0)),
)
conn.commit()
@@ -347,7 +422,7 @@ class AuditStore:
raise
return len(event_rows)
def backfill_rollups(self, limit: int = 10000) -> int:
def backfill_rollups(self, limit: int = 50000) -> int:
"""Backfill one bounded legacy-event chunk without delaying live appends."""
state = self.conn.execute(
"""
@@ -459,7 +534,8 @@ class AuditStore:
)
SELECT ?, occurred_second / 86400, account, {column}
FROM audit_events
WHERE id BETWEEN ? AND ? AND {column} <> ''
WHERE id BETWEEN ? AND ?
AND {column} <> ''
GROUP BY occurred_second / 86400, account, {column}
""",
(kind, start_id, end_id),
@@ -519,15 +595,28 @@ def excluded_account_conditions() -> Tuple[List[str], List[object]]:
return conditions, values
def activity_stream_condition(stream: str) -> str:
if stream == "main":
return "share <> 'FSLogix' COLLATE NOCASE"
if stream == "fslogix":
return "share = 'FSLogix' COLLATE NOCASE"
raise ValueError("Ungültiger Aktivitätsstrom")
def activity_conditions(
start: dt.date,
end: dt.date,
params: Dict[str, List[str]],
*,
stream: str,
include_filters: bool,
include_path: bool = True,
) -> Tuple[List[str], List[object]]:
conditions = ["occurred_second >= ?", "occurred_second < ?"]
conditions = [
"occurred_second >= ?",
"occurred_second < ?",
activity_stream_condition(stream),
]
values: List[object] = [
date_seconds(start),
date_seconds(end + dt.timedelta(days=1)),
@@ -568,14 +657,44 @@ def activity_conditions(
return conditions, values
def rollups_ready(conn: sqlite3.Connection) -> bool:
def rollup_status(conn: sqlite3.Connection) -> Dict[str, object]:
try:
row = conn.execute(
"SELECT ready FROM audit_rollup_state WHERE singleton = 1"
"""
SELECT backfill_next_id, backfill_max_id, ready
FROM audit_rollup_state
WHERE singleton = 1
"""
).fetchone()
except sqlite3.Error:
return False
return row is not None and bool(row[0])
row = None
if row is None:
return {
"ready": False,
"processedEvents": 0,
"totalEvents": 0,
"percent": 0.0,
}
total = max(0, int(row["backfill_max_id"]))
ready = bool(row["ready"])
processed = (
total
if ready
else min(total, max(0, int(row["backfill_next_id"]) - 1))
)
percent = (
100.0
if ready or total == 0
else round(processed * 100 / total, 1)
)
return {
"ready": ready,
"processedEvents": processed,
"totalEvents": total,
"percent": percent,
}
def rollup_range(
@@ -597,44 +716,40 @@ def rollup_matched(
start: dt.date,
end: dt.date,
params: Dict[str, List[str]],
*,
stream: 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)
conditions.append(activity_stream_condition(stream))
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}
FROM audit_daily_counts
WHERE {" AND ".join(conditions)}
""",
values,
@@ -643,56 +758,30 @@ def rollup_matched(
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]],
*,
stream: str,
) -> Dict[str, List[str]]:
conditions, values = activity_conditions(
start, end, params, include_filters=False
)
where_sql = " AND ".join(conditions)
range_conditions, range_values = rollup_range(start, end)
def distinct(column: str) -> List[str]:
conditions = [
*range_conditions,
activity_stream_condition(stream),
f"{column} <> ?",
]
return [
str(row[0])
for row in conn.execute(
f"""
SELECT DISTINCT {column}
FROM audit_events
WHERE {where_sql} AND {column} <> ''
FROM audit_daily_counts
WHERE {" AND ".join(conditions)}
ORDER BY {column} COLLATE NOCASE
""",
values,
(*range_values, ""),
)
]
@@ -764,9 +853,12 @@ def query_activity(
start: dt.date,
end: dt.date,
params: Dict[str, List[str]],
*,
stream: str = "main",
) -> Dict[str, object]:
limit = min(500, max(1, int(params.get("limit", ["100"])[0])))
ready = rollups_ready(conn)
indexing = rollup_status(conn)
ready = bool(indexing["ready"])
path = activity_filters(params)["path"]
path_ids = indexed_path_ids(conn, path) if ready and path else None
if (
@@ -779,6 +871,7 @@ def query_activity(
start,
end,
params,
stream=stream,
include_filters=True,
include_path=path_ids is None,
)
@@ -854,32 +947,23 @@ def query_activity(
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)
if include_count and ready:
matched = rollup_matched(conn, start, end, params, stream=stream)
if include_count and matched is None and not cursor and not has_more:
matched = len(events)
empty_facets = {
"users": [],
"shares": [],
"operations": [],
"actions": [],
}
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": []}
facets = (
rollup_facets(conn, start, end, stream=stream)
if include_facets and ready
else empty_facets
)
return {
"events": events,
@@ -888,40 +972,54 @@ def query_activity(
"matched": matched,
"matchedExact": matched is not None,
"facets": facets,
"indexing": indexing,
"stream": stream,
}
def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, object]:
if rollups_ready(conn):
indexing = rollup_status(conn)
def second_day(value: object) -> Optional[str]:
if value is None:
return None
return dt.datetime.fromtimestamp(
int(value), tz=dt.timezone.utc
).date().isoformat()
if indexing["ready"]:
row = conn.execute(
"""
SELECT count(DISTINCT day), min(day), max(day)
FROM audit_daily_totals
FROM audit_daily_counts
WHERE share <> 'FSLogix' COLLATE NOCASE
"""
).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])
oldest = second_day(None if row[1] is None else int(row[1]) * 86400)
newest = second_day(None if row[2] is None else int(row[2]) * 86400)
else:
row = conn.execute(
oldest_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
SELECT occurred_second FROM audit_events
WHERE share <> 'FSLogix' COLLATE NOCASE
ORDER BY occurred_second, id LIMIT 1
"""
).fetchone()
days = int(row[0] or 0)
oldest = row[1]
newest = row[2]
newest_row = conn.execute(
"""
SELECT occurred_second FROM audit_events
WHERE share <> 'FSLogix' COLLATE NOCASE
ORDER BY occurred_second DESC, id DESC LIMIT 1
"""
).fetchone()
oldest = second_day(oldest_row[0]) if oldest_row else None
newest = second_day(newest_row[0]) if newest_row else None
days = (
(dt.date.fromisoformat(newest) - dt.date.fromisoformat(oldest)).days + 1
if oldest and newest
else 0
)
database_bytes = 0
for suffix in ("", "-wal"):
@@ -934,6 +1032,7 @@ def audit_summary(conn: sqlite3.Connection, database_path: str) -> Dict[str, obj
"bytes": database_bytes,
"oldest": oldest,
"newest": newest,
"indexing": indexing,
}
+3
View File
@@ -204,6 +204,7 @@ class BackupStatus:
self.write(
state="starting",
trigger=trigger,
workerPid=os.getpid(),
startedAt=dt.datetime.now(dt.timezone.utc).isoformat(timespec="seconds"),
finishedAt=None,
destination=destination,
@@ -241,6 +242,7 @@ class BackupStatus:
def complete(self, message: str) -> None:
self.write(
state="completed",
workerPid=None,
finishedAt=dt.datetime.now(dt.timezone.utc).isoformat(timespec="seconds"),
percent=100.0,
transferredBytes=self.value.get("totalBytes", 0),
@@ -252,6 +254,7 @@ class BackupStatus:
def fail(self, message: str) -> None:
self.write(
state="failed",
workerPid=None,
finishedAt=dt.datetime.now(dt.timezone.utc).isoformat(timespec="seconds"),
activeFiles=[],
message=message,
+15 -4
View File
@@ -43,6 +43,7 @@ function backupMessage(data) {
if (message.startsWith("Syncing ")) return `${message.slice(8)} wird synchronisiert.`;
const completed = message.match(/^Backup completed; (\d+) snapshot\(s\) retained$/);
if (completed) return `Sicherung abgeschlossen; ${completed[1]} Sicherungsstände werden aufbewahrt.`;
if (data.interrupted) return "Sicherung wurde durch einen Neustart oder Prozessabbruch unterbrochen.";
if (data.state === "failed") return "Sicherung fehlgeschlagen. Einzelheiten stehen im Sicherungsprotokoll.";
if (data.state === "completed") return "Sicherung abgeschlossen.";
if (data.state === "running") return "Sicherung läuft.";
@@ -102,6 +103,7 @@ function routeFor(path) {
if (path.startsWith("/storage/users")) return "storage-users";
if (path.startsWith("/shares")) return "shares";
if (path.startsWith("/reconciliation")) return "reconciliation";
if (path.startsWith("/activity/fslogix")) return "activity-fslogix";
if (path.startsWith("/activity")) return "activity";
if (path.startsWith("/backup")) return "backup";
if (path.startsWith("/report")) return "report";
@@ -124,6 +126,7 @@ async function navigate(path, replace = false) {
if (route === "storage-data") await renderStorage("data");
if (route === "storage-users") await renderStorage("users");
if (route === "activity") await renderActivity();
if (route === "activity-fslogix") await renderActivity("fslogix");
if (route === "backup") await renderBackup();
if (route === "report") await renderReport();
if (route === "system") await renderSystem();
@@ -238,10 +241,17 @@ async function renderStorage(type) {
draw();
}
async function renderActivity() {
async function renderActivity(stream = "main") {
const today = new Date();
const yesterday = new Date(Date.now() - 86400000);
content.innerHTML = pageHead("Aktivitätsprotokoll", "Lese-, Schreib-, Umbenennungs- und Löschvorgänge nach Datum, Identität, Freigabe, Ergebnis oder Pfad durchsuchen.") + `
const fslogix = stream === "fslogix";
const endpoint = fslogix ? "/api/fslogix-activity" : "/api/activity";
content.innerHTML = pageHead(
fslogix ? "FSLogix-Protokoll" : "Aktivitätsprotokoll",
fslogix
? "FSLogix-Profilcontainer-Aktivitäten getrennt vom regulären Dateiprotokoll durchsuchen."
: "Lese-, Schreib-, Umbenennungs- und Löschvorgänge nach Datum, Identität, Freigabe, Ergebnis oder Pfad durchsuchen.",
) + `
<section class="panel">
<form id="activity-filter" class="filters">
<label>Von (UTC)<input name="from" type="date" value="${isoDay(yesterday)}" required></label>
@@ -270,14 +280,15 @@ async function renderActivity() {
params.set("count", "0");
params.set("facets", "0");
}
const result = await api(`/api/activity?${params}`);
const result = await api(`${endpoint}?${params}`);
state.activity = result;
const rows = document.querySelector("#activity-rows");
if (append) rows.insertAdjacentHTML("beforeend", eventRows(result.events)); else rows.innerHTML = eventRows(result.events);
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 indexNote = result.indexing?.ready === false ? ` · Indexaufbau ${Number(result.indexing.percent || 0).toLocaleString("de-DE")} %` : "";
document.querySelector("#activity-summary").textContent = `${shown} passende Ereignisse${indexNote}`;
const more = document.querySelector("#load-more");
cursor = result.nextCursor || 0;
more.hidden = !result.hasMore;
+1
View File
@@ -32,6 +32,7 @@
<a href="/storage/data" data-route="storage-data">Datenbelegung</a>
<a href="/storage/users" data-route="storage-users">Benutzerbelegung</a>
<a href="/activity" data-route="activity">Aktivitätsprotokoll</a>
<a href="/activity/fslogix" data-route="activity-fslogix">FSLogix-Protokoll</a>
<a href="/backup" data-route="backup">Sicherungen</a>
<a href="/report" data-route="report">PDF-Bericht</a>
<a href="/system" data-route="system">System</a>
+1
View File
@@ -102,6 +102,7 @@ function backupMessage(data) {
if (message.startsWith("Syncing ")) return message.slice(8) + " wird synchronisiert.";
const completed = message.match(/^Backup completed; (\d+) snapshot\(s\) retained$/);
if (completed) return "Sicherung abgeschlossen; " + completed[1] + " Sicherungsstände werden aufbewahrt.";
if (data.interrupted) return "Sicherung wurde durch einen Neustart oder Prozessabbruch unterbrochen.";
if (data.state === "failed") return "Sicherung fehlgeschlagen.";
if (data.state === "completed") return "Sicherung abgeschlossen.";
if (data.state === "running") return "Sicherung läuft.";
+83 -4
View File
@@ -60,7 +60,7 @@ STATIC_ROOT = os.getenv("WEB_STATIC_DIR", "/app/web")
STATE_DB = STATE_DB_PATH
BACKUP_STATUS_FILE = os.getenv("BACKUP_STATUS_FILE", "/state/backup-status.json")
BACKUP_LOG_FILE = os.getenv("BACKUP_LOG_FILE", "/var/log/backup.log")
BACKUP_LOCK_FILE = "/state/backup.lock"
BACKUP_LOCK_FILE = os.path.join(os.getenv("STATE_ROOT", "/state"), "backup.lock")
RECONCILE_STATUS_FILE = os.getenv(
"RECONCILE_STATUS_FILE", "/state/reconcile-status.json"
)
@@ -76,6 +76,7 @@ SID_RE = re.compile(r"S-\d+(?:-\d+)+", re.IGNORECASE)
LOGIN_LIMIT: Dict[str, deque] = {}
LOGIN_LIMIT_LOCK = threading.Lock()
ACTION_LAUNCH_LOCK = threading.Lock()
ACTIVE_PROCESS_STATES = frozenset({"starting", "running"})
def log(message: str) -> None:
@@ -176,6 +177,69 @@ def read_json(path: str, default):
return default
def write_json_atomic(path: str, value: Dict[str, object]) -> None:
directory = os.path.dirname(path)
if directory:
os.makedirs(directory, exist_ok=True)
temp_path = f"{path}.recovery.tmp"
try:
with open(temp_path, "w", encoding="utf-8") as handle:
json.dump(value, handle, separators=(",", ":"), sort_keys=True)
handle.flush()
os.fsync(handle.fileno())
os.replace(temp_path, path)
except (OSError, TypeError, ValueError):
try:
os.remove(temp_path)
except OSError:
pass
raise
def reconcile_backup_process_status(value: Dict[str, object]) -> Dict[str, object]:
"""Replace stale active state when no worker owns the backup lock."""
if str(value.get("state", "")) not in ACTIVE_PROCESS_STATES:
value["processRunning"] = False
return value
try:
lock_dir = os.path.dirname(BACKUP_LOCK_FILE)
if lock_dir:
os.makedirs(lock_dir, exist_ok=True)
with open(BACKUP_LOCK_FILE, "a+", encoding="utf-8") as lock_file:
try:
fcntl.flock(lock_file, fcntl.LOCK_EX | fcntl.LOCK_NB)
except BlockingIOError:
value["processRunning"] = True
return value
latest = read_json(BACKUP_STATUS_FILE, value)
if str(latest.get("state", "")) in ACTIVE_PROCESS_STATES:
latest.update(
{
"state": "failed",
"finishedAt": now_utc().isoformat(timespec="seconds"),
"activeFiles": [],
"currentSource": None,
"workerPid": None,
"interrupted": True,
"message": (
"Backup interrupted because the worker process "
"is no longer running"
),
}
)
write_json_atomic(BACKUP_STATUS_FILE, latest)
log("Recovered stale backup status after worker interruption")
latest["processRunning"] = False
fcntl.flock(lock_file, fcntl.LOCK_UN)
return latest
except OSError as exc:
log(f"Unable to verify backup worker state: {exc}")
value["processRunning"] = None
return value
def base64url(value: bytes) -> str:
return base64.urlsafe_b64encode(value).rstrip(b"=").decode("ascii")
@@ -711,7 +775,9 @@ class DirectoryCache:
def query_audit(params: Dict[str, List[str]]) -> Dict[str, object]:
def query_audit(
params: Dict[str, List[str]], *, stream: str = "main"
) -> Dict[str, object]:
today = now_utc().date()
default_start = today - dt.timedelta(days=1)
try:
@@ -724,7 +790,7 @@ def query_audit(params: Dict[str, List[str]]) -> Dict[str, object]:
raise ValueError(f"Der Datumsbereich darf höchstens {max_days} Tage umfassen")
conn = connect_state_db(STATE_DB, read_only=True)
try:
return query_activity(conn, start, end, params)
return query_activity(conn, start, end, params, stream=stream)
finally:
conn.close()
@@ -747,6 +813,7 @@ def tail_lines(path: str, count: int) -> List[str]:
def backup_payload(include_log: bool = True) -> Dict[str, object]:
value = read_json(BACKUP_STATUS_FILE, {})
value = reconcile_backup_process_status(value)
value.pop("log", None)
configured = bool(os.getenv("BACKUP_DESTINATION", "").strip())
automatic = configured and env_bool("BACKUP_AUTO_ENABLED", True)
@@ -837,7 +904,17 @@ class App:
def overview(self) -> Dict[str, object]:
usage = self.usage.snapshot()
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()}
event_count = recent["matched"]
if event_count is None:
event_count = f"{len(recent['events'])}+"
return {
"usage": usage,
"activeGroups": share_count(),
"recentEvents": recent["events"],
"eventCount": event_count,
"backup": backup_payload(),
"audit": audit_archive_summary(),
}
def system_summary(self) -> Dict[str, object]:
checks = {}
@@ -1026,6 +1103,8 @@ class Handler(BaseHTTPRequestHandler):
self.send_json(APP.usage.snapshot())
elif path == "/api/activity":
self.send_json(query_audit(params))
elif path == "/api/fslogix-activity":
self.send_json(query_audit(params, stream="fslogix"))
elif path == "/api/backup":
self.send_json(
backup_payload(include_log=query_includes_log(params))