7 day trash bin
This commit is contained in:
+14
-10
@@ -49,10 +49,12 @@ def skip_user(user: str) -> bool:
|
||||
account = account_name(user).casefold()
|
||||
return any(account.endswith(suffix) for suffix in skipped_user_suffixes())
|
||||
|
||||
ReadEventKey = Tuple[str, ...]
|
||||
DeduplicationKey = Tuple[str, ...]
|
||||
|
||||
def read_deduplication_key(event: Mapping[str, object]) -> Optional[ReadEventKey]:
|
||||
if str(event.get("action", "")).casefold() != "read":
|
||||
def deduplication_key(event: Mapping[str, object]) -> Optional[DeduplicationKey]:
|
||||
action = str(event.get("action", "")).casefold()
|
||||
share = str(event.get("share", ""))
|
||||
if action != "read" and share.casefold() != "fslogix":
|
||||
return None
|
||||
try:
|
||||
timestamp = dt.datetime.fromisoformat(
|
||||
@@ -71,18 +73,19 @@ def read_deduplication_key(event: Mapping[str, object]) -> Optional[ReadEventKey
|
||||
success = bool(event.get("success", False))
|
||||
return (
|
||||
second,
|
||||
action,
|
||||
str(event.get("user", "")),
|
||||
str(event.get("clientIp", "")),
|
||||
str(event.get("share", "")),
|
||||
share,
|
||||
str(event.get("path", "")),
|
||||
"success" if success else "failure",
|
||||
"" if success else str(event.get("result", "")),
|
||||
)
|
||||
|
||||
def read_deduplication_fingerprint(
|
||||
def deduplication_fingerprint(
|
||||
event: Mapping[str, object]
|
||||
) -> Optional[bytes]:
|
||||
key = read_deduplication_key(event)
|
||||
key = deduplication_key(event)
|
||||
if key is None:
|
||||
return None
|
||||
digest = hashlib.blake2b(digest_size=16)
|
||||
@@ -92,8 +95,9 @@ def read_deduplication_fingerprint(
|
||||
digest.update(encoded)
|
||||
return digest.digest()
|
||||
|
||||
def read_deduplication_fingerprint_values(
|
||||
def deduplication_fingerprint_values(
|
||||
timestamp: object,
|
||||
action: object,
|
||||
user: object,
|
||||
client_ip: object,
|
||||
share: object,
|
||||
@@ -101,9 +105,9 @@ def read_deduplication_fingerprint_values(
|
||||
success: object,
|
||||
result: object,
|
||||
) -> bytes:
|
||||
"""Fingerprint legacy SQLite columns with the live-ingest policy."""
|
||||
fingerprint = read_deduplication_fingerprint({
|
||||
"action": "read",
|
||||
"""Fingerprint SQLite columns with the live-ingest policy."""
|
||||
fingerprint = deduplication_fingerprint({
|
||||
"action": action,
|
||||
"timestamp": timestamp,
|
||||
"user": user,
|
||||
"clientIp": client_ip,
|
||||
|
||||
+50
-44
@@ -11,16 +11,16 @@ from typing import Dict, Iterable, List, Optional, Set, Tuple
|
||||
try:
|
||||
from .audit_policy import (
|
||||
account_name,
|
||||
read_deduplication_fingerprint,
|
||||
read_deduplication_fingerprint_values,
|
||||
deduplication_fingerprint,
|
||||
deduplication_fingerprint_values,
|
||||
skipped_user_suffixes,
|
||||
)
|
||||
from .state_db import connect_state_db
|
||||
except ImportError:
|
||||
from audit_policy import (
|
||||
account_name,
|
||||
read_deduplication_fingerprint,
|
||||
read_deduplication_fingerprint_values,
|
||||
deduplication_fingerprint,
|
||||
deduplication_fingerprint_values,
|
||||
skipped_user_suffixes,
|
||||
)
|
||||
from state_db import connect_state_db
|
||||
@@ -49,13 +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 (
|
||||
CREATE TABLE IF NOT EXISTS audit_event_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_event_dedup_time
|
||||
ON audit_event_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;
|
||||
@@ -147,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 4
|
||||
policy_version INTEGER NOT NULL DEFAULT 5
|
||||
);
|
||||
|
||||
"""
|
||||
@@ -174,6 +174,7 @@ VALUES (?, ?, ?, ?)
|
||||
|
||||
|
||||
def ensure_audit_schema(conn: sqlite3.Connection) -> None:
|
||||
conn.execute("DROP TABLE IF EXISTS audit_read_dedup")
|
||||
conn.executescript(AUDIT_SCHEMA)
|
||||
for legacy_index in (
|
||||
"audit_events_time",
|
||||
@@ -221,7 +222,8 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
|
||||
WHERE singleton = 1
|
||||
"""
|
||||
).fetchone()
|
||||
if state is not None and int(state["policy_version"]) < 4:
|
||||
if state is not None and int(state["policy_version"]) < 5:
|
||||
conn.execute("DELETE FROM audit_event_dedup")
|
||||
for table in (
|
||||
"audit_daily_totals",
|
||||
"audit_daily_counts",
|
||||
@@ -232,7 +234,7 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
|
||||
"""
|
||||
UPDATE audit_rollup_state
|
||||
SET backfill_next_id = 1, backfill_max_id = ?, ready = ?,
|
||||
policy_version = 4
|
||||
policy_version = 5
|
||||
WHERE singleton = 1
|
||||
""",
|
||||
(max_id, int(max_id == 0)),
|
||||
@@ -240,17 +242,21 @@ def ensure_audit_schema(conn: sqlite3.Connection) -> None:
|
||||
conn.commit()
|
||||
|
||||
|
||||
READ_DEDUP_DEFAULT_WINDOW_SECONDS = 2 * 86400
|
||||
DEDUP_DEFAULT_WINDOW_SECONDS = 2 * 86400
|
||||
|
||||
|
||||
def read_deduplication_cutoff() -> int:
|
||||
try:
|
||||
window = int(os.getenv(
|
||||
def deduplication_cutoff() -> int:
|
||||
configured = os.getenv(
|
||||
"AUDIT_DEDUP_WINDOW_SECONDS",
|
||||
os.getenv(
|
||||
"AUDIT_READ_DEDUP_WINDOW_SECONDS",
|
||||
str(READ_DEDUP_DEFAULT_WINDOW_SECONDS),
|
||||
))
|
||||
str(DEDUP_DEFAULT_WINDOW_SECONDS),
|
||||
),
|
||||
)
|
||||
try:
|
||||
window = int(configured)
|
||||
except ValueError:
|
||||
window = READ_DEDUP_DEFAULT_WINDOW_SECONDS
|
||||
window = DEDUP_DEFAULT_WINDOW_SECONDS
|
||||
now = int(dt.datetime.now(dt.timezone.utc).timestamp())
|
||||
return now - max(1, window)
|
||||
|
||||
@@ -285,8 +291,8 @@ 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,
|
||||
"audit_event_fingerprint", 8,
|
||||
deduplication_fingerprint_values, deterministic=True,
|
||||
)
|
||||
ensure_audit_schema(self.conn)
|
||||
|
||||
@@ -314,12 +320,12 @@ class AuditStore:
|
||||
totals = Counter()
|
||||
counts = Counter()
|
||||
facets = set()
|
||||
seen_read_fingerprints = set()
|
||||
seen_fingerprints = set()
|
||||
for event in events:
|
||||
read_fingerprint = read_deduplication_fingerprint(event)
|
||||
fingerprint = deduplication_fingerprint(event)
|
||||
if (
|
||||
read_fingerprint is not None
|
||||
and read_fingerprint in seen_read_fingerprints
|
||||
fingerprint is not None
|
||||
and fingerprint in seen_fingerprints
|
||||
):
|
||||
continue
|
||||
occurred_second = parse_event_second(event["timestamp"])
|
||||
@@ -344,11 +350,11 @@ class AuditStore:
|
||||
success,
|
||||
str(event["path"]),
|
||||
str(event["source"]),
|
||||
read_fingerprint,
|
||||
fingerprint,
|
||||
)
|
||||
)
|
||||
if read_fingerprint is not None:
|
||||
seen_read_fingerprints.add(read_fingerprint)
|
||||
if fingerprint is not None:
|
||||
seen_fingerprints.add(fingerprint)
|
||||
|
||||
self.conn.execute("BEGIN IMMEDIATE")
|
||||
try:
|
||||
@@ -364,7 +370,7 @@ class AuditStore:
|
||||
for row in self.conn.execute(
|
||||
f"""
|
||||
SELECT fingerprint
|
||||
FROM audit_read_dedup
|
||||
FROM audit_event_dedup
|
||||
WHERE fingerprint IN ({placeholders})
|
||||
""",
|
||||
chunk,
|
||||
@@ -453,7 +459,7 @@ class AuditStore:
|
||||
first_inserted_id = last_inserted_id - len(event_rows) + 1
|
||||
self.conn.executemany(
|
||||
"""
|
||||
INSERT INTO audit_read_dedup (
|
||||
INSERT INTO audit_event_dedup (
|
||||
fingerprint, occurred_second, event_id
|
||||
) VALUES (?, ?, ?)
|
||||
""",
|
||||
@@ -474,8 +480,8 @@ class AuditStore:
|
||||
self.conn.executemany(ROLLUP_FACET_SQL, facets)
|
||||
|
||||
self.conn.execute(
|
||||
"DELETE FROM audit_read_dedup WHERE occurred_second < ?",
|
||||
(read_deduplication_cutoff(),),
|
||||
"DELETE FROM audit_event_dedup WHERE occurred_second < ?",
|
||||
(deduplication_cutoff(),),
|
||||
)
|
||||
|
||||
if source_updates:
|
||||
@@ -544,53 +550,53 @@ class AuditStore:
|
||||
end_id = int(chunk[1])
|
||||
self.conn.execute("BEGIN IMMEDIATE")
|
||||
try:
|
||||
dedup_cutoff = read_deduplication_cutoff()
|
||||
dedup_cutoff = deduplication_cutoff()
|
||||
self.conn.execute(
|
||||
"DELETE FROM audit_read_dedup WHERE occurred_second < ?",
|
||||
"DELETE FROM audit_event_dedup WHERE occurred_second < ?",
|
||||
(dedup_cutoff,),
|
||||
)
|
||||
self.conn.execute(
|
||||
"""
|
||||
INSERT INTO audit_read_dedup (
|
||||
INSERT INTO audit_event_dedup (
|
||||
fingerprint, occurred_second, event_id
|
||||
)
|
||||
SELECT audit_read_fingerprint(
|
||||
e.occurred_at, e.user, e.client_ip, e.share, e.path,
|
||||
SELECT audit_event_fingerprint(
|
||||
e.occurred_at, e.action, 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.action = 'read' OR e.share = 'FSLogix' COLLATE NOCASE)
|
||||
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
|
||||
audit_event_dedup.event_id, excluded.event_id
|
||||
)
|
||||
""",
|
||||
(start_id, end_id, dedup_cutoff),
|
||||
)
|
||||
duplicate_read_sql = """
|
||||
duplicate_event_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,
|
||||
JOIN audit_event_dedup AS d
|
||||
ON d.fingerprint = audit_event_fingerprint(
|
||||
e.occurred_at, e.action, e.user, e.client_ip, e.share, e.path,
|
||||
e.success, e.result
|
||||
)
|
||||
WHERE e.id BETWEEN ? AND ?
|
||||
AND e.action = 'read'
|
||||
AND (e.action = 'read' OR e.share = 'FSLogix' COLLATE NOCASE)
|
||||
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})",
|
||||
f"WHERE event_id IN ({duplicate_event_sql})",
|
||||
(start_id, end_id, dedup_cutoff),
|
||||
)
|
||||
self.conn.execute(
|
||||
f"DELETE FROM audit_events "
|
||||
f"WHERE id IN ({duplicate_read_sql})",
|
||||
f"WHERE id IN ({duplicate_event_sql})",
|
||||
(start_id, end_id, dedup_cutoff),
|
||||
)
|
||||
self.conn.execute(
|
||||
|
||||
+9
-1
@@ -32,7 +32,7 @@ require_vfs_modules() {
|
||||
|
||||
local missing=0
|
||||
local module
|
||||
for module in acl_xattr full_audit; do
|
||||
for module in acl_xattr recycle full_audit; do
|
||||
if [[ ! -f "$modules_dir/vfs/${module}.so" ]]; then
|
||||
printf '[init] ERROR: missing VFS module %s at %s/vfs/%s.so\n' "$module" "$modules_dir" "$module" >&2
|
||||
missing=1
|
||||
@@ -251,6 +251,10 @@ write_runtime_env_file() {
|
||||
if [[ -n "${LDAP_BASE_DN:-}" ]]; then
|
||||
printf 'export LDAP_BASE_DN=%q\n' "$LDAP_BASE_DN"
|
||||
fi
|
||||
printf 'export TRASH_RETENTION_DAYS=%q\n' "${TRASH_RETENTION_DAYS:-7}"
|
||||
if [[ -n "${GROUP_ROOT:-}" ]]; then printf 'export GROUP_ROOT=%q\n' "$GROUP_ROOT"; fi
|
||||
if [[ -n "${PRIVATE_ROOT:-}" ]]; then printf 'export PRIVATE_ROOT=%q\n' "$PRIVATE_ROOT"; fi
|
||||
if [[ -n "${FSLOGIX_ROOT:-}" ]]; then printf 'export FSLOGIX_ROOT=%q\n' "$FSLOGIX_ROOT"; fi
|
||||
} > /app/runtime.env
|
||||
chmod 600 /app/runtime.env
|
||||
}
|
||||
@@ -567,6 +571,7 @@ install_cron_job() {
|
||||
SHELL=/bin/bash
|
||||
PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin
|
||||
*/5 * * * * root source /app/runtime.env && RECONCILE_TRIGGER=automatic /usr/bin/python3 /app/reconcile_shares.py
|
||||
17 * * * * root source /app/runtime.env && /usr/bin/python3 /app/trash.py --cleanup
|
||||
EOF
|
||||
|
||||
if [[ -n "${BACKUP_DESTINATION:-}" ]] && env_is_true "${BACKUP_AUTO_ENABLED:-true}"; then
|
||||
@@ -622,6 +627,9 @@ write_runtime_env_file
|
||||
log 'Running startup reconciliation'
|
||||
RECONCILE_TRIGGER=startup python3 /app/reconcile_shares.py
|
||||
|
||||
log 'Preparing seven-day trash repositories'
|
||||
python3 /app/trash.py --cleanup
|
||||
|
||||
start_observability_services
|
||||
|
||||
if [[ -n "${BACKUP_DESTINATION:-}" ]] && env_is_true "${BACKUP_AUTO_ENABLED:-true}"; then
|
||||
|
||||
+416
@@ -0,0 +1,416 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Safe listing, expiry, download, and restoration for Samba recycle bins."""
|
||||
|
||||
import argparse
|
||||
import base64
|
||||
import datetime as dt
|
||||
import os
|
||||
import re
|
||||
import stat
|
||||
from typing import BinaryIO, Dict, List, Optional, Tuple
|
||||
|
||||
|
||||
TRASH_DIRECTORY = ".trash"
|
||||
DEFAULT_RETENTION_DAYS = 7
|
||||
VERSION_PREFIX_RE = re.compile(r"^Copy #\d+ of ", re.IGNORECASE)
|
||||
|
||||
|
||||
def retention_days() -> int:
|
||||
try:
|
||||
value = int(os.getenv("TRASH_RETENTION_DAYS", str(DEFAULT_RETENTION_DAYS)))
|
||||
except ValueError:
|
||||
value = DEFAULT_RETENTION_DAYS
|
||||
return max(1, min(365, value))
|
||||
|
||||
|
||||
def share_roots() -> Dict[str, str]:
|
||||
return {
|
||||
"Data": os.path.abspath(os.getenv("GROUP_ROOT", "/data/groups/data")),
|
||||
"Private": os.path.abspath(os.getenv("PRIVATE_ROOT", "/data/private")),
|
||||
"FSLogix": os.path.abspath(os.getenv("FSLOGIX_ROOT", "/data/fslogix")),
|
||||
}
|
||||
|
||||
|
||||
def trash_root(share_root: str) -> str:
|
||||
return os.path.join(share_root, TRASH_DIRECTORY)
|
||||
|
||||
|
||||
def ensure_trash_roots() -> None:
|
||||
"""Create non-listable sticky repositories users can write through Samba."""
|
||||
for root in share_roots().values():
|
||||
os.makedirs(root, exist_ok=True)
|
||||
repository = trash_root(root)
|
||||
os.makedirs(repository, exist_ok=True)
|
||||
try:
|
||||
os.chown(repository, 0, 0)
|
||||
except PermissionError:
|
||||
if os.geteuid() == 0:
|
||||
raise
|
||||
os.chmod(repository, 0o1733)
|
||||
|
||||
|
||||
def _share_name(value: str) -> str:
|
||||
for name in share_roots():
|
||||
if name.casefold() == value.casefold():
|
||||
return name
|
||||
raise ValueError("Unbekannte Freigabe")
|
||||
|
||||
|
||||
def _relative_parts(value: str) -> List[str]:
|
||||
if not value or value.startswith(("/", "\\")) or "\0" in value:
|
||||
raise ValueError("Ungültiger Papierkorbpfad")
|
||||
parts = value.split("/")
|
||||
if any(part in {"", ".", ".."} for part in parts):
|
||||
raise ValueError("Ungültiger Papierkorbpfad")
|
||||
return parts
|
||||
|
||||
|
||||
def encode_item_id(share: str, relative_path: str) -> str:
|
||||
canonical_share = _share_name(share)
|
||||
_relative_parts(relative_path)
|
||||
payload = f"{canonical_share}\0{relative_path}".encode("utf-8")
|
||||
return base64.urlsafe_b64encode(payload).decode("ascii").rstrip("=")
|
||||
|
||||
|
||||
def decode_item_id(item_id: str) -> Tuple[str, str]:
|
||||
if not item_id or len(item_id) > 8192:
|
||||
raise ValueError("Ungültige Papierkorb-ID")
|
||||
try:
|
||||
padding = "=" * (-len(item_id) % 4)
|
||||
payload = base64.b64decode(
|
||||
item_id + padding,
|
||||
altchars=b"-_",
|
||||
validate=True,
|
||||
).decode("utf-8")
|
||||
share, relative_path = payload.split("\0", 1)
|
||||
except (ValueError, UnicodeDecodeError) as exc:
|
||||
raise ValueError("Ungültige Papierkorb-ID") from exc
|
||||
canonical_share = _share_name(share)
|
||||
_relative_parts(relative_path)
|
||||
return canonical_share, relative_path
|
||||
|
||||
|
||||
def _original_relative(relative_path: str) -> str:
|
||||
parts = _relative_parts(relative_path)
|
||||
if len(parts) < 2:
|
||||
raise ValueError("Papierkorbeintrag enthält keinen Originalpfad")
|
||||
original = parts[1:]
|
||||
original[-1] = VERSION_PREFIX_RE.sub("", original[-1], count=1)
|
||||
if not original[-1]:
|
||||
raise ValueError("Papierkorbeintrag enthält keinen Dateinamen")
|
||||
return "/".join(original)
|
||||
|
||||
|
||||
def _open_directory_chain(
|
||||
root: str,
|
||||
parts: List[str],
|
||||
*,
|
||||
create: bool = False,
|
||||
uid: int = 0,
|
||||
gid: int = 0,
|
||||
mode: int = 0o700,
|
||||
) -> int:
|
||||
flags = os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW
|
||||
current_fd = os.open(root, flags)
|
||||
try:
|
||||
for part in parts:
|
||||
try:
|
||||
next_fd = os.open(part, flags, dir_fd=current_fd)
|
||||
except FileNotFoundError:
|
||||
if not create:
|
||||
raise
|
||||
os.mkdir(part, mode=mode, dir_fd=current_fd)
|
||||
os.chown(
|
||||
part,
|
||||
uid,
|
||||
gid,
|
||||
dir_fd=current_fd,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
os.chmod(
|
||||
part,
|
||||
mode,
|
||||
dir_fd=current_fd,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
next_fd = os.open(part, flags, dir_fd=current_fd)
|
||||
os.close(current_fd)
|
||||
current_fd = next_fd
|
||||
return current_fd
|
||||
except Exception:
|
||||
os.close(current_fd)
|
||||
raise
|
||||
|
||||
|
||||
def _resolved_item(item_id: str) -> Tuple[str, str, str, List[str], os.stat_result]:
|
||||
share, relative_path = decode_item_id(item_id)
|
||||
share_root = share_roots()[share]
|
||||
repository = trash_root(share_root)
|
||||
parts = _relative_parts(relative_path)
|
||||
parent_fd = _open_directory_chain(repository, parts[:-1])
|
||||
try:
|
||||
info = os.stat(parts[-1], dir_fd=parent_fd, follow_symlinks=False)
|
||||
finally:
|
||||
os.close(parent_fd)
|
||||
if not stat.S_ISREG(info.st_mode):
|
||||
raise ValueError("Nur reguläre Dateien können verarbeitet werden")
|
||||
cutoff = dt.datetime.now(dt.timezone.utc).timestamp() - retention_days() * 86400
|
||||
if info.st_mtime < cutoff:
|
||||
raise FileNotFoundError("Papierkorbeintrag ist abgelaufen")
|
||||
return share, share_root, relative_path, parts, info
|
||||
|
||||
|
||||
def _iso_timestamp(seconds: float) -> str:
|
||||
return dt.datetime.fromtimestamp(seconds, dt.timezone.utc).isoformat(
|
||||
timespec="seconds"
|
||||
)
|
||||
|
||||
|
||||
def _item_payload(
|
||||
share: str,
|
||||
relative_path: str,
|
||||
info: os.stat_result,
|
||||
days: int,
|
||||
) -> Dict[str, object]:
|
||||
original = _original_relative(relative_path)
|
||||
deleted_at = float(info.st_mtime)
|
||||
return {
|
||||
"id": encode_item_id(share, relative_path),
|
||||
"share": share,
|
||||
"path": original,
|
||||
"name": original.rsplit("/", 1)[-1],
|
||||
"deletedBy": relative_path.split("/", 1)[0],
|
||||
"deletedAt": _iso_timestamp(deleted_at),
|
||||
"expiresAt": _iso_timestamp(deleted_at + days * 86400),
|
||||
"size": int(info.st_size),
|
||||
}
|
||||
|
||||
|
||||
def list_items(
|
||||
*,
|
||||
share: str = "",
|
||||
path: str = "",
|
||||
limit: int = 200,
|
||||
now: Optional[dt.datetime] = None,
|
||||
) -> Dict[str, object]:
|
||||
days = retention_days()
|
||||
current = now or dt.datetime.now(dt.timezone.utc)
|
||||
cutoff = current.timestamp() - days * 86400
|
||||
selected_share = _share_name(share) if share else ""
|
||||
path_filter = path.strip().casefold()
|
||||
items: List[Dict[str, object]] = []
|
||||
|
||||
for share_name, root in share_roots().items():
|
||||
if selected_share and share_name != selected_share:
|
||||
continue
|
||||
repository = trash_root(root)
|
||||
try:
|
||||
walker = os.walk(repository, topdown=True, followlinks=False)
|
||||
for directory, subdirectories, filenames in walker:
|
||||
safe_subdirectories = []
|
||||
for name in subdirectories:
|
||||
candidate = os.path.join(directory, name)
|
||||
try:
|
||||
if not stat.S_ISLNK(os.lstat(candidate).st_mode):
|
||||
safe_subdirectories.append(name)
|
||||
except OSError:
|
||||
continue
|
||||
subdirectories[:] = safe_subdirectories
|
||||
for filename in filenames:
|
||||
candidate = os.path.join(directory, filename)
|
||||
try:
|
||||
info = os.lstat(candidate)
|
||||
except OSError:
|
||||
continue
|
||||
if not stat.S_ISREG(info.st_mode) or info.st_mtime < cutoff:
|
||||
continue
|
||||
relative_path = os.path.relpath(candidate, repository).replace(
|
||||
os.sep, "/"
|
||||
)
|
||||
try:
|
||||
payload = _item_payload(
|
||||
share_name, relative_path, info, days
|
||||
)
|
||||
except ValueError:
|
||||
continue
|
||||
if path_filter and path_filter not in str(payload["path"]).casefold():
|
||||
continue
|
||||
items.append(payload)
|
||||
except OSError:
|
||||
continue
|
||||
|
||||
items.sort(
|
||||
key=lambda item: (str(item["deletedAt"]), str(item["id"])),
|
||||
reverse=True,
|
||||
)
|
||||
maximum = max(1, min(500, int(limit)))
|
||||
return {
|
||||
"items": items[:maximum],
|
||||
"matched": len(items),
|
||||
"truncated": len(items) > maximum,
|
||||
"retentionDays": days,
|
||||
"scannedAt": current.isoformat(timespec="seconds"),
|
||||
}
|
||||
|
||||
|
||||
def open_download(item_id: str) -> Tuple[BinaryIO, Dict[str, object]]:
|
||||
share, _share_root, relative_path, parts, _info = _resolved_item(item_id)
|
||||
repository = trash_root(share_roots()[share])
|
||||
parent_fd = _open_directory_chain(repository, parts[:-1])
|
||||
try:
|
||||
file_fd = os.open(parts[-1], os.O_RDONLY | os.O_NOFOLLOW, dir_fd=parent_fd)
|
||||
finally:
|
||||
os.close(parent_fd)
|
||||
try:
|
||||
info = os.fstat(file_fd)
|
||||
if not stat.S_ISREG(info.st_mode):
|
||||
raise ValueError("Nur reguläre Dateien können heruntergeladen werden")
|
||||
payload = _item_payload(share, relative_path, info, retention_days())
|
||||
return os.fdopen(file_fd, "rb"), payload
|
||||
except Exception:
|
||||
os.close(file_fd)
|
||||
raise
|
||||
|
||||
|
||||
def _restore_directory_mode(share: str) -> int:
|
||||
if share == "Data":
|
||||
return 0o2770
|
||||
return 0o700
|
||||
|
||||
|
||||
def _remove_empty_trash_parents(repository: str, relative_path: str) -> None:
|
||||
current = os.path.dirname(os.path.join(repository, relative_path))
|
||||
repository = os.path.abspath(repository)
|
||||
while (
|
||||
os.path.commonpath((repository, current)) == repository
|
||||
and current != repository
|
||||
):
|
||||
try:
|
||||
os.rmdir(current)
|
||||
except OSError:
|
||||
break
|
||||
current = os.path.dirname(current)
|
||||
|
||||
|
||||
def restore_item(item_id: str) -> Dict[str, object]:
|
||||
share, share_root, relative_path, parts, info = _resolved_item(item_id)
|
||||
original = _original_relative(relative_path)
|
||||
original_parts = _relative_parts(original)
|
||||
if original_parts[0] == TRASH_DIRECTORY:
|
||||
raise ValueError("Ungültiger Wiederherstellungspfad")
|
||||
|
||||
source_parent_fd = _open_directory_chain(
|
||||
trash_root(share_root), parts[:-1]
|
||||
)
|
||||
destination_parent_fd = _open_directory_chain(
|
||||
share_root,
|
||||
original_parts[:-1],
|
||||
create=True,
|
||||
uid=info.st_uid,
|
||||
gid=info.st_gid,
|
||||
mode=_restore_directory_mode(share),
|
||||
)
|
||||
linked = False
|
||||
try:
|
||||
os.link(
|
||||
parts[-1],
|
||||
original_parts[-1],
|
||||
src_dir_fd=source_parent_fd,
|
||||
dst_dir_fd=destination_parent_fd,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
linked = True
|
||||
linked_info = os.stat(
|
||||
original_parts[-1],
|
||||
dir_fd=destination_parent_fd,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
if (
|
||||
not stat.S_ISREG(linked_info.st_mode)
|
||||
or linked_info.st_dev != info.st_dev
|
||||
or linked_info.st_ino != info.st_ino
|
||||
):
|
||||
os.unlink(original_parts[-1], dir_fd=destination_parent_fd)
|
||||
linked = False
|
||||
raise RuntimeError("Papierkorbeintrag wurde währenddessen verändert")
|
||||
try:
|
||||
os.unlink(parts[-1], dir_fd=source_parent_fd)
|
||||
except Exception:
|
||||
os.unlink(original_parts[-1], dir_fd=destination_parent_fd)
|
||||
linked = False
|
||||
raise
|
||||
finally:
|
||||
os.close(source_parent_fd)
|
||||
os.close(destination_parent_fd)
|
||||
|
||||
if not linked:
|
||||
raise RuntimeError("Wiederherstellung konnte nicht abgeschlossen werden")
|
||||
_remove_empty_trash_parents(trash_root(share_root), relative_path)
|
||||
return {
|
||||
"restored": True,
|
||||
"share": share,
|
||||
"path": original,
|
||||
}
|
||||
|
||||
|
||||
def cleanup_expired(now: Optional[dt.datetime] = None) -> Dict[str, int]:
|
||||
days = retention_days()
|
||||
current = now or dt.datetime.now(dt.timezone.utc)
|
||||
cutoff = current.timestamp() - days * 86400
|
||||
removed = 0
|
||||
removed_bytes = 0
|
||||
|
||||
ensure_trash_roots()
|
||||
for root in share_roots().values():
|
||||
repository = trash_root(root)
|
||||
for directory, subdirectories, filenames in os.walk(
|
||||
repository, topdown=False, followlinks=False
|
||||
):
|
||||
for filename in filenames:
|
||||
candidate = os.path.join(directory, filename)
|
||||
try:
|
||||
info = os.lstat(candidate)
|
||||
if info.st_mtime >= cutoff:
|
||||
continue
|
||||
if not (
|
||||
stat.S_ISREG(info.st_mode) or stat.S_ISLNK(info.st_mode)
|
||||
):
|
||||
continue
|
||||
os.unlink(candidate)
|
||||
removed += 1
|
||||
if stat.S_ISREG(info.st_mode):
|
||||
removed_bytes += int(info.st_size)
|
||||
except OSError:
|
||||
continue
|
||||
for name in subdirectories:
|
||||
candidate = os.path.join(directory, name)
|
||||
try:
|
||||
info = os.lstat(candidate)
|
||||
if stat.S_ISLNK(info.st_mode):
|
||||
if info.st_mtime < cutoff:
|
||||
os.unlink(candidate)
|
||||
removed += 1
|
||||
continue
|
||||
os.rmdir(candidate)
|
||||
except OSError:
|
||||
continue
|
||||
return {"removed": removed, "removedBytes": removed_bytes}
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser()
|
||||
parser.add_argument("--cleanup", action="store_true")
|
||||
args = parser.parse_args()
|
||||
if not args.cleanup:
|
||||
parser.error("--cleanup is required")
|
||||
result = cleanup_expired()
|
||||
print(
|
||||
f"[trash] Removed {result['removed']} expired item(s) "
|
||||
f"({result['removedBytes']} bytes)",
|
||||
flush=True,
|
||||
)
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -105,6 +105,7 @@ function routeFor(path) {
|
||||
if (path.startsWith("/reconciliation")) return "reconciliation";
|
||||
if (path.startsWith("/activity/fslogix")) return "activity-fslogix";
|
||||
if (path.startsWith("/activity")) return "activity";
|
||||
if (path.startsWith("/trash")) return "trash";
|
||||
if (path.startsWith("/backup")) return "backup";
|
||||
if (path.startsWith("/report")) return "report";
|
||||
if (path.startsWith("/system")) return "system";
|
||||
@@ -127,6 +128,7 @@ async function navigate(path, replace = false) {
|
||||
if (route === "storage-users") await renderStorage("users");
|
||||
if (route === "activity") await renderActivity();
|
||||
if (route === "activity-fslogix") await renderActivity("fslogix");
|
||||
if (route === "trash") await renderTrash();
|
||||
if (route === "backup") await renderBackup();
|
||||
if (route === "report") await renderReport();
|
||||
if (route === "system") await renderSystem();
|
||||
@@ -303,6 +305,68 @@ async function renderActivity(stream = "main") {
|
||||
await load(false);
|
||||
}
|
||||
|
||||
async function renderTrash() {
|
||||
content.innerHTML = pageHead(
|
||||
"Papierkorb",
|
||||
"Gelöschte Dateien sieben Tage lang herunterladen oder am Originalpfad wiederherstellen.",
|
||||
) + `
|
||||
<section class="panel">
|
||||
<form id="trash-filter" class="filters">
|
||||
<label>Freigabe<select name="share"><option value="">Alle Freigaben</option><option>Data</option><option>Private</option><option>FSLogix</option></select></label>
|
||||
<label class="wide">Pfad enthält<input name="path" placeholder="Ordner oder Dateiname"></label>
|
||||
<button type="submit">Aktualisieren</button>
|
||||
</form>
|
||||
<p id="trash-summary" class="muted"></p>
|
||||
<div class="table-wrap"><table><thead><tr><th>Gelöscht (UTC)</th><th>Benutzer</th><th>Freigabe</th><th>Originalpfad</th><th>Größe</th><th>Verfügbar bis</th><th>Aktionen</th></tr></thead><tbody id="trash-rows"></tbody></table></div>
|
||||
</section>`;
|
||||
const form = document.querySelector("#trash-filter");
|
||||
const load = async () => {
|
||||
const params = new URLSearchParams(new FormData(form));
|
||||
params.set("limit", "500");
|
||||
const result = await api(`/api/trash?${params}`);
|
||||
const body = document.querySelector("#trash-rows");
|
||||
if (!result.items.length) {
|
||||
body.innerHTML = '<tr><td colspan="7" class="empty">Keine gelöschten Dateien</td></tr>';
|
||||
} else {
|
||||
body.innerHTML = result.items.map(item => `<tr>
|
||||
<td class="timestamp">${esc(utcTime(item.deletedAt))}</td>
|
||||
<td>${esc(item.deletedBy)}</td>
|
||||
<td>${esc(item.share)}</td>
|
||||
<td class="path">${esc(item.path)}</td>
|
||||
<td class="numeric">${bytes(item.size)}</td>
|
||||
<td class="timestamp">${esc(utcTime(item.expiresAt))}</td>
|
||||
<td><div class="row-actions"><a class="button" href="/api/trash/download?id=${encodeURIComponent(item.id)}" download>Herunterladen</a><button type="button" data-restore="${esc(item.id)}">Wiederherstellen</button></div></td>
|
||||
</tr>`).join("");
|
||||
}
|
||||
const shown = result.items.length.toLocaleString("de-DE");
|
||||
const total = Number(result.matched || 0).toLocaleString("de-DE");
|
||||
const suffix = result.truncated ? ` · ${shown} von ${total} angezeigt` : "";
|
||||
document.querySelector("#trash-summary").textContent = `${total} Dateien · Aufbewahrung ${result.retentionDays} Tage${suffix}`;
|
||||
};
|
||||
form.addEventListener("submit", async event => {
|
||||
event.preventDefault();
|
||||
try { await load(); } catch (error) { notice(error.message); }
|
||||
});
|
||||
document.querySelector("#trash-rows").addEventListener("click", async event => {
|
||||
const button = event.target.closest("button[data-restore]");
|
||||
if (!button) return;
|
||||
if (!window.confirm("Datei am Originalpfad wiederherstellen?")) return;
|
||||
button.disabled = true;
|
||||
try {
|
||||
const result = await api("/api/trash/restore", {
|
||||
method: "POST",
|
||||
body: JSON.stringify({id: button.dataset.restore}),
|
||||
});
|
||||
notice(`${result.share}: ${result.path} wurde wiederhergestellt.`);
|
||||
await load();
|
||||
} catch (error) {
|
||||
notice(error.message);
|
||||
button.disabled = false;
|
||||
}
|
||||
});
|
||||
await load();
|
||||
}
|
||||
|
||||
function triggerLabel(value) {
|
||||
return ({automatic: "Automatisch", web: "Weboberfläche", manual: "Befehlszeile", startup: "Systemstart"}[value] || value || "—");
|
||||
}
|
||||
|
||||
@@ -33,6 +33,7 @@
|
||||
<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="/trash" data-route="trash">Papierkorb</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>
|
||||
|
||||
+2
-1
@@ -8,7 +8,7 @@
|
||||
body { margin: 0; min-height: 100vh; }
|
||||
a { color: #0645ad; }
|
||||
button, input, select { font: inherit; }
|
||||
button { padding: .4rem .7rem; border: 1px solid #777; color: #111; background: #eee; cursor: pointer; }
|
||||
button, .button { display: inline-block; padding: .4rem .7rem; border: 1px solid #777; color: #111; background: #eee; cursor: pointer; text-decoration: none; white-space: nowrap; }
|
||||
button:disabled { color: #777; cursor: wait; }
|
||||
input, select { width: 100%; padding: .4rem; border: 1px solid #999; background: #fff; }
|
||||
button:focus, input:focus, select:focus, a:focus { outline: 2px solid #0645ad; outline-offset: 1px; }
|
||||
@@ -64,6 +64,7 @@ th, td { padding: .5rem; border-bottom: 1px solid #ccc; vertical-align: top; }
|
||||
.badge.warn { color: #750; }
|
||||
.toolbar { display: flex; flex-wrap: wrap; gap: .5rem; }
|
||||
.toolbar input { max-width: 340px; }
|
||||
.row-actions { display: flex; flex-wrap: wrap; gap: .4rem; }
|
||||
.list { margin: 0; padding: 0; list-style: none; }
|
||||
.select-row { width: 100%; display: grid; grid-template-columns: 1fr auto; gap: .5rem; border: 0; border-bottom: 1px solid #ccc; background: #fff; text-align: left; }
|
||||
.select-row.active { font-weight: bold; background: #eee; }
|
||||
|
||||
@@ -27,8 +27,10 @@ from typing import Dict, List, Optional, Tuple
|
||||
|
||||
try:
|
||||
from app import reconcile_shares as directory
|
||||
from app import trash
|
||||
except ImportError: # Container execution uses /app as the import root.
|
||||
import reconcile_shares as directory
|
||||
import trash
|
||||
|
||||
try:
|
||||
from app.audit_store import (
|
||||
@@ -460,6 +462,8 @@ def scan_children(root: str) -> List[Dict[str, object]]:
|
||||
except OSError:
|
||||
return rows
|
||||
for entry in entries:
|
||||
if entry.name == trash.TRASH_DIRECTORY:
|
||||
continue
|
||||
try:
|
||||
if not entry.is_dir(follow_symlinks=False):
|
||||
continue
|
||||
@@ -1032,6 +1036,43 @@ class Handler(BaseHTTPRequestHandler):
|
||||
raise ValueError("Ein JSON-Objekt ist erforderlich")
|
||||
return value
|
||||
|
||||
def send_trash_download(self, params: Dict[str, List[str]]) -> None:
|
||||
item_id = params.get("id", [""])[0]
|
||||
try:
|
||||
handle, item = trash.open_download(item_id)
|
||||
except FileNotFoundError:
|
||||
self.send_error_json(HTTPStatus.NOT_FOUND, "Datei nicht gefunden oder abgelaufen")
|
||||
return
|
||||
except ValueError as exc:
|
||||
self.send_error_json(HTTPStatus.BAD_REQUEST, str(exc))
|
||||
return
|
||||
except OSError as exc:
|
||||
log(f"Trash download failed: {exc}")
|
||||
self.send_error_json(HTTPStatus.INTERNAL_SERVER_ERROR, "Download konnte nicht geöffnet werden")
|
||||
return
|
||||
|
||||
filename = str(item["name"])
|
||||
encoded_name = urllib.parse.quote(filename, safe="")
|
||||
with handle:
|
||||
self.send_response(HTTPStatus.OK)
|
||||
self.security_headers()
|
||||
self.send_header("Content-Type", "application/octet-stream")
|
||||
self.send_header("Content-Length", str(item["size"]))
|
||||
self.send_header(
|
||||
"Content-Disposition",
|
||||
f"attachment; filename*=UTF-8{chr(39) * 2}{encoded_name}",
|
||||
)
|
||||
self.end_headers()
|
||||
try:
|
||||
while True:
|
||||
chunk = handle.read(1024 * 1024)
|
||||
if not chunk:
|
||||
break
|
||||
self.wfile.write(chunk)
|
||||
except (BrokenPipeError, ConnectionResetError):
|
||||
pass
|
||||
|
||||
|
||||
def do_POST(self) -> None: # pylint: disable=invalid-name
|
||||
parsed = urllib.parse.urlparse(self.path)
|
||||
if parsed.path == "/api/login":
|
||||
@@ -1060,6 +1101,42 @@ class Handler(BaseHTTPRequestHandler):
|
||||
cookie = f"{JWT_COOKIE}=; Path=/; Max-Age=0; HttpOnly; Secure; SameSite=Strict"
|
||||
self.send_json({"ok": True}, cookie=cookie)
|
||||
return
|
||||
if parsed.path == "/api/trash/restore":
|
||||
user = self.require_user()
|
||||
if user is None:
|
||||
return
|
||||
try:
|
||||
body = self.read_json_body()
|
||||
result = trash.restore_item(str(body.get("id", "")))
|
||||
except FileExistsError:
|
||||
self.send_error_json(
|
||||
HTTPStatus.CONFLICT,
|
||||
"Am Originalpfad existiert bereits eine Datei",
|
||||
)
|
||||
return
|
||||
except FileNotFoundError:
|
||||
self.send_error_json(
|
||||
HTTPStatus.NOT_FOUND,
|
||||
"Datei nicht gefunden oder abgelaufen",
|
||||
)
|
||||
return
|
||||
except ValueError as exc:
|
||||
self.send_error_json(HTTPStatus.BAD_REQUEST, str(exc))
|
||||
return
|
||||
except (OSError, RuntimeError) as exc:
|
||||
log(f"Trash restore failed: {exc}")
|
||||
self.send_error_json(
|
||||
HTTPStatus.INTERNAL_SERVER_ERROR,
|
||||
"Datei konnte nicht wiederhergestellt werden",
|
||||
)
|
||||
return
|
||||
log(
|
||||
f"{user['sub']} restored "
|
||||
f"{result['share']}:{result['path']}"
|
||||
)
|
||||
self.send_json(result)
|
||||
return
|
||||
|
||||
if parsed.path.startswith("/api/actions/"):
|
||||
user = self.require_user()
|
||||
if user is None:
|
||||
@@ -1105,6 +1182,16 @@ class Handler(BaseHTTPRequestHandler):
|
||||
self.send_json(query_audit(params))
|
||||
elif path == "/api/fslogix-activity":
|
||||
self.send_json(query_audit(params, stream="fslogix"))
|
||||
elif path == "/api/trash/download":
|
||||
self.send_trash_download(params)
|
||||
elif path == "/api/trash":
|
||||
self.send_json(
|
||||
trash.list_items(
|
||||
share=params.get("share", [""])[0],
|
||||
path=params.get("path", [""])[0],
|
||||
limit=int(params.get("limit", ["200"])[0]),
|
||||
)
|
||||
)
|
||||
elif path == "/api/backup":
|
||||
self.send_json(
|
||||
backup_payload(include_log=query_includes_log(params))
|
||||
|
||||
Reference in New Issue
Block a user