Files

7582 lines
337 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import os
import asyncio
import shutil
import uuid
import threading
import tempfile
import subprocess
import shutil
import io
import json
import logging
import traceback
import hashlib
import hmac
import secrets
import subprocess
import base64
import re
import unicodedata
import socket
import ssl
import time
import urllib.error
import urllib.request
from typing import Any
from concurrent.futures import ThreadPoolExecutor
from urllib.parse import quote
import urllib.parse
from datetime import datetime, date, timedelta
from pathlib import Path
from fastapi import Body, Depends, FastAPI, File, Form, HTTPException, Request, UploadFile
from fastapi.responses import RedirectResponse, JSONResponse, StreamingResponse, Response, FileResponse, HTMLResponse
from starlette.exceptions import HTTPException as StarletteHTTPException
from starlette.middleware.sessions import SessionMiddleware
from starlette.requests import ClientDisconnect
from fastapi.staticfiles import StaticFiles
from fastapi.templating import Jinja2Templates
from sqlalchemy import and_, case, func, or_, false, distinct, cast, String
from sqlalchemy.orm import Session, aliased, joinedload, selectinload, load_only, contains_eager
from sqlalchemy.exc import IntegrityError
from jinja2 import pass_context
from .database import Base, engine, get_db, SessionLocal
from .fields import ASSET_FIELDS, COMPUTER_FIELDS, DEFAULT_FIELDS, FIELD_DATA_TYPES, SYSTEM_READONLY_FIELDS
from .models import Asset, AssetHistory, Category, CategoryField, StatusOption, SyncRun, User, ChartDefinition, Language, Translation, FieldDefinition, AssetFieldValue, MeshFieldMapping, SoftwarePackage, SoftwarePackageParameter, SoftwareJob, InstalledSoftware, JobEvent, JobDefinition, JobDefinitionRevision, JobDefinitionParameter, JobDefinitionAssetParameter, AssetJobState
from .config import load_config, public_config, save_config
from .meshcentral import synchronize, fetch_device_summaries_for_linking
from .migrations import apply_lightweight_migrations
from .i18n import seed_i18n, translate, dictionary as translation_dictionary, languages as i18n_languages, clear_translation_cache
from .software_control import token_hash, detect_platform, execute_job, build_registry_user_script
from .job_state import backfill_asset_job_states, filter_state_key, sync_asset_job_state
from .presence import mesh_presence_loop
from .version import APP_VERSION
from .backup import (BACKUP_DIR, BACKUP_INTERVAL_HOURS, BACKUP_RETENTION_DAYS, backup_path, create_backup, delete_backup, list_backups, restore_backup, store_uploaded_backup, automatic_backup_loop, system_storage_information)
from openpyxl import Workbook, load_workbook
from openpyxl.styles import Font, PatternFill, Alignment
BASE_DIR = Path(__file__).resolve().parent
UPLOAD_DIR = BASE_DIR / "static" / "uploads"
UPLOAD_DIR.mkdir(parents=True, exist_ok=True)
LOG_DIR = Path(os.getenv("SYNC_LOG_DIR", "/app/data/logs/sync"))
LOG_DIR.mkdir(parents=True, exist_ok=True)
APP_LOG_DIR = Path(os.getenv("APP_LOG_DIR", "/app/data/logs"))
APP_LOG_DIR.mkdir(parents=True, exist_ok=True)
SOFTWARE_CALLBACK_DEBUG_LOG = APP_LOG_DIR / "software-callback-debug.log"
_SOFTWARE_LOG_LOCK = threading.Lock()
def _optional_form_int(value: Any, *, field_name: str, minimum: int | None = None, maximum: int | None = None) -> int | None:
"""Convert an optional HTML form value to int.
Browsers submit an empty numeric input as an empty string. FastAPI/Pydantic
cannot coerce that value directly to ``int | None``, so optional numeric
form fields are normalized here before validation and persistence.
"""
if value is None:
return None
text = str(value).strip()
if not text:
return None
try:
parsed = int(text)
except (TypeError, ValueError) as exc:
raise HTTPException(400, f"Invalid integer value for {field_name}.") from exc
if minimum is not None:
parsed = max(minimum, parsed)
if maximum is not None:
parsed = min(maximum, parsed)
return parsed
def _record_job_event(db: Session, job: SoftwareJob, event_type: str, status: str | None = None, message: str | None = None, created_by: str | None = None) -> JobEvent:
event = JobEvent(
job_id=job.id,
attempt=int(job.attempt_count or 0),
event_type=event_type,
status=status if status is not None else job.status,
message=(message or "")[:4000] or None,
created_by=created_by,
)
db.add(event)
return event
_SOFTWARE_EXCLUSION_PLATFORMS = {"all", "windows", "linux", "macos", "unknown"}
def _normalize_software_exclusion_rules(value: Any) -> list[dict[str, str]]:
"""Return safe, deduplicated substring rules from configuration or forms."""
if not isinstance(value, list):
return []
normalized: list[dict[str, str]] = []
seen: set[tuple[str, str]] = set()
for raw_rule in value[:500]:
if not isinstance(raw_rule, dict):
continue
platform = str(raw_rule.get("platform") or "all").strip().lower()
if platform not in _SOFTWARE_EXCLUSION_PLATFORMS:
platform = "all"
name_contains = re.sub(
r"\s+",
" ",
str(raw_rule.get("name_contains") or raw_rule.get("name") or "").strip(),
)[:300]
if not name_contains:
continue
identity = (platform, name_contains.casefold())
if identity in seen:
continue
seen.add(identity)
normalized.append({"platform": platform, "name_contains": name_contains})
return normalized
def _software_name_is_excluded(
name: str,
platform: str,
rules: list[dict[str, str]],
) -> bool:
"""Match a software name against case-insensitive platform substring rules."""
name_key = str(name or "").casefold()
platform_key = str(platform or "unknown").strip().lower() or "unknown"
for rule in rules:
rule_platform = rule.get("platform", "all")
if rule_platform not in {"all", platform_key}:
continue
needle = str(rule.get("name_contains") or "").casefold()
if needle and needle in name_key:
return True
return False
def _software_settings() -> dict[str, Any]:
config = load_config()
software = dict(config.get("software", {}))
# Compatibility with the callback field that existed in the general section.
if not software.get("callback_base_url"):
software["callback_base_url"] = str(config.get("general", {}).get("callback_base_url", "") or "")
software["inventory_exclusion_rules"] = _normalize_software_exclusion_rules(
software.get("inventory_exclusion_rules", [])
)
return software
def _software_debug_log(message: str, *, force: bool = False) -> None:
settings = _software_settings()
if not force and not bool(settings.get("debug_logging", True)):
return
line = f"{datetime.utcnow().isoformat(timespec='milliseconds')}Z | {message.rstrip()}\n"
try:
with _SOFTWARE_LOG_LOCK:
SOFTWARE_CALLBACK_DEBUG_LOG.parent.mkdir(parents=True, exist_ok=True)
with SOFTWARE_CALLBACK_DEBUG_LOG.open("a", encoding="utf-8") as handle:
handle.write(line)
max_lines = max(100, min(int(settings.get("debug_log_max_lines", 1000) or 1000), 10000))
lines = SOFTWARE_CALLBACK_DEBUG_LOG.read_text(encoding="utf-8", errors="replace").splitlines()
if len(lines) > max_lines:
SOFTWARE_CALLBACK_DEBUG_LOG.write_text("\n".join(lines[-max_lines:]) + "\n", encoding="utf-8")
except Exception:
logger.exception("Software callback debug log could not be written")
def _software_debug_tail(max_lines: int = 400) -> str:
try:
lines = SOFTWARE_CALLBACK_DEBUG_LOG.read_text(encoding="utf-8", errors="replace").splitlines()
return "\n".join(lines[-max(20, min(max_lines, 2000)):])
except OSError:
return ""
logger = logging.getLogger("assetmanager")
logger.setLevel(logging.INFO)
if not logger.handlers:
file_handler = logging.FileHandler(APP_LOG_DIR / "errors.log", encoding="utf-8")
file_handler.setFormatter(logging.Formatter("%(asctime)s | %(levelname)s | %(message)s"))
logger.addHandler(file_handler)
ldap_logger = logging.getLogger("assetmanager.ldap")
ldap_logger.setLevel(logging.DEBUG)
ldap_logger.propagate = False
if not ldap_logger.handlers:
ldap_file_handler = logging.FileHandler(APP_LOG_DIR / "ldap.log", encoding="utf-8")
ldap_file_handler.setFormatter(logging.Formatter("%(asctime)s | %(levelname)s | %(message)s"))
ldap_logger.addHandler(ldap_file_handler)
ldap_console_handler = logging.StreamHandler()
ldap_console_handler.setFormatter(logging.Formatter("LDAP | %(levelname)s | %(message)s"))
ldap_logger.addHandler(ldap_console_handler)
app = FastAPI(title=os.getenv("APP_TITLE", "AssetManager"))
_session_cfg = load_config().get("authentication", {})
_session_env = _session_cfg.get("session_secret_env", "SESSION_SECRET")
_session_secret = os.getenv(_session_env) or os.getenv("SESSION_SECRET") or "assetmanager-change-this-session-secret"
app.mount("/static", StaticFiles(directory=BASE_DIR / "static"), name="static")
templates = Jinja2Templates(directory=BASE_DIR / "templates")
templates.env.globals["application_config"] = load_config
def application_version() -> str:
"""Read VERSION when available and otherwise use the embedded build version."""
candidates = [
BASE_DIR.parent / "VERSION",
BASE_DIR / "VERSION",
]
for candidate in candidates:
try:
value = candidate.read_text(encoding="utf-8").strip()
if value:
return value
except OSError:
continue
return APP_VERSION
templates.env.globals["application_version"] = application_version
def application_info() -> dict[str, str]:
"""Read persistent application metadata.
Primary file: /app/config/APPINFO.json (or APPINFO_PATH).
A legacy APPINFO.json is migrated once when possible.
"""
defaults = {
"name": "AssetManager", "author": "", "organization": "",
"contact": "", "website": "", "repository": "",
"license": "", "copyright": "", "description": "",
}
persistent = Path(os.getenv("APPINFO_PATH", "/app/config/APPINFO.json"))
legacy_candidates = [
Path("/app/data/config/APPINFO.json"), # path used by v0.3.14.2
BASE_DIR.parent / "APPINFO.json", # /app/APPINFO.json
BASE_DIR / "APPINFO.json",
Path("/APPINFO.json"), # legacy file in container root
]
if not persistent.exists():
for legacy in legacy_candidates:
try:
if legacy.exists():
persistent.parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(legacy, persistent)
break
except OSError:
continue
candidates = [persistent, *legacy_candidates]
for candidate in candidates:
try:
payload = json.loads(candidate.read_text(encoding="utf-8"))
if isinstance(payload, dict):
return {**defaults, **{str(k): str(v or "") for k, v in payload.items()}}
except (OSError, ValueError, TypeError):
continue
return defaults
templates.env.globals["application_info"] = application_info
def manufacturer_logo_path() -> str:
"""Return the first fixed manufacturer logo found in the static image folder.
The file is intentionally not configurable through the GUI. Administrators
can replace it only in the application files using the fixed base name
``manufacturer-logo`` and one of the supported image extensions.
"""
image_dir = BASE_DIR / "static" / "img"
for extension in (".svg", ".png", ".jpg", ".jpeg", ".webp"):
candidate = image_dir / f"manufacturer-logo{extension}"
try:
if candidate.is_file():
return f"/static/img/{candidate.name}"
except OSError:
continue
return ""
templates.env.globals["manufacturer_logo_path"] = manufacturer_logo_path
@pass_context
def _template_translate(context, key: str, default: str | None = None, **values):
"""Translate template text without opening a database session for every cell.
The previous implementation created and closed a SQLAlchemy session for every
``t(...)`` call. A job table with hundreds of rows therefore caused thousands
of avoidable pool checkouts. Translation dictionaries are already cached by
the i18n module; keep the two dictionaries on the request and format locally.
"""
request = context.get("request")
language_code = "en"
if request:
user = request.session.get("user") or {}
language_code = user.get("language_code") or request.session.get("language_code") or load_config().get("general", {}).get("default_language", "en")
language_code = str(language_code or "en").lower()
cached = getattr(request.state, "template_translation_dictionaries", None)
if not cached or cached[0] != language_code:
db = SessionLocal()
try:
own = translation_dictionary(db, language_code)
english = translation_dictionary(db, "en")
finally:
db.close()
cached = (language_code, own, english)
request.state.template_translation_dictionaries = cached
_, own, english = cached
text_value = own.get(key) or english.get(key) or default or key
try:
return text_value.format(**values)
except (KeyError, ValueError):
return text_value
db = SessionLocal()
try:
return translate(db, language_code, key, default, **values)
finally:
db.close()
def _template_languages():
db = SessionLocal()
try:
return [{"code": item.code, "name": item.name, "native_name": item.native_name} for item in i18n_languages(db)]
finally:
db.close()
templates.env.globals["t"] = _template_translate
templates.env.globals["available_languages"] = _template_languages
def _utc_iso(value):
"""Serialize project timestamps as explicit UTC for browser conversion."""
if value is None:
return ""
if isinstance(value, datetime):
if value.tzinfo is None:
return value.isoformat(timespec="seconds") + "Z"
return value.astimezone(timezone.utc).isoformat(timespec="seconds").replace("+00:00", "Z")
return str(value)
templates.env.globals["utc_iso"] = _utc_iso
def _translate_request(request: Request, key: str, default: str | None = None, **values) -> str:
user = _session_user(request) or {}
language_code = user.get("language_code") or request.session.get("language_code") or load_config().get("general", {}).get("default_language", "en")
db = SessionLocal()
try:
return translate(db, language_code, key, default, **values)
finally:
db.close()
@app.get("/custom.css", include_in_schema=False)
def custom_css():
css = str(load_config().get("general", {}).get("custom_css", "") or "")
return Response(content=css, media_type="text/css", headers={"Cache-Control": "no-store, max-age=0"})
SYNC_JOBS: dict[str, dict] = {}
SYNC_JOBS_LOCK = threading.Lock()
BACKUP_STOP_EVENT = threading.Event()
BACKUP_THREAD: threading.Thread | None = None
PRESENCE_STOP_EVENT = threading.Event()
PRESENCE_THREAD: threading.Thread | None = None
SOFTWARE_TIMEOUT_STOP_EVENT = threading.Event()
SOFTWARE_TIMEOUT_THREAD: threading.Thread | None = None
SOFTWARE_JOB_TERMINAL_STATES = {"success", "failed", "partial", "timeout", "cancelled"}
SOFTWARE_JOB_AUTOMATIC_RETRY_STATES = {"failed", "partial", "timeout", "sent"}
def _configured_callback_worker_count() -> int:
try:
value = int(_software_settings().get("callback_worker_count", 4) or 4)
except (TypeError, ValueError):
value = 4
return max(1, min(value, 8))
SOFTWARE_CALLBACK_WORKER_COUNT = _configured_callback_worker_count()
SOFTWARE_CALLBACK_EXECUTOR = ThreadPoolExecutor(
max_workers=SOFTWARE_CALLBACK_WORKER_COUNT,
thread_name_prefix="assetmanager-callback",
)
def _expire_stale_software_jobs(db: Session, job_id: int | None = None) -> int:
"""Mark callback jobs as timed out independently of page views."""
query = db.query(SoftwareJob).filter(
SoftwareJob.callback_completed_at.is_(None),
SoftwareJob.callback_expires_at < datetime.utcnow(),
SoftwareJob.status.notin_(SOFTWARE_JOB_TERMINAL_STATES),
)
if job_id is not None:
query = query.filter(SoftwareJob.id == job_id)
rows = query.all()
if not rows:
return 0
language_code = str(load_config().get("general", {}).get("default_language", "de") or "de")
timeout_message = translate(db, language_code, "software.status.timeout", "Timed out")
now = datetime.utcnow()
for job in rows:
job.status = "timeout"
job.message = timeout_message
job.finished_at = now
_record_job_event(db, job, "timeout", "timeout", timeout_message)
sync_asset_job_state(db, job, status_at=now)
db.commit()
return len(rows)
def _automatic_retry_runtime_settings() -> tuple[bool, int, int, int]:
settings = _software_settings()
enabled = bool(settings.get("automatic_retry_enabled", True))
try:
max_age_hours = int(settings.get("automatic_retry_max_age_hours", 12) or 12)
except (TypeError, ValueError):
max_age_hours = 12
try:
retry_interval_seconds = int(settings.get("automatic_retry_interval_seconds", 60) or 60)
except (TypeError, ValueError):
retry_interval_seconds = 60
try:
max_retries = int(settings.get("automatic_retry_max_retries", 3) or 3)
except (TypeError, ValueError):
max_retries = 3
return (
enabled,
max(1, min(max_age_hours, 168)),
max(10, min(retry_interval_seconds, 86400)),
max(1, min(max_retries, 10)),
)
def _job_callback_base_for_background(job: SoftwareJob) -> str:
parameters = job.parameters if isinstance(job.parameters, dict) else {}
stored = str(parameters.get("_callback_base_url") or "").strip().rstrip("/")
if stored:
return stored
configured = str(_software_settings().get("callback_base_url", "") or "").strip().rstrip("/")
if configured:
return configured
# Compatibility for jobs created before the callback base was stored in
# the job snapshot. The resolved script contains the complete callback URL.
match = re.search(
r"(https?://[^\s'\"<>]+?)/api/software-jobs/\d+/callback(?:\?[^\s'\"<>]*)?",
str(job.resolved_script or ""),
flags=re.IGNORECASE,
)
return match.group(1).rstrip("/") if match else ""
def _queue_automatic_job_retries(db: Session) -> list[tuple[int, str, str]]:
enabled, max_age_hours, retry_interval_seconds, max_retries = _automatic_retry_runtime_settings()
if not enabled:
return []
now = datetime.utcnow()
oldest_created_at = now - timedelta(hours=max_age_hours)
retry_ready_at = now - timedelta(seconds=retry_interval_seconds)
candidates = (
db.query(SoftwareJob)
.options(joinedload(SoftwareJob.asset), joinedload(SoftwareJob.package))
.filter(
SoftwareJob.status.in_(SOFTWARE_JOB_AUTOMATIC_RETRY_STATES),
SoftwareJob.created_at >= oldest_created_at,
func.coalesce(SoftwareJob.finished_at, SoftwareJob.sent_at, SoftwareJob.created_at) <= retry_ready_at,
)
.order_by(func.coalesce(SoftwareJob.finished_at, SoftwareJob.sent_at, SoftwareJob.created_at), SoftwareJob.id)
.limit(100)
.all()
)
if not candidates:
return []
candidate_ids = [job.id for job in candidates]
retry_counts = dict(
db.query(JobEvent.job_id, func.count(JobEvent.id))
.filter(
JobEvent.job_id.in_(candidate_ids),
JobEvent.event_type == "automatic_retry_requested",
)
.group_by(JobEvent.job_id)
.all()
)
queued: list[tuple[int, str, str]] = []
for job in candidates:
if int(retry_counts.get(job.id, 0) or 0) >= max_retries:
continue
if not job.asset or not job.asset.mesh_node_id:
continue
callback_base = _job_callback_base_for_background(job)
if not callback_base:
logger.warning(
"Automatic retry skipped for software job %s because no callback base URL could be determined",
job.id,
)
continue
token = _prepare_existing_job_restart(db, None, job, automatic=True)
queued.append((job.id, token, callback_base))
if queued:
db.commit()
return queued
def _software_job_timeout_loop(stop_event: threading.Event) -> None:
while not stop_event.is_set():
queued: list[tuple[int, str, str]] = []
db = SessionLocal()
try:
expired = _expire_stale_software_jobs(db)
if expired:
logger.info("%s software jobs automatically marked as timed out", expired)
queued = _queue_automatic_job_retries(db)
except Exception:
db.rollback()
logger.exception("Software job timeout or retry check failed")
finally:
db.close()
for job_id, token, callback_base in queued:
threading.Thread(
target=execute_job,
args=(job_id, token, callback_base),
daemon=True,
name=f"assetmanager-auto-retry-{job_id}",
).start()
if queued:
logger.info("%s software jobs automatically queued again", len(queued))
_, _, retry_interval_seconds, _ = _automatic_retry_runtime_settings()
stop_event.wait(min(30, max(5, retry_interval_seconds)))
@app.get('/job-definitions', response_class=HTMLResponse)
def job_definitions_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
search_text = str(request.query_params.get('q') or '').strip()
active_only = str(request.query_params.get('active_only', '1')).lower() not in {'0', 'false', 'off', 'no'}
query = db.query(JobDefinition)
if active_only:
query = query.filter(JobDefinition.enabled.is_(True))
if search_text:
pattern = f"%{search_text}%"
query = query.filter(or_(JobDefinition.name.ilike(pattern), JobDefinition.description.ilike(pattern)))
rows = query.order_by(JobDefinition.enabled.desc(), JobDefinition.name.asc()).all()
referenced_definition_ids = set()
for job in db.query(SoftwareJob).filter(SoftwareJob.job_type == 'job_definition').all():
try:
referenced_definition_ids.add(int((job.parameters or {}).get('definition_id')))
except (TypeError, ValueError, AttributeError):
continue
definition_rows = [
{
'item': row,
'can_delete': (not row.is_system and row.id not in referenced_definition_ids),
'delete_reason': (
'system' if row.is_system else
'referenced' if row.id in referenced_definition_ids else
None
),
}
for row in rows
]
return templates.TemplateResponse(
'job_definitions.html',
{
'request': request,
'definition_rows': definition_rows,
'active_only': active_only,
'search_text': search_text,
},
)
@app.get('/job-definitions/new', response_class=HTMLResponse)
def job_definition_new(request: Request):
_require_admin(request)
return templates.TemplateResponse('job_definition_form.html', {'request':request,'definition':None,'test_result':None,'script_preview':None,'preview_values':None})
@app.get('/job-definitions/{definition_id}/edit', response_class=HTMLResponse)
def job_definition_edit(definition_id:int, request:Request, db:Session=Depends(get_db)):
_require_admin(request); row=db.get(JobDefinition,definition_id)
if not row: raise HTTPException(404)
return templates.TemplateResponse('job_definition_form.html', {'request':request,'definition':row,'test_result':None,'script_preview':None,'preview_values':None})
def _parse_registry_entries_json(raw: str) -> list[dict]:
try:
items = json.loads(raw or "[]")
except (TypeError, ValueError):
raise HTTPException(400, "Invalid registry entry definition.")
allowed_actions={"set","delete_value","delete_key"}
allowed_hives={"HKLM","HKCU","HKCR","HKU","HKCC"}
allowed_types={"REG_SZ","REG_EXPAND_SZ","REG_MULTI_SZ","REG_DWORD","REG_QWORD","REG_BINARY"}
result=[]
for item in items if isinstance(items,list) else []:
if not isinstance(item,dict): continue
raw_action=str(item.get("action") or "set")
# Compatibility for the short-lived 0.5.4.1 permission-only action.
legacy_permission_action = raw_action == "set_permissions"
action = "set" if legacy_permission_action else raw_action
hive=str(item.get("hive") or "HKLM").upper()
path=str(item.get("path") or "").strip().strip("\\")
name=str(item.get("name") or "").strip()
value_type=str(item.get("value_type") or "REG_SZ").upper()
if action not in allowed_actions or hive not in allowed_hives or not path: continue
if value_type not in allowed_types: value_type="REG_SZ"
result.append({
"action":action,"hive":hive,"path":path,"name":name,
"value_type":value_type,"value":str(item.get("value") or ""),
"set_value":False if legacy_permission_action else bool(item.get("set_value", True)),
"only_if_missing":bool(item.get("only_if_missing")),
"permissions_enabled":True if legacy_permission_action else bool(item.get("permissions_enabled")),
"principal":str(item.get("principal") or "S-1-5-32-545").strip(),
"rights":str(item.get("rights") or "FullControl").strip(),
"access_type":"Deny" if str(item.get("access_type") or "Allow").lower()=="deny" else "Allow",
"inherit_subkeys":bool(item.get("inherit_subkeys", True)),
"replace_existing":bool(item.get("replace_existing", True)),
})
return result
def _apply_job_definition_form(row, request, name, description, platform, interpreter, source_type, source_path, inline_script, arguments, timeout_seconds, success_codes, enabled, job_kind='script', registry_entries_json='[]', registry_backup=None):
kind = job_kind if job_kind in {'script','registry'} else 'script'
row.name=name.strip();row.description=description.strip() or None;row.job_kind=kind
row.platform='windows' if kind=='registry' else platform
row.interpreter='powershell' if kind=='registry' else interpreter
row.source_type='inline' if kind=='registry' else source_type
row.source_path=None if kind=='registry' else (source_path.strip() or None)
row.inline_script=None if kind=='registry' else (inline_script if source_type=='inline' else None)
row.arguments=None if kind=='registry' else (arguments.strip() or None)
row.configuration={"registry_entries":_parse_registry_entries_json(registry_entries_json),"registry_backup":registry_backup=='on'} if kind=='registry' else {}
row.timeout_seconds=max(1,min(int(timeout_seconds or 300),86400));row.success_codes=success_codes.strip() or '0';row.enabled=enabled=='on';row.updated_by=(request.session.get('user') or {}).get('username')
def _parse_job_parameters_json(raw: str) -> list[dict]:
try:
items = json.loads(raw or "[]")
except (TypeError, ValueError):
raise HTTPException(400, "Invalid job parameter definition.")
result=[]; seen=set()
for index,item in enumerate(items if isinstance(items,list) else []):
if not isinstance(item,dict): continue
key=re.sub(r"[^A-Za-z0-9_]", "_", str(item.get("parameter_key") or "").strip()).strip("_")
if not key or key in seen: continue
seen.add(key)
result.append({"parameter_key":key,"label":str(item.get("label") or key).strip()[:140],"description":str(item.get("description") or "").strip() or None,"data_type":str(item.get("data_type") or "text") if str(item.get("data_type") or "text") in {"text","number","boolean","path"} else "text","default_value":str(item.get("default_value") or ""),"required":bool(item.get("required")),"sort_order":index})
return result
def _sync_job_definition_parameters(db: Session, definition: JobDefinition, items: list[dict]) -> None:
existing={x.parameter_key:x for x in definition.parameters}
wanted={x["parameter_key"] for x in items}
for row in list(definition.parameters):
if row.parameter_key not in wanted:
db.delete(row)
for item in items:
row=existing.get(item["parameter_key"])
if row is None:
row=JobDefinitionParameter(definition_id=definition.id, parameter_key=item["parameter_key"]); db.add(row)
for key,value in item.items(): setattr(row,key,value)
db.flush()
def _job_effective_parameters(db: Session, definition: JobDefinition, asset: Asset) -> dict[str,str]:
parameter_ids=[x.id for x in definition.parameters]
overrides={}
if parameter_ids:
overrides={x.definition_parameter_id:x.value_text for x in db.query(JobDefinitionAssetParameter).filter(JobDefinitionAssetParameter.asset_id==asset.id, JobDefinitionAssetParameter.definition_parameter_id.in_(parameter_ids)).all()}
return {x.parameter_key: (overrides[x.id] if x.id in overrides else (x.default_value or "")) for x in definition.parameters}
@app.post('/job-definitions')
def job_definition_create(request:Request,name:str=Form(...),description:str=Form(''),job_kind:str=Form('script'),platform:str=Form('windows'),interpreter:str=Form('powershell'),source_type:str=Form('inline'),source_path:str=Form(''),inline_script:str=Form(''),arguments:str=Form(''),timeout_seconds:int=Form(300),success_codes:str=Form('0'),enabled:str|None=Form(None),parameters_json:str=Form('[]'),registry_entries_json:str=Form('[]'),registry_backup:str|None=Form(None),db:Session=Depends(get_db)):
_require_admin(request); row=JobDefinition(created_by=(request.session.get('user') or {}).get('username'))
_apply_job_definition_form(row,request,name,description,platform,interpreter,source_type,source_path,inline_script,arguments,timeout_seconds,success_codes,enabled,job_kind,registry_entries_json,registry_backup);db.add(row);db.flush();_sync_job_definition_parameters(db,row,_parse_job_parameters_json(parameters_json));db.commit();db.refresh(row)
db.add(JobDefinitionRevision(definition_id=row.id,revision=row.revision,snapshot=_job_definition_snapshot(row),changed_by=row.updated_by));db.commit()
return RedirectResponse('/job-definitions?toast_success='+quote(_translate_request(request,'jobdefs.saved','Job definition saved.')),303)
@app.post('/job-definitions/{definition_id}')
def job_definition_update(definition_id:int,request:Request,name:str=Form(...),description:str=Form(''),job_kind:str=Form('script'),platform:str=Form('windows'),interpreter:str=Form('powershell'),source_type:str=Form('inline'),source_path:str=Form(''),inline_script:str=Form(''),arguments:str=Form(''),timeout_seconds:int=Form(300),success_codes:str=Form('0'),enabled:str|None=Form(None),parameters_json:str=Form('[]'),registry_entries_json:str=Form('[]'),registry_backup:str|None=Form(None),db:Session=Depends(get_db)):
_require_admin(request);row=db.get(JobDefinition,definition_id)
if not row: raise HTTPException(404)
row.revision=int(row.revision or 0)+1;_apply_job_definition_form(row,request,name,description,platform,interpreter,source_type,source_path,inline_script,arguments,timeout_seconds,success_codes,enabled,job_kind,registry_entries_json,registry_backup);_sync_job_definition_parameters(db,row,_parse_job_parameters_json(parameters_json));db.flush();db.add(JobDefinitionRevision(definition_id=row.id,revision=row.revision,snapshot=_job_definition_snapshot(row),changed_by=row.updated_by));db.commit()
return RedirectResponse(f'/job-definitions/{row.id}/edit?toast_success='+quote(_translate_request(request,'jobdefs.saved','Job definition saved.')),303)
@app.post('/job-definitions/{definition_id}/delete')
def job_definition_delete(definition_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
row = db.get(JobDefinition, definition_id)
if not row:
raise HTTPException(404)
if row.is_system:
return RedirectResponse('/job-definitions?toast_error=' + quote(_translate_request(request, 'jobdefs.delete.system_forbidden', 'System job definitions cannot be deleted.')), 303)
referenced_jobs = []
for job in db.query(SoftwareJob).filter(SoftwareJob.job_type == 'job_definition').all():
try:
referenced_id = int((job.parameters or {}).get('definition_id'))
except (TypeError, ValueError):
continue
if referenced_id == row.id:
referenced_jobs.append(job.id)
if len(referenced_jobs) >= 1:
break
if referenced_jobs:
message = _translate_request(
request,
'jobdefs.delete.referenced',
'This job definition is referenced by existing jobs and cannot be deleted. Deactivate it instead.',
)
return RedirectResponse('/job-definitions?toast_error=' + quote(message), 303)
name = row.name
db.delete(row)
db.commit()
message = _translate_request(request, 'jobdefs.deleted', 'Job definition "{name}" was deleted.').format(name=name)
return RedirectResponse('/job-definitions?toast_success=' + quote(message), 303)
@app.post('/job-definitions/{definition_id}/test', response_class=HTMLResponse)
def job_definition_test(definition_id:int,request:Request,name:str=Form(...),description:str=Form(''),job_kind:str=Form('script'),platform:str=Form('windows'),interpreter:str=Form('powershell'),source_type:str=Form('inline'),source_path:str=Form(''),inline_script:str=Form(''),arguments:str=Form(''),timeout_seconds:int=Form(300),success_codes:str=Form('0'),enabled:str|None=Form(None),parameters_json:str=Form('[]'),registry_entries_json:str=Form('[]'),registry_backup:str|None=Form(None),db:Session=Depends(get_db)):
_require_admin(request);row=db.get(JobDefinition,definition_id)
if not row: raise HTTPException(404)
_apply_job_definition_form(row,request,name,description,platform,interpreter,source_type,source_path,inline_script,arguments,timeout_seconds,success_codes,enabled,job_kind,registry_entries_json,registry_backup)
_sync_job_definition_parameters(db,row,_parse_job_parameters_json(parameters_json))
result=_validate_job_definition(row, request)
preview=''
if row.source_type=='inline':
preview_interpreter='powershell' if row.job_kind=='registry' else row.interpreter
preview_source=build_registry_user_script(row.configuration) if row.job_kind=='registry' else (row.inline_script or '')
wrapper=_generate_job_script(preview_interpreter,'wrapper')
start_marker,end_marker=JOB_SCRIPT_MARKERS[preview_interpreter]
block=start_marker+'\n\n'+end_marker
indent={'powershell':' ','python':' ','bash':'','cmd':''}[preview_interpreter]
commands='\n'.join(indent+line if line.strip() else line for line in preview_source.splitlines())
preview=wrapper.replace(block,start_marker+'\n'+commands+'\n'+end_marker,1)
values=_job_preview_values()
for parameter in row.parameters: values[parameter.parameter_key]=parameter.default_value or ''
preview=_replace_job_placeholders(preview,values)
db.rollback()
return templates.TemplateResponse('job_definition_form.html', {'request':request,'definition':row,'test_result':result,'script_preview':preview,'preview_values':_job_preview_values()})
JOB_SCRIPT_TEMPLATE_DIR = BASE_DIR / "templates" / "jobs"
JOB_SCRIPT_TEMPLATE_FILES = {
"powershell": ("powershell_wrapper.ps1", "powershell_test.txt"),
"python": ("python_wrapper.py.txt", "python_test.txt"),
"bash": ("bash_wrapper.sh", "bash_test.txt"),
"cmd": ("cmd_wrapper.cmd", "cmd_test.txt"),
}
JOB_SCRIPT_MARKERS = {
"powershell": (" # --- Eigene Befehle hier einfügen ---", " # --- Ende eigene Befehle ---"),
"python": (" # --- Eigene Befehle hier einfügen ---", " # --- Ende eigene Befehle ---"),
"bash": ("# --- Eigene Befehle hier einfügen ---", "# --- Ende eigene Befehle ---"),
"cmd": ("rem --- Eigene Befehle hier einfügen ---", "rem --- Ende eigene Befehle ---"),
}
def _read_job_script_template(filename: str) -> str:
path = (JOB_SCRIPT_TEMPLATE_DIR / filename).resolve()
root = JOB_SCRIPT_TEMPLATE_DIR.resolve()
if root not in path.parents or not path.is_file():
raise HTTPException(500, "Job script template is missing.")
return path.read_text(encoding="utf-8").replace("\r\n", "\n")
def _generate_job_script(interpreter: str, action: str) -> str:
interpreter = str(interpreter or "").strip().lower()
action = str(action or "").strip().lower()
if interpreter not in JOB_SCRIPT_TEMPLATE_FILES:
raise HTTPException(400, "Unsupported interpreter.")
if action not in {"wrapper", "test"}:
raise HTTPException(400, "Unsupported generator action.")
wrapper_file, test_file = JOB_SCRIPT_TEMPLATE_FILES[interpreter]
wrapper = _read_job_script_template(wrapper_file)
if action == "wrapper":
return wrapper
commands = _read_job_script_template(test_file).rstrip("\n")
start_marker, end_marker = JOB_SCRIPT_MARKERS[interpreter]
marker_block = start_marker + "\n\n" + end_marker
replacement = start_marker + "\n" + commands + "\n" + end_marker
if marker_block not in wrapper:
raise HTTPException(500, "Job script template markers are invalid.")
return wrapper.replace(marker_block, replacement, 1)
@app.post('/api/job-definitions/script-generator')
async def job_definition_script_generator(request: Request):
_require_admin(request)
try:
payload = await request.json()
except Exception as exc:
raise HTTPException(400, "Invalid JSON request.") from exc
interpreter = str((payload or {}).get("interpreter") or "")
action = str((payload or {}).get("action") or "")
logger.info(
"Job script generator requested | user=%s interpreter=%s action=%s",
(_session_user(request) or {}).get("username"),
interpreter,
action,
)
script = _generate_job_script(interpreter, action)
logger.info(
"Job script generator completed | interpreter=%s action=%s bytes=%s",
interpreter,
action,
len(script.encode("utf-8")),
)
return JSONResponse({"script": script, "interpreter": interpreter, "action": action})
@app.on_event("startup")
def start_backup_scheduler():
global BACKUP_THREAD
storage = system_storage_information()
backup_status = storage["backups"]
logger.info(
"Backup directory: %s | exists=%s | writable=%s | mount=%s",
backup_status["path"], backup_status["exists"], backup_status["writable"], backup_status["is_mount"],
)
if not backup_status["writable"]:
logger.error("Backup directory is not writable: %s", backup_status.get("error") or backup_status["path"])
if BACKUP_THREAD is None or not BACKUP_THREAD.is_alive():
BACKUP_STOP_EVENT.clear()
BACKUP_THREAD = threading.Thread(target=automatic_backup_loop, args=(BACKUP_STOP_EVENT,), daemon=True, name="assetmanager-backup")
BACKUP_THREAD.start()
@app.on_event("shutdown")
def stop_backup_scheduler():
BACKUP_STOP_EVENT.set()
def _hash_password(password: str) -> str:
salt = secrets.token_bytes(16)
iterations = 310_000
digest = hashlib.pbkdf2_hmac("sha256", password.encode("utf-8"), salt, iterations)
return f"pbkdf2_sha256${iterations}${salt.hex()}${digest.hex()}"
def _verify_password(password: str, encoded: str | None) -> bool:
if not encoded:
return False
try:
algorithm, iterations_text, salt_hex, digest_hex = encoded.split("$", 3)
if algorithm != "pbkdf2_sha256":
return False
candidate = hashlib.pbkdf2_hmac(
"sha256", password.encode("utf-8"), bytes.fromhex(salt_hex), int(iterations_text)
)
return hmac.compare_digest(candidate.hex(), digest_hex)
except (ValueError, TypeError):
return False
def _ensure_local_emergency_admin(db: Session) -> None:
"""Create or protect the configured local break-glass administrator.
The password from the environment is used only when the account is first
created or when an existing local account has no password hash. Existing
password hashes are never overwritten during container startup.
"""
username = str(os.getenv("LOCAL_ADMIN_USERNAME", "") or "").strip()
password = str(os.getenv("LOCAL_ADMIN_PASSWORD", "") or "")
if not username and not password:
return
if not username or not password:
logger.error(
"Lokaler Notfall-Administrator wurde nicht eingerichtet: "
"LOCAL_ADMIN_USERNAME und LOCAL_ADMIN_PASSWORD müssen gemeinsam gesetzt sein."
)
return
if len(username) > 120:
logger.error("LOCAL_ADMIN_USERNAME ist länger als 120 Zeichen; Konto wurde nicht eingerichtet.")
return
if len(password) < 12:
logger.error("LOCAL_ADMIN_PASSWORD muss mindestens 12 Zeichen lang sein; Konto wurde nicht eingerichtet.")
return
matches = db.query(User).filter(func.lower(User.username) == username.casefold()).all()
if len(matches) > 1:
logger.error(
"Lokaler Notfall-Administrator konnte nicht eingerichtet werden: "
"Der Benutzername %r existiert mehrfach mit unterschiedlicher Schreibweise.",
username,
)
return
user = matches[0] if matches else None
if user and user.auth_source != "local":
logger.error(
"Lokaler Notfall-Administrator konnte nicht eingerichtet werden: "
"Benutzername %r gehört bereits zu einem %s-Konto.",
username,
user.auth_source,
)
return
created = False
password_initialized = False
if not user:
user = User(
username=username,
display_name="Lokaler Notfall-Administrator",
password_hash=_hash_password(password),
auth_source="local",
is_active=True,
is_admin=True,
is_protected=True,
allowed_mesh_groups=[],
asset_access_scopes=["all"],
language_code=load_config().get("general", {}).get("default_language", "de"),
)
db.add(user)
created = True
password_initialized = True
else:
user.is_active = True
user.is_admin = True
user.is_protected = True
user.auth_source = "local"
if not user.display_name:
user.display_name = "Lokaler Notfall-Administrator"
if not user.password_hash:
user.password_hash = _hash_password(password)
password_initialized = True
db.commit()
logger.info(
"Lokaler Notfall-Administrator %r ist aktiv und geschützt "
"(neu=%s, Passwort initialisiert=%s).",
user.username,
created,
password_initialized,
)
def _session_user(request: Request) -> dict | None:
value = request.session.get("user")
return value if isinstance(value, dict) else None
def _changed_by(request: Request) -> str | None:
user = _session_user(request)
return (user.get("display_name") or user.get("username")) if user else None
def _login_user(request: Request, user: User) -> None:
request.session["user"] = {
"id": user.id,
"username": user.username,
"display_name": user.display_name or user.username,
"email": user.email or "",
"auth_source": user.auth_source,
"is_admin": bool(user.is_admin),
"is_protected": bool(getattr(user, "is_protected", False)),
"allowed_mesh_groups": user.allowed_mesh_groups or [],
"asset_access_scopes": user.asset_access_scopes or [],
"profile_department": user.profile_department or "",
"profile_location": user.profile_location or "",
"language_code": user.language_code or "en",
"persist_browser_state": bool(getattr(user, "persist_browser_state", False)),
"toast_position": getattr(user, "toast_position", "top-right") or "top-right",
"toast_duration_seconds": int(getattr(user, "toast_duration_seconds", 6) or 6),
"theme_mode": getattr(user, "theme_mode", "light") if getattr(user, "theme_mode", "light") in {"light", "dark"} else "light",
}
def _is_admin(request: Request) -> bool:
user = _session_user(request)
return bool(user and user.get("is_admin"))
def _require_admin(request: Request) -> None:
if not _is_admin(request):
raise HTTPException(403, "Diese Funktion ist nur für Administratoren verfügbar.")
ASSET_ACCESS_SCOPES = {"self", "department", "location", "all"}
def _allowed_groups(request: Request) -> list[str] | None:
user = _session_user(request)
if not user or user.get("is_admin"):
return None
return [str(x).strip() for x in (user.get("allowed_mesh_groups") or []) if str(x).strip()]
def _asset_access_conditions(request: Request):
"""Build additive (OR-linked) asset visibility conditions for a non-admin user."""
user = _session_user(request)
if not user or user.get("is_admin"):
return None
scopes = {str(value).strip().lower() for value in (user.get("asset_access_scopes") or [])}
scopes &= ASSET_ACCESS_SCOPES
if "all" in scopes:
return None
conditions = []
username = str(user.get("username") or "").strip()
department = str(user.get("profile_department") or "").strip()
location = str(user.get("profile_location") or "").strip()
if "self" in scopes and username:
conditions.append(func.lower(func.trim(Asset.assigned_to)) == username.casefold())
if "department" in scopes and department:
conditions.append(func.lower(func.trim(Asset.department)) == department.casefold())
if "location" in scopes and location:
conditions.append(func.lower(func.trim(Asset.location)) == location.casefold())
groups = _allowed_groups(request) or []
if groups:
conditions.append(Asset.mesh_group.in_(groups))
return conditions
def _apply_asset_access(query, request: Request):
conditions = _asset_access_conditions(request)
if conditions is None:
return query
return query.filter(or_(*conditions) if conditions else false())
def _get_visible_asset(db: Session, request: Request, asset_id: int) -> Asset:
query = _apply_asset_access(db.query(Asset), request)
asset = query.filter(Asset.id == asset_id).first()
if not asset:
raise HTTPException(404, "Asset nicht gefunden oder keine Berechtigung")
return asset
def _action_status(db: Session, action: str) -> StatusOption | None:
column = StatusOption.use_for_issue if action == "issue" else StatusOption.use_for_return
return db.query(StatusOption).filter(StatusOption.active.is_(True), column.is_(True)).first()
def _ldap_authenticate(username: str, password: str, auth_config: dict) -> dict | None:
from ldap3 import ALL, Connection, Server, SUBTREE
from ldap3.core.exceptions import LDAPException
from ldap3.utils.conv import escape_filter_chars
ldap = auth_config.get("ldap", {})
server_name = str(ldap.get("server", "")).strip()
base_dn = str(ldap.get("user_base_dn", "")).strip()
use_ssl = bool(ldap.get("use_ssl", False))
start_tls = bool(ldap.get("start_tls", False))
port = int(ldap.get("port") or (636 if use_ssl else 389))
bind_dn = str(ldap.get("bind_dn", "")).strip()
env_name = str(ldap.get("bind_password_env", "LDAP_BIND_PASSWORD")).strip()
bind_password = os.getenv(env_name, "") if env_name else ""
if not bind_password:
bind_password = str(ldap.get("bind_password", ""))
ldap_logger.info(
"Anmeldeversuch Benutzer=%r Server=%s Port=%s SSL=%s StartTLS=%s BaseDN=%r BindDN=%r Passwortvariable=%r gesetzt=%s",
username, server_name, port, use_ssl, start_tls, base_dn, bind_dn, env_name, bool(bind_password),
)
if not server_name:
ldap_logger.error("LDAP-Server ist nicht konfiguriert")
return None
if not base_dn:
ldap_logger.error("Benutzer-Basis-DN ist nicht konfiguriert")
return None
if not password:
ldap_logger.warning("Leeres Benutzerpasswort für Benutzer=%r", username)
return None
if bind_dn and not bind_password:
ldap_logger.error("Bind-DN ist gesetzt, aber das Bind-Passwort fehlt. Erwartete Umgebungsvariable=%r", env_name)
return None
server = Server(server_name, port=port, use_ssl=use_ssl, get_info=ALL, connect_timeout=10)
search_connection = None
user_connection = None
try:
search_connection = Connection(
server,
user=bind_dn or None,
password=bind_password or None,
auto_bind=False,
receive_timeout=15,
raise_exceptions=False,
)
if start_tls and not use_ssl:
search_connection.open()
ldap_logger.debug("TCP/LDAP-Verbindung geöffnet")
if not search_connection.start_tls():
ldap_logger.error("StartTLS fehlgeschlagen: result=%r last_error=%r", search_connection.result, search_connection.last_error)
return None
ldap_logger.debug("StartTLS erfolgreich")
if not search_connection.bind():
ldap_logger.error("Dienstkonto-Bind fehlgeschlagen: result=%r last_error=%r", search_connection.result, search_connection.last_error)
return None
ldap_logger.info("Dienstkonto-Bind erfolgreich")
raw_filter = str(ldap.get("user_filter", "(sAMAccountName={username})"))
escaped_username = escape_filter_chars(username)
user_filter = raw_filter.replace("{username}", escaped_username).replace("{{username}}", escaped_username)
display_attr = str(ldap.get("display_name_attribute", "displayName"))
email_attr = str(ldap.get("email_attribute", "mail"))
attributes = list(dict.fromkeys([display_attr, email_attr, "distinguishedName"]))
ldap_logger.info("Benutzersuche BaseDN=%r Filter=%r Attribute=%r", base_dn, user_filter, attributes)
search_ok = search_connection.search(
base_dn,
user_filter,
search_scope=SUBTREE,
attributes=attributes,
)
ldap_logger.info(
"Benutzersuche beendet: success=%s Treffer=%s result=%r last_error=%r",
search_ok, len(search_connection.entries), search_connection.result, search_connection.last_error,
)
if len(search_connection.entries) != 1:
if len(search_connection.entries) == 0:
ldap_logger.warning("Kein LDAP-Benutzer für Benutzer=%r gefunden", username)
else:
ldap_logger.warning("Mehrere LDAP-Benutzer für Benutzer=%r gefunden: %s", username, len(search_connection.entries))
return None
entry = search_connection.entries[0]
user_dn = entry.entry_dn
display_name = str(getattr(entry, display_attr, "") or username)
email = str(getattr(entry, email_attr, "") or "")
ldap_logger.info("Benutzer gefunden: DN=%r Anzeigename=%r E-Mail=%r", user_dn, display_name, email)
user_connection = Connection(
server,
user=user_dn,
password=password,
auto_bind=False,
receive_timeout=15,
raise_exceptions=False,
)
if start_tls and not use_ssl:
user_connection.open()
if not user_connection.start_tls():
ldap_logger.error("StartTLS beim Benutzer-Bind fehlgeschlagen: result=%r last_error=%r", user_connection.result, user_connection.last_error)
return None
if not user_connection.bind():
ldap_logger.warning("Benutzer-Bind fehlgeschlagen für DN=%r: result=%r last_error=%r", user_dn, user_connection.result, user_connection.last_error)
return None
ldap_logger.info("LDAP-Anmeldung erfolgreich für Benutzer=%r DN=%r", username, user_dn)
return {"username": username, "display_name": display_name, "email": email}
except LDAPException as exc:
ldap_logger.exception("LDAP-Ausnahme für Benutzer=%r: %s", username, exc)
return None
except Exception as exc:
ldap_logger.exception("Unerwarteter LDAP-Fehler für Benutzer=%r: %s", username, exc)
return None
finally:
for connection in (user_connection, search_connection):
if connection is not None:
try:
connection.unbind()
except Exception:
pass
@app.middleware("http")
async def authentication_gate(request: Request, call_next):
mode = str(load_config().get("authentication", {}).get("mode", "none"))
public_paths = {"/login", "/logout"}
is_agent_callback = (request.url.path.startswith("/api/software-jobs/") and request.url.path.endswith("/callback")) or request.url.path == "/api/software-callback/health"
if mode in {"local", "ldap"} and not request.session.get("user") and request.url.path not in public_paths and not request.url.path.startswith("/static/") and not is_agent_callback:
if request.method == "GET":
return RedirectResponse("/login?toast_warning=Bitte zuerst anmelden", status_code=303)
return JSONResponse({"detail": "Anmeldung erforderlich"}, status_code=401)
return await call_next(request)
app.add_middleware(
SessionMiddleware,
secret_key=_session_secret,
same_site="lax",
https_only=False,
)
def _is_browser_request(request: Request) -> bool:
accept = request.headers.get("accept", "")
return "text/html" in accept or request.method != "GET"
def _toast_redirect(request: Request, message: str, level: str = "error") -> RedirectResponse:
referer = request.headers.get("referer")
target = referer if referer and referer.startswith(str(request.base_url).rstrip("/")) else "/"
separator = "&" if "?" in target else "?"
return RedirectResponse(f"{target}{separator}toast_{level}={quote(message)}", status_code=303)
@app.exception_handler(StarletteHTTPException)
async def http_exception_handler(request: Request, exc: StarletteHTTPException):
message = str(exc.detail) if exc.detail else f"HTTP-Fehler {exc.status_code}"
logger.error("HTTP %s | %s %s | %s", exc.status_code, request.method, request.url.path, message)
if _is_browser_request(request):
return _toast_redirect(request, message)
return JSONResponse({"detail": message}, status_code=exc.status_code)
@app.exception_handler(Exception)
async def unhandled_exception_handler(request: Request, exc: Exception):
logger.error(
"Unerwarteter Fehler | %s %s | %s\n%s",
request.method,
request.url.path,
exc,
traceback.format_exc(),
)
if _is_browser_request(request):
return _toast_redirect(request, "Ein unerwarteter Fehler ist aufgetreten. Details wurden protokolliert.")
return JSONResponse({"detail": "Interner Serverfehler"}, status_code=500)
def _job_log(job_id: str, level: str, message: str) -> None:
entry = {
"time": datetime.now().strftime("%H:%M:%S"),
"level": level,
"message": str(message),
}
with SYNC_JOBS_LOCK:
job = SYNC_JOBS.get(job_id)
if job is not None:
job["logs"].append(entry)
job["logs"] = job["logs"][-10000:]
try:
with (LOG_DIR / f"{job_id}.jsonl").open("a", encoding="utf-8") as handle:
import json
handle.write(json.dumps(entry, ensure_ascii=False) + "\n")
except OSError:
pass
def _run_sync_job(job_id: str) -> None:
db = SessionLocal()
try:
result = synchronize(db, lambda level, message: _job_log(job_id, level, message))
with SYNC_JOBS_LOCK:
job = SYNC_JOBS[job_id]
job["status"] = result.run.status
job["run_id"] = result.run.id
job["finished"] = True
except Exception as exc:
_job_log(job_id, "error", f"Unerwarteter Fehler: {exc}")
with SYNC_JOBS_LOCK:
job = SYNC_JOBS[job_id]
job["status"] = "error"
job["finished"] = True
finally:
db.close()
def _seed_field_definitions(db: Session) -> None:
existing = {item.field_name: item for item in db.query(FieldDefinition).all()}
for order, (field_name, label) in enumerate(ASSET_FIELDS.items()):
item = existing.get(field_name)
if not item:
item = FieldDefinition(
field_name=field_name, label=label, description=None,
data_type=("boolean" if field_name == "mesh_online" else "integer" if field_name == "mesh_mtype" else "long_text" if field_name in {"notes","storage_details","memory_details","gpu_name","antivirus"} else "date" if field_name in {"purchase_date","warranty_until","assigned_on"} else "status" if field_name == "status" else "ip" if field_name == "ip_address" else "short_text"),
readonly=field_name in SYSTEM_READONLY_FIELDS, is_system=True,
translation_key=f"field.{field_name}", sort_order=order,
excel_import=field_name not in SYSTEM_READONLY_FIELDS,
bulk_edit=field_name not in SYSTEM_READONLY_FIELDS,
is_default=field_name in DEFAULT_FIELDS,
)
db.add(item); db.flush(); existing[field_name] = item
elif field_name in SYSTEM_READONLY_FIELDS:
item.readonly = True; item.excel_import = False; item.bulk_edit = False
db.flush()
for link in db.query(CategoryField).all():
definition = existing.get(link.field_name)
if definition and link.field_definition_id != definition.id:
link.field_definition_id = definition.id
# Every category gets every available field; new links are inactive.
categories = db.query(Category).all()
definitions = db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).order_by(FieldDefinition.sort_order).all()
for category in categories:
linked = {x.field_definition_id for x in category.visible_fields if x.field_definition_id}
linked_names = {x.field_name for x in category.visible_fields}
next_order = max((x.sort_order for x in category.visible_fields), default=-1) + 1
for definition in definitions:
if definition.id in linked or definition.field_name in linked_names:
continue
db.add(CategoryField(category_id=category.id, field_name=definition.field_name, label=definition.label, field_definition_id=definition.id, active=False, required=False, show_in_list=False, readonly=definition.readonly, sort_order=next_order))
next_order += 1
db.commit()
DEFAULT_MESH_FIELD_MAPPINGS = [
# source_path, target, transform, multi mode, separator, enabled, priority, description
("node._id|node.nodeid|node.nodeId|node.id|_id|nodeid|nodeId|id", "mesh_node_id", "string", "first", "\n", True, 10, "Eindeutige MeshCentral Node-ID"),
("node.name|node.hostname|node.rname|name|hostname|rname", "name", "string", "first", "\n", True, 20, "Gerätename aus MeshCentral"),
("node.hostname|node.name|node.rname|hostname|name|rname", "hostname", "string", "first", "\n", True, 30, "Computername/Hostname"),
("node.osdesc|node.os|node.operatingSystem|osdesc|os|operatingSystem", "operating_system", "string", "first", "\n", True, 40, "Betriebssystembeschreibung"),
("sys.hardware.windows.osinfo.BuildNumber|sys.hardware.windows.osinfo.Version", "os_patch_level", "string", "first", "\n", True, 50, "Windows Build- oder Versionsnummer"),
("net.netif2", "ip_address", "active_ipv4_address", "first", "\n", True, 60, "IPv4-Adresse des ersten aktiven Nicht-Loopback-Adapters"),
("net.netif2", "mac_address", "active_ipv4_mac", "first", "\n", True, 70, "MAC-Adresse desselben aktiven IPv4-Adapters"),
("sys.hardware.identifiers.board_vendor|sys.hardware.identifiers.bios_vendor", "manufacturer", "string", "first", "\n", True, 80, "System- oder BIOS-Hersteller"),
("sys.hardware.identifiers.product_name", "model", "string", "first", "\n", True, 90, "Produkt-/Modellbezeichnung"),
("sys.hardware.identifiers.bios_serial|sys.hardware.windows.bios.SerialNumber|node.serial_number|node.serialNumber", "serial_number", "normalize_serial", "first", "\n", True, 100, "BIOS-/Geräteseriennummer"),
("sys.hardware.identifiers.cpu_name|sys.hardware.windows.cpu[].Name", "cpu", "string", "first", "\n", True, 110, "CPU-Bezeichnung"),
("sys.hardware.windows.memory", "memory_gb", "sum_memory_gb", "first", "\n", True, 120, "Summe der Speichermodule in GiB"),
("sys.hardware.windows.memory", "memory_details", "format_memory_modules", "first", "\n", True, 130, "Formatierte Liste der Speichermodule"),
("sys.hardware.windows.drives", "storage_gb", "sum_storage_gb", "first", "\n", True, 140, "Gesamtkapazität der Datenträger in GB"),
("sys.hardware.windows.drives", "storage_details", "format_drives", "first", "\n", True, 150, "Formatierte Liste der Datenträger"),
("sys.hardware.windows.gpu", "gpu_name", "gpu_names", "first", "\n", True, 160, "Namen der erkannten Grafikkarten"),
("sys.hardware.tpm.SpecVersion", "tpm_version", "string", "first", "\n", True, 170, "TPM-Spezifikationsversion"),
("node.av|av", "antivirus", "active_antivirus", "first", "\n", True, 180, "Aktive Antivirenprodukte"),
("node.groupname|groupname", "mesh_group", "string", "first", "\n", True, 190, "MeshCentral-Gerätegruppe"),
("node.mtype|node.icon|mtype|icon", "mesh_mtype", "integer", "first", "\n", True, 200, "MeshCentral Gerätetyp (mtype)"),
]
def _seed_mesh_field_mappings(db: Session) -> None:
"""Create missing defaults without overwriting persisted user settings."""
rows = db.query(MeshFieldMapping).order_by(MeshFieldMapping.id).all()
by_target: dict[str, list[MeshFieldMapping]] = {}
for row in rows:
by_target.setdefault(row.target_field_name, []).append(row)
# Import legacy config mappings only when that exact database mapping does
# not already exist. Existing database rows remain the source of truth.
config_mappings = load_config().get("meshcentral", {}).get("field_mappings", []) or []
for index, item in enumerate(config_mappings):
if not isinstance(item, dict):
continue
target = str(item.get("target_field") or item.get("target_field_name") or "").strip()
source = str(item.get("source_path") or "").strip()
if not target or not source:
continue
if any(row.source_path == source and row.target_field_name == target for row in rows):
continue
row = MeshFieldMapping(
source_path=source,
target_field_name=target,
transform=str(item.get("transform") or "raw"),
multi_value_mode=str(item.get("multi_value_mode") or "first"),
separator=str(item.get("separator") or "\\n"),
enabled=bool(item.get("enabled", True)),
priority=int(item.get("priority") or (1000 + index)),
description=item.get("description"),
update_rule=str(item.get("update_rule") or "fill_empty"),
is_system=False,
)
db.add(row)
rows.append(row)
by_target.setdefault(target, []).append(row)
legacy_rules = load_config().get("meshcentral", {}).get("field_rules", {}) or {}
for source, target, transform, multi, separator, enabled, priority, description in DEFAULT_MESH_FIELD_MAPPINGS:
candidates = by_target.get(target, [])
system_row = next((row for row in candidates if bool(getattr(row, "is_system", False))), None)
if system_row is None:
# A legacy database may contain exactly one mapping for the target.
# Adopt it as the system row without changing its configured values.
system_row = candidates[0] if candidates else None
if system_row is None:
system_row = MeshFieldMapping(
source_path=source,
target_field_name=target,
transform=transform,
multi_value_mode=multi,
separator=separator,
enabled=enabled,
priority=priority,
description=description,
update_rule=str(legacy_rules.get(target) or "fill_empty"),
is_system=True,
)
db.add(system_row)
by_target.setdefault(target, []).append(system_row)
continue
# Only fill genuinely missing values. Never reset a user's persisted
# "who wins" rule, enabled state, source path, transform or priority.
system_row.is_system = True
if not str(system_row.source_path or "").strip():
system_row.source_path = source
if not str(system_row.transform or "").strip():
system_row.transform = transform
if not str(system_row.multi_value_mode or "").strip():
system_row.multi_value_mode = multi
if system_row.separator is None:
system_row.separator = separator
if system_row.priority is None:
system_row.priority = priority
if not str(system_row.description or "").strip():
system_row.description = description
if not str(system_row.update_rule or "").strip():
configured_rule = legacy_rules.get(target)
system_row.update_rule = configured_rule if configured_rule in {"fill_empty", "meshcentral_wins", "local_wins"} else "fill_empty"
db.commit()
def _definition_label(field: CategoryField) -> str:
return field.definition.label if field.definition else field.label
def _definition_readonly(field: CategoryField) -> bool:
return bool(field.definition.readonly) if field.definition else bool(getattr(field, "readonly", False))
def _custom_value_get(asset: Asset, definition: FieldDefinition, db: Session):
row = db.query(AssetFieldValue).filter(AssetFieldValue.asset_id == asset.id, AssetFieldValue.field_definition_id == definition.id).first()
if not row: return None
return row.value_json if definition.data_type == "json" else row.value_boolean if definition.data_type == "boolean" else row.value_datetime if definition.data_type == "datetime" else row.value_date if definition.data_type == "date" else row.value_number if definition.data_type in {"integer","decimal"} else row.value_text
def _custom_value_set(asset: Asset, definition: FieldDefinition, raw_value, db: Session) -> None:
row = db.query(AssetFieldValue).filter(AssetFieldValue.asset_id == asset.id, AssetFieldValue.field_definition_id == definition.id).first()
if not row:
row = AssetFieldValue(asset_id=asset.id, field_definition_id=definition.id); db.add(row)
row.value_text=row.value_number=row.value_date=None; row.value_boolean=None; row.value_datetime=None; row.value_json=None
value = None if raw_value in (None, "") else raw_value
if definition.data_type == "boolean": row.value_boolean = str(value).lower() in {"1","true","on","yes"} if value is not None else None
elif definition.data_type in {"integer","decimal"}: row.value_number = str(value) if value is not None else None
elif definition.data_type == "date": row.value_date = str(value) if value is not None else None
elif definition.data_type == "datetime": row.value_datetime = datetime.fromisoformat(str(value)) if value else None
elif definition.data_type == "json":
import json; row.value_json = json.loads(str(value)) if value else None
else: row.value_text = str(value) if value is not None else None
templates.env.globals["field_label"] = _definition_label
templates.env.globals["field_readonly"] = _definition_readonly
SOFTWARE_INVENTORY_WINDOWS_SCRIPT = r'''
$paths = @(
'HKLM:\Software\Microsoft\Windows\CurrentVersion\Uninstall\*',
'HKLM:\Software\WOW6432Node\Microsoft\Windows\CurrentVersion\Uninstall\*'
)
$software = Get-ItemProperty $paths -ErrorAction SilentlyContinue |
Where-Object { $_.DisplayName } |
ForEach-Object {
[PSCustomObject]@{
name = [string]$_.DisplayName
version = [string]$_.DisplayVersion
publisher = [string]$_.Publisher
install_date = [string]$_.InstallDate
architecture = $(if ($_.PSPath -like '*WOW6432Node*') { 'x86' } else { 'x64' })
}
} | Sort-Object name, version -Unique
Write-JobLog ("Software inventory completed; entries=" + $software.Count)
$JobResult = @{ software = $software }
'''
SOFTWARE_INVENTORY_LINUX_SCRIPT = r'''
import os
import subprocess
items = []
if os.path.exists("/usr/bin/dpkg-query"):
output = subprocess.check_output(
["dpkg-query", "-W", "-f=${binary:Package}\t${Version}\t${Maintainer}\n"],
text=True,
errors="replace",
)
for line in output.splitlines():
parts = line.split("\t")
items.append({
"name": parts[0],
"version": parts[1] if len(parts) > 1 else "",
"publisher": parts[2] if len(parts) > 2 else "",
"architecture": "",
"install_date": None,
})
elif os.path.exists("/usr/bin/rpm"):
output = subprocess.check_output(
["rpm", "-qa", "--qf", "%{NAME}\t%{VERSION}-%{RELEASE}\t%{VENDOR}\t%{ARCH}\n"],
text=True,
errors="replace",
)
for line in output.splitlines():
parts = line.split("\t")
items.append({
"name": parts[0],
"version": parts[1] if len(parts) > 1 else "",
"publisher": parts[2] if len(parts) > 2 else "",
"architecture": parts[3] if len(parts) > 3 else "",
"install_date": None,
})
elif os.path.exists("/sbin/apk") or os.path.exists("/usr/sbin/apk"):
output = subprocess.check_output(["apk", "info", "-v"], text=True, errors="replace")
for line in output.splitlines():
items.append({"name": line, "version": "", "publisher": "", "architecture": "", "install_date": None})
print("Software inventory completed; entries=" + str(len(items)))
job_result = {"software": items}
'''
def _seed_job_definitions(db: Session) -> None:
defaults = [
dict(name='Software-Inventur Windows', system_key='software_inventory_windows', platform='windows', interpreter='powershell', inline_script=SOFTWARE_INVENTORY_WINDOWS_SCRIPT),
dict(name='Software-Inventur Linux', system_key='software_inventory_linux', platform='linux', interpreter='python', inline_script=SOFTWARE_INVENTORY_LINUX_SCRIPT),
]
for data in defaults:
row=db.query(JobDefinition).filter(JobDefinition.system_key==data['system_key']).first()
if row is None:
row=JobDefinition(description='Systemdefinition für die Softwareinventur. Enthält ausschließlich den fachlichen Skriptcode; der technische Rahmen wird zur Laufzeit erzeugt.', source_type='inline', enabled=True, is_system=True, timeout_seconds=600, success_codes='0', **data)
db.add(row)
else:
# System definitions are maintained by the application. This one-time
# normalization removes previously stored callback/wrapper code.
row.name=data['name']; row.platform=data['platform']; row.interpreter=data['interpreter']
row.source_type='inline'; row.inline_script=data['inline_script']; row.source_path=None
row.description='Systemdefinition für die Softwareinventur. Enthält ausschließlich den fachlichen Skriptcode; der technische Rahmen wird zur Laufzeit erzeugt.'
row.enabled=True; row.is_system=True; row.timeout_seconds=600; row.success_codes='0'
db.commit()
def _job_definition_snapshot(row: JobDefinition) -> dict:
return {**{k:getattr(row,k) for k in ('name','description','job_kind','configuration','system_key','platform','interpreter','source_type','source_path','inline_script','arguments','timeout_seconds','success_codes','enabled','is_system','revision')}, 'parameters':[{'parameter_key':x.parameter_key,'label':x.label,'description':x.description,'data_type':x.data_type,'default_value':x.default_value,'required':x.required,'sort_order':x.sort_order} for x in row.parameters]}
def _job_preview_values() -> dict[str, str]:
return {
"JobId": "12345",
"CallbackUrl": "https://assetmanager.example.invalid/api/software-jobs/12345/callback?token=preview-token",
"AssetName": "DEMO-PC-001",
"Hostname": "DEMO-PC-001",
"IPAddress": "192.0.2.25",
"CurrentUser": "demo.user",
"SerialNumber": "DEMO-SERIAL-001",
"Manufacturer": "Example Manufacturer",
"Model": "Example Model",
"Department": "IT",
"Location": "Head Office",
"MeshNodeId": "node/demo-preview"
}
def _replace_job_placeholders(script: str, values: dict[str, str] | None = None) -> str:
rendered = str(script or "")
for key, value in (values or _job_preview_values()).items():
rendered = rendered.replace("{{" + key + "}}", str(value))
return rendered
def _validate_job_definition(row: JobDefinition, request: Request | None = None) -> str:
def tr(key: str, default: str, **values) -> str:
return _translate_request(request, key, default, **values) if request else default.format(**values)
errors=[]; notes=[]
script=build_registry_user_script(row.configuration) if row.job_kind=='registry' else ((row.inline_script or '') if row.source_type=='inline' else '')
if row.job_kind=='registry' and not (row.configuration or {}).get('registry_entries'): errors.append(tr('jobdefs.validation.registry_empty','At least one registry entry is required.'))
if row.job_kind!='registry' and row.source_type=='inline' and not script.strip(): errors.append(tr('jobdefs.validation.inline_empty','The inline script is empty.'))
if row.source_type=='mounted':
path=Path(row.source_path or '')
if not path.is_absolute(): errors.append(tr('jobdefs.validation.mounted_absolute','The mounted path must be absolute.'))
elif not path.is_file(): errors.append(tr('jobdefs.validation.file_missing','File not found: {path}',path=path))
else:
script=path.read_text(encoding='utf-8',errors='replace')
notes.append(tr('jobdefs.validation.file_readable','File readable: {path} ({count} characters)',path=path,count=len(script)))
if row.source_type=='unc':
path=str(row.source_path or '')
if not path.startswith('\\\\'): errors.append(tr('jobdefs.validation.unc_prefix','A UNC path must begin with \\\\.'))
else: notes.append(tr('jobdefs.validation.unc_ok','UNC syntax is valid. Reachability is checked on the target device.'))
rendered = _replace_job_placeholders(script)
if rendered and row.interpreter=='python':
try: compile(rendered,'<job-definition>','exec');notes.append(tr('jobdefs.validation.python_ok','Python syntax is valid.'))
except SyntaxError as exc: errors.append(tr('jobdefs.validation.python_error','Python syntax error on line {line}: {message}',line=exc.lineno,message=exc.msg))
elif rendered and row.interpreter=='bash' and shutil.which('bash'):
with tempfile.NamedTemporaryFile('w',suffix='.sh',delete=False,encoding='utf-8') as f:f.write(rendered);name=f.name
try:
r=subprocess.run(['bash','-n',name],capture_output=True,text=True,timeout=10)
if r.returncode: errors.append(r.stderr.strip())
else: notes.append(tr('jobdefs.validation.bash_ok','Bash syntax is valid.'))
finally: Path(name).unlink(missing_ok=True)
elif rendered and row.interpreter in {'powershell','cmd'}: notes.append(tr('jobdefs.validation.windows_later','Content is present. Full syntax validation is performed later on a Windows test device.'))
if errors: return 'ERROR\n'+'\n'.join('- '+x for x in errors)+'\n\n'+'\n'.join(x for x in notes if x)
return 'OK\n'+'\n'.join(x for x in notes if x)
def _seed_software_inventory_package(db: Session) -> SoftwarePackage:
package = db.query(SoftwarePackage).filter(SoftwarePackage.name == "Software-Inventur").first()
if package is None:
package = SoftwarePackage(
name="Software-Inventur",
description="Liest installierte Software über den vorhandenen MeshCentral-Agenten aus.",
package_type="inventory",
enabled=True,
is_system=True,
callback_timeout_minutes=30,
)
db.add(package)
db.commit()
db.refresh(package)
return package
def _seed_generic_job_package(db: Session) -> SoftwarePackage:
package = db.query(SoftwarePackage).filter(SoftwarePackage.name == "Jobdefinition").first()
if package is None:
package = SoftwarePackage(
name="Jobdefinition",
description="Technisches Systempaket für ausführbare Jobdefinitionen.",
package_type="job_definition",
enabled=True,
is_system=True,
callback_timeout_minutes=30,
)
db.add(package)
db.commit()
db.refresh(package)
return package
@app.on_event("startup")
def startup():
Base.metadata.create_all(bind=engine)
apply_lightweight_migrations()
from .database import SessionLocal
db = SessionLocal()
try:
_ensure_local_emergency_admin(db)
seed_i18n(db)
_seed_field_definitions(db)
_seed_mesh_field_mappings(db)
_seed_software_inventory_package(db)
_seed_generic_job_package(db)
_seed_job_definitions(db)
backfill_asset_job_states(db)
default_statuses = ["Aktiv", "Defekt", "im Lager", "Ausgegeben", "Verschrottet"]
if db.query(StatusOption).count() == 0:
for order, status_name in enumerate(default_statuses):
db.add(StatusOption(name=status_name, sort_order=order, active=True))
db.commit()
if not db.query(StatusOption).filter(StatusOption.use_for_issue.is_(True)).first():
item = db.query(StatusOption).filter(StatusOption.name == "Ausgegeben").first()
if item: item.use_for_issue = True
if not db.query(StatusOption).filter(StatusOption.use_for_return.is_(True)).first():
item = db.query(StatusOption).filter(StatusOption.name == "im Lager").first()
if item: item.use_for_return = True
db.commit()
if db.query(Category).count() == 0:
computer = Category(name="Computer", description="Computer und Notebooks aus MeshCentral oder manueller Erfassung")
db.add(computer)
db.flush()
for order, field_name in enumerate(COMPUTER_FIELDS):
db.add(CategoryField(
category_id=computer.id,
field_name=field_name,
label=ASSET_FIELDS[field_name],
active=True,
required=field_name == "name",
show_in_list=field_name in {"asset_tag", "name", "hostname", "operating_system", "os_patch_level", "mac_address", "status"},
sort_order=order,
))
db.commit()
else:
computer = db.query(Category).options(joinedload(Category.visible_fields)).filter(Category.name == "Computer").first()
if computer:
existing = {field.field_name for field in computer.visible_fields}
next_order = max((field.sort_order for field in computer.visible_fields), default=-1) + 1
important_list_fields = {"hostname", "manufacturer", "model", "serial_number", "operating_system", "ip_address", "mac_address", "storage_details", "gpu_name", "tpm_version", "antivirus", "mesh_group"}
changed = False
for field_name in COMPUTER_FIELDS:
if field_name not in existing:
db.add(CategoryField(
category_id=computer.id, field_name=field_name, label=ASSET_FIELDS[field_name],
active=True, required=False, show_in_list=field_name in important_list_fields,
sort_order=next_order,
))
next_order += 1
changed = True
if changed:
db.commit()
_seed_field_definitions(db)
# Frühere Excel-Importe konnten Datumswerte als
# "YYYY-MM-DD 00:00:00" speichern. Browser zeigen solche Werte in
# <input type="date"> leer an. Bestehende System-Datumsfelder werden
# deshalb einmalig auf YYYY-MM-DD normalisiert.
date_repairs = 0
for asset in db.query(Asset).filter(
(Asset.purchase_date.like("% %")) |
(Asset.purchase_date.like("%T%")) |
(Asset.warranty_until.like("% %")) |
(Asset.warranty_until.like("%T%")) |
(Asset.assigned_on.like("% %")) |
(Asset.assigned_on.like("%T%"))
).all():
for field_name in ("purchase_date", "warranty_until", "assigned_on"):
old_value = getattr(asset, field_name, None)
new_value = _normalize_date_value(old_value)
if old_value and new_value and old_value != new_value:
setattr(asset, field_name, new_value)
date_repairs += 1
if date_repairs:
db.commit()
logger.info("%s ältere Datumswerte auf YYYY-MM-DD normalisiert", date_repairs)
finally:
db.close()
@app.on_event("startup")
def start_software_timeout_service():
global SOFTWARE_TIMEOUT_THREAD
if SOFTWARE_TIMEOUT_THREAD is None or not SOFTWARE_TIMEOUT_THREAD.is_alive():
SOFTWARE_TIMEOUT_STOP_EVENT.clear()
SOFTWARE_TIMEOUT_THREAD = threading.Thread(
target=_software_job_timeout_loop,
args=(SOFTWARE_TIMEOUT_STOP_EVENT,),
daemon=True,
name="assetmanager-software-timeouts",
)
SOFTWARE_TIMEOUT_THREAD.start()
@app.on_event("shutdown")
def stop_software_timeout_service():
SOFTWARE_TIMEOUT_STOP_EVENT.set()
@app.on_event("startup")
def start_mesh_presence_service():
global PRESENCE_THREAD
if PRESENCE_THREAD is None or not PRESENCE_THREAD.is_alive():
PRESENCE_STOP_EVENT.clear()
PRESENCE_THREAD = threading.Thread(
target=mesh_presence_loop,
args=(PRESENCE_STOP_EVENT,),
daemon=True,
name="assetmanager-mesh-presence",
)
PRESENCE_THREAD.start()
@app.on_event("shutdown")
def stop_mesh_presence_service():
PRESENCE_STOP_EVENT.set()
def delete_asset_image(image_path: str | None) -> None:
"""Löscht nur individuell hochgeladene Asset-Bilder aus dem Upload-Ordner."""
if not image_path or not image_path.startswith("/static/uploads/"):
return
filename = Path(image_path).name
target = UPLOAD_DIR / filename
try:
target.unlink(missing_ok=True)
except OSError:
pass
def save_upload(upload: UploadFile | None) -> str | None:
if not upload or not upload.filename:
return None
suffix = Path(upload.filename).suffix.lower()
if suffix not in {".png", ".jpg", ".jpeg", ".webp", ".gif"}:
raise HTTPException(400, "Nur PNG, JPG, WEBP oder GIF sind erlaubt.")
filename = f"{uuid.uuid4().hex}{suffix}"
target = UPLOAD_DIR / filename
with target.open("wb") as buffer:
shutil.copyfileobj(upload.file, buffer)
return f"/static/uploads/{filename}"
def _display_value(value) -> str:
if value is None:
return ""
if isinstance(value, (dict, list)):
import json
return json.dumps(value, ensure_ascii=False, sort_keys=True)
return str(value)
def _record_asset_history(db: Session, asset: Asset, changes: dict, source: str = "manual", changed_by: str | None = None) -> None:
if not changes:
return
db.add(AssetHistory(asset_id=asset.id, source=source, changed_by=changed_by, changes=changes))
def _xlsx_response(filename: str, headers: list[str], rows: list[list]) -> Response:
"""Erzeugt eine echte XLSX-Datei und liefert vollständige Bytes statt eines Streams.
Das verhindert beschädigte Downloads bei einzelnen Reverse-Proxies und Browsern.
"""
workbook = Workbook()
sheet = workbook.active
sheet.title = "Export"
sheet.freeze_panes = "A2"
sheet.auto_filter.ref = f"A1:{sheet.cell(row=1, column=max(1, len(headers))).column_letter}1"
sheet.append(headers)
header_fill = PatternFill("solid", fgColor="D9EAF7")
for cell in sheet[1]:
cell.font = Font(bold=True)
cell.fill = header_fill
cell.alignment = Alignment(vertical="top")
for row in rows:
sheet.append([_display_value(value) for value in row])
for column in sheet.columns:
max_length = min(max((len(str(cell.value or "")) for cell in column), default=10) + 2, 60)
sheet.column_dimensions[column[0].column_letter].width = max_length
output = io.BytesIO()
workbook.save(output)
payload = output.getvalue()
ascii_name = "".join(ch if ch.isascii() and (ch.isalnum() or ch in "._-") else "_" for ch in filename)
disposition = f"attachment; filename={ascii_name}; filename*=UTF-8''{quote(filename)}"
return Response(
content=payload,
media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
headers={
"Content-Disposition": disposition,
"Content-Length": str(len(payload)),
"X-Content-Type-Options": "nosniff",
"Cache-Control": "no-store",
},
)
def _import_field_map(category: Category) -> dict[str, str]:
result: dict[str, str] = {"asset-id": "asset_id", "asset_id": "asset_id", "id": "asset_id", "bezeichnung": "name", "name": "name", "kategorie": "category"}
for field in category.visible_fields:
if field.active and hasattr(Asset, field.field_name):
result[field.label.strip().casefold()] = field.field_name
result[field.field_name.strip().casefold()] = field.field_name
return result
def _excel_cell_value(value) -> str | None:
if value is None:
return None
if isinstance(value, datetime):
return value.strftime("%Y-%m-%d %H:%M:%S")
if isinstance(value, date):
return value.isoformat()
text = str(value).strip()
return text or None
def _normalize_date_value(value: Any) -> str | None:
"""Normalisiert Excel-/Datenbankwerte für HTML-Datumsfelder auf YYYY-MM-DD."""
if value in (None, ""):
return None
if isinstance(value, datetime):
return value.date().isoformat()
if isinstance(value, date):
return value.isoformat()
text = str(value).strip()
if not text:
return None
# ISO-Datum oder ISO-Zeitstempel: Der Datumsteil sind die ersten 10 Zeichen.
if re.match(r"^\d{4}-\d{2}-\d{2}(?:[ T].*)?$", text):
return text[:10]
# Deutsche Schreibweise aus manuellen Excel-Dateien.
match = re.match(r"^(\d{1,2})\.(\d{1,2})\.(\d{4})(?:\s+.*)?$", text)
if match:
day, month, year = match.groups()
return f"{year}-{int(month):02d}-{int(day):02d}"
return text
def _normalize_datetime_value(value: Any) -> str | None:
"""Normalisiert Werte für <input type=datetime-local>."""
if value in (None, ""):
return None
if isinstance(value, datetime):
return value.strftime("%Y-%m-%dT%H:%M")
if isinstance(value, date):
return f"{value.isoformat()}T00:00"
text = str(value).strip()
if not text:
return None
text = text.replace(" ", "T")
if re.match(r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}", text):
return text[:16]
if re.match(r"^\d{4}-\d{2}-\d{2}$", text):
return f"{text}T00:00"
return text
def _html_field_value(value: Any, data_type: str | None = None) -> str:
if data_type == "date":
return _normalize_date_value(value) or ""
if data_type == "datetime":
return _normalize_datetime_value(value) or ""
return _display_value(value)
templates.env.globals["html_field_value"] = _html_field_value
def _request_language(request: Request | None) -> str:
if request is None:
return "en"
user = request.session.get("user") or {}
return str(user.get("language_code") or request.session.get("language_code") or load_config().get("general", {}).get("default_language", "en")).lower()
def _localized_text_key(value: Any, language_code: str = "en") -> str:
text = str(value or "").strip().casefold()
# German users expect umlauts near their expanded base letters instead of after Z.
if language_code.startswith("de"):
text = text.replace("ä", "ae").replace("ö", "oe").replace("ü", "ue").replace("ß", "ss")
return "".join(ch for ch in unicodedata.normalize("NFKD", text) if not unicodedata.combining(ch))
def _sort_localized(items: list[Any], attribute: str, request: Request | None) -> list[Any]:
language_code = _request_language(request)
return sorted(items, key=lambda item: _localized_text_key(getattr(item, attribute, ""), language_code))
def _active_status_options(db: Session, request: Request | None = None) -> list[StatusOption]:
items = db.query(StatusOption).filter(StatusOption.active.is_(True)).all()
return _sort_localized(items, "name", request)
def _asset_action_changes(asset: Asset, action: str, status_name: str, assigned_to: str | None = None, assigned_on: str | None = None) -> dict[str, dict[str, str]]:
changes: dict[str, dict[str, str]] = {}
if _display_value(asset.status) != status_name:
changes[ASSET_FIELDS["status"]] = {"old": _display_value(asset.status), "new": status_name}
asset.status = status_name
if action == "issue":
new_name = (assigned_to or "").strip() or None
new_date = (assigned_on or date.today().isoformat()).strip()
if _display_value(asset.assigned_to) != _display_value(new_name):
changes[ASSET_FIELDS["assigned_to"]] = {"old": _display_value(asset.assigned_to), "new": _display_value(new_name)}
asset.assigned_to = new_name
if _display_value(asset.assigned_on) != new_date:
changes[ASSET_FIELDS["assigned_on"]] = {"old": _display_value(asset.assigned_on), "new": new_date}
asset.assigned_on = new_date
else:
if asset.assigned_to:
changes[ASSET_FIELDS["assigned_to"]] = {"old": _display_value(asset.assigned_to), "new": ""}
asset.assigned_to = None
if asset.assigned_on:
changes[ASSET_FIELDS["assigned_on"]] = {"old": _display_value(asset.assigned_on), "new": ""}
asset.assigned_on = None
return changes
def _normalize_optional_text(value: str | None) -> str | None:
text = (value or "").strip()
return text or None
def _validate_asset_identifiers(db: Session, *, asset_tag: str | None, manufacturer: str | None, serial_number: str | None, exclude_id: int | None = None) -> None:
asset_tag = _normalize_optional_text(asset_tag)
manufacturer = _normalize_optional_text(manufacturer)
serial_number = _normalize_optional_text(serial_number)
if asset_tag:
query = db.query(Asset).filter(Asset.asset_tag == asset_tag)
if exclude_id is not None:
query = query.filter(Asset.id != exclude_id)
if query.first():
raise HTTPException(400, f"Die Inventarnummer {asset_tag} ist bereits vergeben.")
if manufacturer and serial_number:
query = db.query(Asset).filter(Asset.manufacturer == manufacturer, Asset.serial_number == serial_number)
if exclude_id is not None:
query = query.filter(Asset.id != exclude_id)
if query.first():
raise HTTPException(400, f"Die Seriennummer {serial_number} ist beim Hersteller {manufacturer} bereits vorhanden.")
def _would_create_parent_cycle(db: Session, asset_id: int, parent_id: int | None) -> bool:
seen = {asset_id}
current = parent_id
while current:
if current in seen:
return True
seen.add(current)
parent = db.query(Asset.parent_asset_id).filter(Asset.id == current).scalar()
current = parent
return False
def _list_fields_for_categories(categories: list[Category]) -> list[CategoryField]:
"""Bildet eine eindeutige, sortierte Vereinigungsmenge der konfigurierten Listenspalten."""
result: list[CategoryField] = []
seen: set[str] = set()
for category in categories:
for field in sorted(category.visible_fields, key=lambda item: item.sort_order):
if not field.active or not field.show_in_list or field.field_name in seen:
continue
result.append(field)
seen.add(field.field_name)
return result
CHART_FIELDS = {
"category": "Kategorie",
"status": "Status",
"operating_system": "Betriebssystem",
"manufacturer": "Hersteller",
"model": "Modell",
"location": "Standort",
"room_number": "Raumnummer",
"department": "Abteilung",
"mesh_group": "MeshCentral-Gruppe",
"antivirus": "Antivirus",
"tpm_version": "TPM-Version",
}
CHART_TYPES = {
"donut": "Donutdiagramm",
"pie": "Kuchendiagramm",
"bar": "Balkendiagramm",
"stacked_bar": "Gestapeltes Balkendiagramm",
}
def _parse_value_mappings(text_value: str) -> list[dict[str, str]]:
"""Liest Regeln im Format Suchtext => Anzeigename (eine Regel je Zeile)."""
mappings: list[dict[str, str]] = []
for line_number, raw_line in enumerate((text_value or "").splitlines(), start=1):
line = raw_line.strip()
if not line or line.startswith("#"):
continue
separator = "=>" if "=>" in line else "=" if "=" in line else None
if not separator:
raise HTTPException(400, f"Ungültige Wertzuordnung in Zeile {line_number}. Erwartet: Suchtext => Anzeigename")
search, label = (part.strip() for part in line.split(separator, 1))
if not search or not label:
raise HTTPException(400, f"Unvollständige Wertzuordnung in Zeile {line_number}.")
mappings.append({"search": search, "label": label})
return mappings
def _value_mappings_text(chart: ChartDefinition | None) -> str:
if not chart:
return ""
return "\n".join(
f"{entry.get('search', '')} => {entry.get('label', '')}"
for entry in (chart.value_mappings or [])
if isinstance(entry, dict) and entry.get("search") and entry.get("label")
)
def _automatic_chart_value(field: str, value: str) -> str:
"""Fasst technische Rohwerte für Diagramme zusammen, ohne die Assetdaten zu verändern."""
text = value.strip()
lower = text.casefold()
if field == "operating_system":
if "windows server" in lower:
match = re.search(r"windows server\s*(20\d{2})", text, re.IGNORECASE)
return f"Windows Server {match.group(1)}" if match else "Windows Server"
if "windows 11" in lower:
return "Windows 11"
if "windows 10" in lower:
return "Windows 10"
if "ubuntu" in lower:
return "Ubuntu"
if "debian" in lower:
return "Debian"
if "red hat" in lower or "rhel" in lower:
return "Red Hat Enterprise Linux"
if "centos" in lower:
return "CentOS"
if "fedora" in lower:
return "Fedora"
if "macos" in lower or "mac os" in lower:
return "macOS"
if "linux" in lower:
return "Linux"
if field == "manufacturer":
if lower.startswith("hewlett") or lower in {"hp", "hp inc.", "hp inc"}:
return "HP"
if lower.startswith("dell"):
return "Dell"
if lower.startswith("lenovo"):
return "Lenovo"
if lower.startswith("microsoft"):
return "Microsoft"
return text
def _transform_chart_value(chart: ChartDefinition, field: str, value: str) -> str:
# Eigene Regeln haben Vorrang; gesucht wird ohne Beachtung der Groß-/Kleinschreibung.
lower = value.casefold()
for entry in chart.value_mappings or []:
if not isinstance(entry, dict):
continue
search = str(entry.get("search") or "").strip()
label = str(entry.get("label") or "").strip()
if search and label and search.casefold() in lower:
return label
if (chart.value_mode or "raw") == "automatic":
return _automatic_chart_value(field, value)
return value
def _chart_value(asset: Asset, chart: ChartDefinition, field: str, show_empty: bool) -> str | None:
if field == "category":
value = asset.category.name if asset.category else None
else:
value = getattr(asset, field, None)
if isinstance(value, str):
value = value.strip()
if value in (None, ""):
return "Nicht angegeben" if show_empty else None
return _transform_chart_value(chart, field, str(value))
def _chart_payload(chart: ChartDefinition, assets: list[Asset]) -> dict:
grouped: dict[str, dict[str, int] | int] = {}
stack_values: set[str] = set()
excluded_statuses = {str(value).casefold() for value in (chart.excluded_statuses or []) if value not in (None, "")}
for asset in assets:
if excluded_statuses and str(asset.status or "").casefold() in excluded_statuses:
continue
group = _chart_value(asset, chart, chart.group_field, chart.show_empty)
if group is None:
continue
if chart.chart_type == "stacked_bar" and chart.stack_field:
stack = _chart_value(asset, chart, chart.stack_field, chart.show_empty)
if stack is None:
continue
stack_values.add(stack)
bucket = grouped.setdefault(group, {})
bucket[stack] = bucket.get(stack, 0) + 1
else:
grouped[group] = int(grouped.get(group, 0)) + 1
def total(item):
value = item[1]
return sum(value.values()) if isinstance(value, dict) else value
items = list(grouped.items())
if chart.sort_mode == "alpha":
items.sort(key=lambda item: item[0].casefold())
else:
items.sort(key=total, reverse=True)
limit = max(1, min(int(chart.max_items or 12), 50))
if len(items) > limit and chart.chart_type != "stacked_bar":
kept, rest = items[:limit-1], items[limit-1:]
kept.append(("Sonstige", sum(int(v) for _, v in rest)))
items = kept
else:
items = items[:limit]
maximum = max((total(item) for item in items), default=1)
total_count = sum(total(item) for item in items)
donut_stops = []
current = 0.0
for index, (_, value) in enumerate(items):
count = total(("", value))
end = current + (count / total_count * 100 if total_count else 0)
donut_stops.append(f"var(--chart-color-{index % 12}) {current:.2f}% {end:.2f}%")
current = end
return {
"chart": chart,
"items": items,
"stack_values": sorted(stack_values, key=str.casefold),
"maximum": maximum,
"total": total_count,
"donut_gradient": ", ".join(donut_stops) or "#e5e7eb 0 100%",
}
def _seed_default_charts(db: Session) -> None:
defaults = [
{"name": "Betriebssysteme", "chart_type": "donut", "group_field": "operating_system", "value_mode": "automatic", "dashboard_visible": True, "max_items": 10},
{"name": "Geräte nach Kategorie und Status", "chart_type": "stacked_bar", "group_field": "category", "stack_field": "status", "dashboard_visible": True, "max_items": 15},
{"name": "Hersteller", "chart_type": "bar", "group_field": "manufacturer", "dashboard_visible": False, "max_items": 10},
{"name": "Statusübersicht", "chart_type": "donut", "group_field": "status", "dashboard_visible": False, "max_items": 10},
]
changed = False
for values in defaults:
chart = db.query(ChartDefinition).filter(ChartDefinition.name == values["name"]).first()
if chart is None:
chart = ChartDefinition(**values, is_system=True)
db.add(chart)
changed = True
elif not chart.is_system:
chart.is_system = True
changed = True
if changed:
db.commit()
@app.get("/")
def dashboard(request: Request, db: Session = Depends(get_db)):
categories = db.query(Category).options(joinedload(Category.visible_fields)).all()
assets = (
_apply_asset_access(db.query(Asset), request)
.options(joinedload(Asset.category).joinedload(Category.visible_fields))
.order_by(Asset.updated_at.desc())
.limit(10)
.all()
)
recent_categories = []
seen_category_ids = set()
for asset in assets:
if asset.category_id not in seen_category_ids:
recent_categories.append(asset.category)
seen_category_ids.add(asset.category_id)
fields = _list_fields_for_categories(recent_categories)
_seed_default_charts(db)
chart_assets = _apply_asset_access(db.query(Asset), request).options(joinedload(Asset.category)).all()
dashboard_charts = [_chart_payload(chart, chart_assets) for chart in db.query(ChartDefinition).filter(ChartDefinition.dashboard_visible.is_(True)).order_by(ChartDefinition.id).all()]
return templates.TemplateResponse(
"dashboard.html",
{"request": request, "categories": categories, "assets": assets, "fields": fields, "is_admin": _is_admin(request), "issue_status_name": (_action_status(db, "issue").name if _action_status(db, "issue") else None), "dashboard_charts": dashboard_charts},
)
BULK_EDIT_EXCLUDED_FIELDS = {"asset_tag", "serial_number", "mesh_node_id"}
def _bulk_edit_fields_for_assets(assets: list[Asset]) -> list[dict]:
"""Return active fields shared by all selected asset categories."""
if not assets:
return []
category_maps = []
for asset in assets:
fields = {
field.field_name: _definition_label(field)
for field in asset.category.visible_fields
if field.active and hasattr(asset, field.field_name) and field.field_name not in BULK_EDIT_EXCLUDED_FIELDS and not _definition_readonly(field) and (not field.definition or field.definition.bulk_edit)
}
category_maps.append(fields)
common = set(category_maps[0])
for item in category_maps[1:]:
common &= set(item)
return [{"field_name": name, "label": category_maps[0][name]} for name in sorted(common, key=lambda n: category_maps[0][n].casefold())]
def _selected_asset_ids(values: list[str]) -> list[int]:
result = []
for value in values:
try:
number = int(value)
except (TypeError, ValueError):
continue
if number > 0 and number not in result:
result.append(number)
return result
ASSET_JOB_FILTER_STATUSES = (
"never_run",
"created",
"sending",
"sent",
"running",
"success",
"partial",
"failed",
"timeout",
"cancelled",
)
def _resolve_asset_job_filter(db: Session, selection: str) -> dict[str, Any] | None:
selection = str(selection or "").strip()
if selection == "software_inventory":
return {
"selection": selection,
"state_key": filter_state_key("software_inventory"),
"job_type": "software_inventory",
"definition": None,
"definition_id": None,
"current_revision": None,
"label": "Softwareinventur",
"label_key": "jobs.type.software_inventory",
}
if not selection.startswith("definition:"):
return None
try:
definition_id = int(selection.split(":", 1)[1])
except (TypeError, ValueError):
return None
definition = db.get(JobDefinition, definition_id)
if not definition:
return None
return {
"selection": selection,
"state_key": filter_state_key("job_definition", definition.id),
"job_type": "job_definition",
"definition": definition,
"definition_id": definition.id,
"current_revision": int(definition.revision or 0),
"label": definition.name,
"label_key": None,
}
def _normalize_job_filter_selections(value: Any) -> list[str]:
"""Normalize repeated query/form values while preserving their order."""
if value is None:
raw_values: list[Any] = []
elif isinstance(value, (list, tuple, set)):
raw_values = list(value)
else:
raw_values = [value]
result: list[str] = []
for raw in raw_values:
for item in str(raw or "").split(","):
selection = item.strip()
if not selection or selection in result:
continue
if selection == "software_inventory" or re.fullmatch(r"definition:\d+", selection):
result.append(selection)
return result
def _legacy_job_filter_exclusions(condition: str) -> list[str]:
"""Preserve links generated by 0.5.5.17 while switching to exclusions."""
include_by_condition = {
"successful": {"success"},
"not_successful": set(ASSET_JOB_FILTER_STATUSES) - {"success"},
"never_run": {"never_run"},
"last_failed": {"failed"},
"last_partial": {"partial"},
"timeout": {"timeout"},
"callback_pending": {"sent"},
"active": {"created", "sending", "running"},
}
included = include_by_condition.get(str(condition or "").strip())
if included is None:
return []
return [status for status in ASSET_JOB_FILTER_STATUSES if status not in included]
def _normalize_job_filter_exclusions(value: Any) -> list[str]:
"""Normalize repeated form/query values or a comma-separated status list."""
if value is None:
raw_values: list[Any] = []
elif isinstance(value, (list, tuple, set)):
raw_values = list(value)
else:
raw_values = [value]
requested: set[str] = set()
for raw in raw_values:
for item in str(raw or "").split(","):
status = item.strip().lower()
if status in ASSET_JOB_FILTER_STATUSES:
requested.add(status)
return [status for status in ASSET_JOB_FILTER_STATUSES if status in requested]
def _normalize_asset_job_filter(
db: Session,
*,
selection: Any,
excluded_statuses: Any = None,
revision_mode: str = "all",
) -> dict[str, Any] | None:
selections = _normalize_job_filter_selections(selection)
specs: list[dict[str, Any]] = []
state_keys: set[str] = set()
for selected_value in selections:
spec = _resolve_asset_job_filter(db, selected_value)
if spec is None or spec["state_key"] in state_keys:
continue
specs.append(spec)
state_keys.add(spec["state_key"])
if not specs:
return None
revision_mode = "current" if str(revision_mode or "").strip() == "current" else "all"
if not any(spec["job_type"] == "job_definition" for spec in specs):
revision_mode = "all"
return {
"selections": [spec["selection"] for spec in specs],
"specs": specs,
"state_keys": [spec["state_key"] for spec in specs],
"labels": [spec["label"] for spec in specs],
"startable_specs": [spec for spec in specs if spec.get("definition") is None or spec["definition"].enabled],
"excluded_statuses": _normalize_job_filter_exclusions(excluded_statuses),
"revision_mode": revision_mode,
}
def _asset_job_filter_url(job_filter: dict[str, Any], category_id: int | None = None, **toast: str) -> str:
params: list[tuple[str, Any]] = []
params.extend(("job_filter", selection) for selection in job_filter.get("selections") or [])
params.append(("job_filter_revision", job_filter["revision_mode"]))
params.extend(("job_filter_exclude", status) for status in job_filter.get("excluded_statuses") or [])
if category_id:
params.append(("category_id", category_id))
params.extend((key, value) for key, value in toast.items() if value)
return "/assets?" + urllib.parse.urlencode(params, doseq=True)
def _asset_job_filter_view(state: AssetJobState | None, job_filter: dict[str, Any] | None = None) -> dict[str, Any]:
effective_never_run = state is None
if state is not None and job_filter and job_filter.get("revision_mode") == "current":
current_revision = job_filter.get("current_revision")
effective_never_run = current_revision is not None and state.definition_revision != current_revision
if state is None:
return {
"status": "never_run",
"last_execution": None,
"last_success": None,
"definition_revision": None,
"last_success_revision": None,
"execution_count": 0,
}
return {
"status": "never_run" if effective_never_run else (state.last_status or "created"),
"last_execution": state.last_started_at or state.last_sent_at or state.last_status_at or state.last_created_at,
"last_success": state.last_success_at,
"definition_revision": state.definition_revision,
"last_success_revision": state.last_success_revision,
"execution_count": int(state.execution_count or 0),
}
def _combined_asset_job_filter_view(
states_by_key: dict[str, AssetJobState],
job_filter: dict[str, Any],
) -> dict[str, Any]:
"""Combine alternative selected jobs using the most recently executed state.
This is intentional for cases such as the generic software inventory job
and platform-specific inventory definitions: an asset counts as successful
when its latest execution among the selected alternatives was successful.
"""
candidates: list[tuple[datetime, dict[str, Any]]] = []
latest_success: datetime | None = None
total_executions = 0
for spec in job_filter.get("specs") or []:
state = states_by_key.get(spec["state_key"])
view = _asset_job_filter_view(
state,
{
"revision_mode": job_filter.get("revision_mode"),
"current_revision": spec.get("current_revision"),
},
)
total_executions += int(view.get("execution_count") or 0)
if state is not None:
success_is_current = True
if job_filter.get("revision_mode") == "current" and spec.get("current_revision") is not None:
success_is_current = state.last_success_revision == spec.get("current_revision")
if success_is_current and state.last_success_at and (latest_success is None or state.last_success_at > latest_success):
latest_success = state.last_success_at
if view["status"] == "never_run":
continue
timestamp = view.get("last_execution") or (state.updated_at if state is not None else None) or datetime.min
view.update({
"job_label": spec["label"],
"job_selection": spec["selection"],
})
candidates.append((timestamp, view))
if not candidates:
return {
"status": "never_run",
"last_execution": None,
"last_success": latest_success,
"definition_revision": None,
"last_success_revision": None,
"execution_count": total_executions,
"job_label": None,
"job_selection": None,
}
result = max(candidates, key=lambda item: item[0])[1]
result["last_success"] = latest_success
result["execution_count"] = total_executions
return result
def _filter_assets_by_job_states(
db: Session,
assets: list[Asset],
job_filter: dict[str, Any],
) -> tuple[list[Asset], dict[int, dict[str, Any]]]:
if not assets:
return [], {}
asset_ids = [asset.id for asset in assets]
rows = (
db.query(AssetJobState)
.filter(
AssetJobState.asset_id.in_(asset_ids),
AssetJobState.state_key.in_(job_filter.get("state_keys") or [""]),
)
.all()
)
states_by_asset: dict[int, dict[str, AssetJobState]] = {}
for state in rows:
states_by_asset.setdefault(state.asset_id, {})[state.state_key] = state
excluded = set(job_filter.get("excluded_statuses") or [])
filtered_assets: list[Asset] = []
views: dict[int, dict[str, Any]] = {}
for asset in assets:
view = _combined_asset_job_filter_view(states_by_asset.get(asset.id, {}), job_filter)
if view["status"] in excluded:
continue
filtered_assets.append(asset)
views[asset.id] = view
return filtered_assets, views
@app.get("/assets")
def assets_list(request: Request, category_id: int | None = None, db: Session = Depends(get_db)):
categories = (
db.query(Category)
.options(selectinload(Category.visible_fields).joinedload(CategoryField.definition))
.order_by(Category.name)
.all()
)
requested_exclusions = request.query_params.getlist("job_filter_exclude")
if not requested_exclusions and request.query_params.get("job_filter_condition"):
requested_exclusions = _legacy_job_filter_exclusions(request.query_params.get("job_filter_condition", ""))
job_filter = _normalize_asset_job_filter(
db,
selection=request.query_params.getlist("job_filter"),
excluded_statuses=requested_exclusions,
revision_mode=request.query_params.get("job_filter_revision", "all"),
)
query = _apply_asset_access(db.query(Asset), request)
selected = None
if category_id:
selected = next((category for category in categories if category.id == category_id), None)
if not selected:
raise HTTPException(404, "Kategorie nicht gefunden")
query = query.filter(Asset.category_id == category_id)
assets = (
query
.options(joinedload(Asset.category))
.order_by(func.lower(Asset.name), Asset.id)
.all()
)
job_state_by_asset: dict[int, dict[str, Any]] = {}
if job_filter:
assets, job_state_by_asset = _filter_assets_by_job_states(db, assets, job_filter)
fields = (
[field for field in sorted(selected.visible_fields, key=lambda item: item.sort_order) if field.active and field.show_in_list]
if selected
else _list_fields_for_categories(categories)
)
issue_status = _action_status(db, "issue")
enabled_job_definitions = (
db.query(JobDefinition)
.filter(JobDefinition.enabled.is_(True))
.order_by(func.lower(JobDefinition.name), JobDefinition.id)
.all()
if _is_admin(request)
else []
)
filter_job_definitions = (
db.query(JobDefinition)
.filter(JobDefinition.enabled.is_(True))
.order_by(func.lower(JobDefinition.name), JobDefinition.id)
.all()
if _is_admin(request)
else []
)
return templates.TemplateResponse(
"assets.html",
{
"request": request,
"assets": assets,
"categories": categories,
"selected": selected,
"fields": fields,
"is_admin": _is_admin(request),
"issue_status_name": issue_status.name if issue_status else None,
"total_assets": len(assets),
"job_definitions": enabled_job_definitions,
"job_filter_definitions": filter_job_definitions,
"job_filter": job_filter,
"job_filter_statuses": ASSET_JOB_FILTER_STATUSES,
"job_state_by_asset": job_state_by_asset,
"asset_platforms": {asset.id: detect_platform(asset) for asset in assets},
},
)
@app.get("/charts")
def charts_page(request: Request, db: Session = Depends(get_db)):
_seed_default_charts(db)
assets = _apply_asset_access(db.query(Asset), request).options(joinedload(Asset.category)).all()
charts = [_chart_payload(chart, assets) for chart in db.query(ChartDefinition).order_by(ChartDefinition.id).all()]
return templates.TemplateResponse("charts.html", {"request": request, "charts": charts, "is_admin": _is_admin(request)})
@app.get("/charts/manage")
def charts_manage(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
_seed_default_charts(db)
charts = db.query(ChartDefinition).order_by(ChartDefinition.name).all()
return templates.TemplateResponse("charts_manage.html", {"request": request, "charts": charts, "chart_fields": CHART_FIELDS, "chart_types": CHART_TYPES})
def _chart_form_context(request: Request, db: Session, chart: ChartDefinition | None = None) -> dict[str, Any]:
language_code = _request_language(request)
statuses = db.query(StatusOption).filter(StatusOption.active.is_(True)).all()
statuses.sort(key=lambda item: _localized_text_key(item.name, language_code))
return {
"request": request,
"chart": chart,
"chart_fields": CHART_FIELDS,
"chart_types": CHART_TYPES,
"value_mappings_text": _value_mappings_text(chart),
"status_options": statuses,
"excluded_statuses": set(chart.excluded_statuses or []) if chart else set(),
}
@app.get("/charts/new")
def chart_new(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
return templates.TemplateResponse("chart_form.html", _chart_form_context(request, db))
@app.post("/charts/new")
def chart_create(request: Request, name: str = Form(...), description: str = Form(""), chart_type: str = Form("bar"), group_field: str = Form("category"), stack_field: str = Form(""), dashboard_visible: str | None = Form(None), show_empty: str | None = Form(None), max_items: int = Form(12), sort_mode: str = Form("count_desc"), value_mode: str = Form("raw"), value_mappings_text: str = Form(""), excluded_statuses: list[str] = Form([]), db: Session = Depends(get_db)):
_require_admin(request)
if chart_type not in CHART_TYPES or group_field not in CHART_FIELDS:
raise HTTPException(400, "Ungültige Diagrammkonfiguration")
if chart_type == "stacked_bar" and stack_field not in CHART_FIELDS:
raise HTTPException(400, "Für ein gestapeltes Diagramm ist eine gültige Unterteilung erforderlich.")
if value_mode not in {"raw", "automatic", "custom"}:
raise HTTPException(400, "Ungültige Wertedarstellung")
mappings = _parse_value_mappings(value_mappings_text)
valid_statuses = {row.name for row in db.query(StatusOption).all()}
selected_statuses = [value for value in excluded_statuses if value in valid_statuses]
chart = ChartDefinition(name=name.strip(), description=description.strip() or None, chart_type=chart_type, group_field=group_field, stack_field=stack_field if chart_type == "stacked_bar" else None, dashboard_visible=dashboard_visible == "on", show_empty=show_empty == "on", max_items=max(1, min(max_items, 50)), sort_mode=sort_mode, value_mode=value_mode, value_mappings=mappings, excluded_statuses=selected_statuses, created_by=_changed_by(request), is_system=False)
db.add(chart)
try:
db.commit()
except IntegrityError:
db.rollback(); raise HTTPException(400, "Ein Diagramm mit diesem Namen existiert bereits.")
return RedirectResponse("/charts/manage?toast_success=Diagramm angelegt", status_code=303)
@app.get("/charts/{chart_id}/edit")
def chart_edit(chart_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
chart = db.get(ChartDefinition, chart_id)
if not chart: raise HTTPException(404, "Diagramm nicht gefunden")
return templates.TemplateResponse("chart_form.html", _chart_form_context(request, db, chart))
@app.post("/charts/{chart_id}/edit")
def chart_update(chart_id: int, request: Request, name: str = Form(...), description: str = Form(""), chart_type: str = Form("bar"), group_field: str = Form("category"), stack_field: str = Form(""), dashboard_visible: str | None = Form(None), show_empty: str | None = Form(None), max_items: int = Form(12), sort_mode: str = Form("count_desc"), value_mode: str = Form("raw"), value_mappings_text: str = Form(""), excluded_statuses: list[str] = Form([]), db: Session = Depends(get_db)):
_require_admin(request)
chart = db.get(ChartDefinition, chart_id)
if not chart: raise HTTPException(404, "Diagramm nicht gefunden")
if chart_type not in CHART_TYPES or group_field not in CHART_FIELDS or (chart_type == "stacked_bar" and stack_field not in CHART_FIELDS):
raise HTTPException(400, "Ungültige Diagrammkonfiguration")
if value_mode not in {"raw", "automatic", "custom"}:
raise HTTPException(400, "Ungültige Wertedarstellung")
mappings = _parse_value_mappings(value_mappings_text)
valid_statuses = {row.name for row in db.query(StatusOption).all()}
chart.name=name.strip(); chart.description=description.strip() or None; chart.chart_type=chart_type; chart.group_field=group_field; chart.stack_field=stack_field if chart_type == "stacked_bar" else None; chart.dashboard_visible=dashboard_visible == "on"; chart.show_empty=show_empty == "on"; chart.max_items=max(1,min(max_items,50)); chart.sort_mode=sort_mode; chart.value_mode=value_mode; chart.value_mappings=mappings; chart.excluded_statuses=[value for value in excluded_statuses if value in valid_statuses]
try:
db.commit()
except IntegrityError:
db.rollback(); raise HTTPException(400, "Ein Diagramm mit diesem Namen existiert bereits.")
return RedirectResponse("/charts/manage?toast_success=Diagramm gespeichert", status_code=303)
@app.post("/charts/{chart_id}/duplicate")
def chart_duplicate(chart_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
source = db.get(ChartDefinition, chart_id)
if not source: raise HTTPException(404, "Diagramm nicht gefunden")
base_name = f"{source.name} (Kopie)"
name = base_name
number = 2
while db.query(ChartDefinition).filter(ChartDefinition.name == name).first():
name = f"{source.name} (Kopie {number})"
number += 1
duplicate = ChartDefinition(name=name, description=source.description, chart_type=source.chart_type, group_field=source.group_field, stack_field=source.stack_field, dashboard_visible=False, show_empty=source.show_empty, max_items=source.max_items, sort_mode=source.sort_mode, value_mode=source.value_mode, value_mappings=list(source.value_mappings or []), excluded_statuses=list(source.excluded_statuses or []), created_by=_changed_by(request), is_system=False)
db.add(duplicate); db.commit(); db.refresh(duplicate)
return RedirectResponse(f"/charts/{duplicate.id}/edit?toast_success=Diagramm dupliziert", status_code=303)
@app.post("/charts/{chart_id}/delete")
def chart_delete(chart_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
chart=db.get(ChartDefinition, chart_id)
if not chart: raise HTTPException(404, "Diagramm nicht gefunden")
if chart.is_system:
return RedirectResponse("/charts/manage?toast_warning=Standarddiagramme können bearbeitet, aber nicht gelöscht werden", status_code=303)
db.delete(chart); db.commit()
return RedirectResponse("/charts/manage?toast_success=Diagramm gelöscht", status_code=303)
@app.get("/assets/new")
def asset_new(request: Request, category_id: int | None = None, db: Session = Depends(get_db)):
_require_admin(request)
categories = db.query(Category).options(joinedload(Category.visible_fields)).all()
selected = next((c for c in categories if c.id == category_id), categories[0] if categories else None)
return templates.TemplateResponse("asset_form.html", {"request": request, "asset": None, "categories": categories, "selected": selected, "fields": ASSET_FIELDS, "status_options": _active_status_options(db, request), "duplicate_mode": False, "parent_candidates": db.query(Asset).order_by(Asset.name).all(), "custom_values": {}})
@app.get("/assets/{asset_id}/duplicate")
def asset_duplicate(asset_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
original = db.query(Asset).options(joinedload(Asset.category).joinedload(Category.visible_fields)).filter(Asset.id == asset_id).first()
if not original:
raise HTTPException(404, "Asset nicht gefunden")
categories = db.query(Category).options(joinedload(Category.visible_fields)).all()
class DuplicateValues:
pass
duplicate = DuplicateValues()
for field_name in ASSET_FIELDS:
setattr(duplicate, field_name, getattr(original, field_name, None))
duplicate.asset_tag = None
duplicate.mesh_node_id = None
duplicate.name = f"{original.name} Kopie"
return templates.TemplateResponse("asset_form.html", {
"request": request, "asset": duplicate, "categories": categories, "selected": original.category,
"fields": ASSET_FIELDS, "status_options": _active_status_options(db, request), "duplicate_mode": True, "parent_candidates": db.query(Asset).filter(Asset.id != original.id).order_by(Asset.name).all(), "custom_values": {},
})
@app.post("/assets/new")
async def asset_create(request: Request, category_id: int = Form(...), image: UploadFile | None = File(None), db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
category = db.query(Category).options(joinedload(Category.visible_fields)).filter(Category.id == category_id).first()
if not category:
raise HTTPException(404, "Kategorie nicht gefunden")
data = {f.field_name: form.get(f.field_name) or None for f in category.visible_fields if f.active and hasattr(Asset, f.field_name) and f.field_name != "parent_asset_id" and not _definition_readonly(f)}
for field in category.visible_fields:
if field.field_name not in data or not field.definition:
continue
if field.definition.data_type == "date":
data[field.field_name] = _normalize_date_value(data[field.field_name])
elif field.definition.data_type == "datetime":
data[field.field_name] = _normalize_datetime_value(data[field.field_name])
data["room_number"] = _normalize_optional_text(form.get("room_number"))
parent_raw = _normalize_optional_text(form.get("parent_asset_id"))
data["parent_asset_id"] = int(parent_raw) if parent_raw else None
if not data.get("name"):
raise HTTPException(400, "Bezeichnung fehlt")
_validate_asset_identifiers(db, asset_tag=data.get("asset_tag"), manufacturer=data.get("manufacturer"), serial_number=data.get("serial_number"))
asset = Asset(category_id=category_id, image_path=save_upload(image), **data)
db.add(asset)
try:
db.flush()
for field in category.visible_fields:
if field.active and field.definition and not field.definition.is_system and not field.definition.readonly:
_custom_value_set(asset, field.definition, form.get(field.field_name), db)
created_changes = {_definition_label(field): {"old": "", "new": _display_value(getattr(asset, field.field_name, None))} for field in category.visible_fields if field.active and hasattr(asset, field.field_name) and getattr(asset, field.field_name, None) not in (None, "")}
_record_asset_history(db, asset, created_changes, source="manual-create", changed_by=_changed_by(request))
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Asset konnte nicht gespeichert werden. Inventarnummer, Seriennummer oder MeshCentral Node-ID ist möglicherweise bereits vergeben.") from exc
return RedirectResponse(f"/assets/{asset.id}", status_code=303)
@app.get("/assets/import")
def assets_import_form(request: Request, category_id: int | None = None, db: Session = Depends(get_db)):
categories = _sort_localized(db.query(Category).options(joinedload(Category.visible_fields)).all(), "name", request)
selected = next((category for category in categories if category.id == category_id), categories[0] if categories else None)
return templates.TemplateResponse(
"assets_import.html",
{"request": request, "categories": categories, "selected": selected, "result": None, "import_only_active_fields": False},
)
@app.get("/assets/import/template.xlsx")
def assets_import_template(request: Request, category_id: int, db: Session = Depends(get_db)):
_require_admin(request)
category = db.query(Category).filter(Category.id == category_id).first()
if not category:
raise HTTPException(404, "Kategorie nicht gefunden")
definitions = (
db.query(FieldDefinition)
.filter(
FieldDefinition.is_active.is_(True),
FieldDefinition.excel_import.is_(True),
FieldDefinition.readonly.is_(False),
)
.order_by(FieldDefinition.sort_order, FieldDefinition.label)
.all()
)
import_definitions = [definition for definition in definitions if definition.field_name != "name"]
headers = ["Asset-ID", "Kategorie", "Bezeichnung"] + [definition.label for definition in import_definitions]
return _xlsx_response(f"asset-import-{category.name}.xlsx", headers, [])
@app.post("/assets/import")
async def assets_import(
request: Request,
category_id: int = Form(...),
file: UploadFile = File(...),
duplicate_action: str = Form("skip"),
import_only_active_fields: str | None = Form(None),
import_mode: str = Form("import"),
db: Session = Depends(get_db),
):
_require_admin(request)
dry_run = import_mode == "dry_run"
category = db.query(Category).options(joinedload(Category.visible_fields)).filter(Category.id == category_id).first()
categories = _sort_localized(db.query(Category).options(joinedload(Category.visible_fields)).all(), "name", request)
if not category:
raise HTTPException(404, "Kategorie nicht gefunden")
if not file.filename or Path(file.filename).suffix.lower() != ".xlsx":
raise HTTPException(400, "Bitte eine XLSX-Datei auswählen.")
try:
raw = await file.read()
workbook = load_workbook(io.BytesIO(raw), read_only=True, data_only=True)
sheet = workbook.active
except Exception as exc:
logger.exception("Excel-Import konnte nicht geöffnet werden")
raise HTTPException(400, "Die Excel-Datei ist ungültig oder beschädigt.") from exc
rows_iter = sheet.iter_rows(values_only=True)
header_row = next(rows_iter, None)
if not header_row:
raise HTTPException(400, "Die Excel-Datei enthält keine Kopfzeile.")
# Excel columns are recognized from the global field catalog. Category
# visibility is only applied when the optional checkbox is enabled.
import_definitions = (
db.query(FieldDefinition)
.filter(
FieldDefinition.is_active.is_(True),
FieldDefinition.excel_import.is_(True),
FieldDefinition.readonly.is_(False),
)
.order_by(FieldDefinition.sort_order, FieldDefinition.label)
.all()
)
definitions_by_name = {definition.field_name: definition for definition in import_definitions}
field_map: dict[str, str] = {
"asset-id": "asset_id",
"asset_id": "asset_id",
"id": "asset_id",
"bezeichnung": "name",
"name": "name",
"kategorie": "category",
"category": "category",
"category_id": "category",
}
for definition in import_definitions:
field_map[definition.label.strip().casefold()] = definition.field_name
field_map[definition.field_name.strip().casefold()] = definition.field_name
categories_by_id = {item.id: item for item in categories}
categories_by_name = {item.name.strip().casefold(): item for item in categories}
columns: list[str | None] = []
unknown_headers: list[str] = []
for header in header_row:
label = _excel_cell_value(header) or ""
field_name = field_map.get(label.casefold())
columns.append(field_name)
if label and not field_name:
unknown_headers.append(label)
if "name" not in columns:
raise HTTPException(400, "Die Spalte Bezeichnung bzw. Name fehlt.")
created = 0
updated = 0
skipped = 0
errors: list[str] = []
row_results: list[dict[str, Any]] = []
changed_by = _changed_by(request)
for excel_row_number, values in enumerate(rows_iter, start=2):
raw_data: dict[str, str | None] = {}
for index, value in enumerate(values):
if index >= len(columns):
break
field_name = columns[index]
if field_name:
raw_data[field_name] = _excel_cell_value(value)
if not any(value not in (None, "") for value in raw_data.values()):
continue
asset_id_value = raw_data.pop("asset_id", None)
category_value = raw_data.pop("category", None)
category_was_supplied = category_value not in (None, "")
row_category = category
if category_was_supplied:
try:
row_category = categories_by_id.get(int(category_value))
except (TypeError, ValueError):
row_category = categories_by_name.get(str(category_value).strip().casefold())
if row_category is None:
message = f"Kategorie {category_value} wurde nicht gefunden"
errors.append(f"Zeile {excel_row_number}: {message}")
row_results.append({"row": excel_row_number, "asset_id": "", "action": "Fehler", "reason": message})
continue
if import_only_active_fields == "on":
allowed_fields = {
field.field_name
for field in row_category.visible_fields
if field.active
and field.definition
and field.definition.is_active
and field.definition.excel_import
and not field.definition.readonly
} | {"name"}
else:
allowed_fields = set(definitions_by_name) | {"name"}
# Leere Zellen werden bewusst nicht in Updates übernommen.
data = {
name: value
for name, value in raw_data.items()
if name in allowed_fields and value not in (None, "")
}
for field_name, value in list(data.items()):
definition = definitions_by_name.get(field_name)
if definition and definition.data_type == "date":
data[field_name] = _normalize_date_value(value)
elif definition and definition.data_type == "datetime":
data[field_name] = _normalize_datetime_value(value)
system_data = {name: value for name, value in data.items() if hasattr(Asset, name)}
custom_data = {
name: value
for name, value in data.items()
if name in definitions_by_name and not definitions_by_name[name].is_system
}
if asset_id_value not in (None, ""):
try:
parsed_asset_id = int(str(asset_id_value).strip())
except (TypeError, ValueError):
message = f"Asset-ID {asset_id_value} ist ungültig"
errors.append(f"Zeile {excel_row_number}: {message}")
row_results.append({"row": excel_row_number, "asset_id": str(asset_id_value or ""), "action": "Fehler", "reason": message})
continue
existing = db.query(Asset).filter(Asset.id == parsed_asset_id).first()
if existing is None:
message = f"Asset-ID {parsed_asset_id} wurde nicht gefunden"
errors.append(f"Zeile {excel_row_number}: {message}")
row_results.append({"row": excel_row_number, "asset_id": parsed_asset_id, "action": "Fehler", "reason": message})
continue
changes: dict[str, dict[str, str]] = {}
if category_was_supplied and existing.category_id != row_category.id:
changes["Kategorie"] = {
"old": existing.category.name if existing.category else str(existing.category_id),
"new": row_category.name,
}
existing.category_id = row_category.id
for field_name, value in system_data.items():
old_value = getattr(existing, field_name, None)
if _display_value(old_value) != _display_value(value):
definition = definitions_by_name.get(field_name)
label = definition.label if definition else ASSET_FIELDS.get(field_name, field_name)
changes[label] = {
"old": _display_value(old_value),
"new": _display_value(value),
}
setattr(existing, field_name, value)
for field_name, value in custom_data.items():
definition = definitions_by_name[field_name]
old_value = _custom_value_get(existing, definition, db)
if _display_value(old_value) != _display_value(value):
changes[definition.label] = {
"old": _display_value(old_value),
"new": _display_value(value),
}
_custom_value_set(existing, definition, value, db)
if changes:
_record_asset_history(
db,
existing,
changes,
source="excel-import-update",
changed_by=changed_by,
)
updated += 1
changed_labels = ", ".join(changes.keys())
row_results.append({
"row": excel_row_number,
"asset_id": existing.id,
"action": "Aktualisiert",
"reason": f"Geänderte Felder: {changed_labels}",
})
else:
skipped += 1
row_results.append({
"row": excel_row_number,
"asset_id": existing.id,
"action": "Übersprungen",
"reason": "Keine Änderungen; alle ausgefüllten Werte entsprechen bereits dem Datensatz.",
})
continue
# Ohne Asset-ID wird immer ein neuer Datensatz erzeugt.
if not system_data.get("name"):
message = "Für ein neues Asset fehlt die Bezeichnung"
errors.append(f"Zeile {excel_row_number}: {message}")
row_results.append({"row": excel_row_number, "asset_id": "", "action": "Fehler", "reason": message})
continue
try:
with db.begin_nested():
asset = Asset(category_id=row_category.id, **system_data)
db.add(asset)
db.flush()
for field_name, value in custom_data.items():
_custom_value_set(asset, definitions_by_name[field_name], value, db)
created_changes = {
"Kategorie": {"old": "", "new": row_category.name},
**{
(definitions_by_name[name].label if name in definitions_by_name else ASSET_FIELDS.get(name, name)): {"old": "", "new": _display_value(value)}
for name, value in data.items()
},
}
_record_asset_history(
db,
asset,
created_changes,
source="excel-import-create",
changed_by=changed_by,
)
created += 1
row_results.append({
"row": excel_row_number,
"asset_id": asset.id,
"action": "Neu angelegt",
"reason": f"Neues Asset in Kategorie {row_category.name}",
})
except IntegrityError:
message = (
"Das neue Asset verletzt eine Eindeutigkeitsregel "
"(zum Beispiel Inventarnummer, Seriennummer oder MeshCentral Node-ID)"
)
errors.append(f"Zeile {excel_row_number}: {message}")
row_results.append({"row": excel_row_number, "asset_id": "", "action": "Fehler", "reason": message})
continue
if dry_run:
# Flushes and nested savepoints above provide the same validation as a
# real import. The outer transaction is rolled back so no asset,
# custom field value or history entry remains stored.
db.rollback()
else:
try:
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Der Import konnte wegen doppelter eindeutiger Werte nicht gespeichert werden.") from exc
result = {
"created": created,
"updated": updated,
"skipped": skipped,
"errors": errors,
"unknown_headers": unknown_headers,
"row_results": row_results,
"dry_run": dry_run,
}
return templates.TemplateResponse(
"assets_import.html",
{"request": request, "categories": categories, "selected": category, "result": result, "import_only_active_fields": import_only_active_fields == "on", "import_mode": import_mode},
)
@app.get("/asset-tree")
def asset_tree(
request: Request,
root_id: str | None = None,
name_filter: str = "",
hide_without_parent: str | None = None,
db: Session = Depends(get_db),
):
"""Zeigt die gefilterte Gerätehierarchie.
Mehrere Kategorien werden als wiederholte Query-Parameter ``category_id``
übergeben. Filter werden auf jedes Gerät angewendet, nicht nur auf die
Wurzelknoten. Wenn der direkte Elternknoten ausgefiltert ist, wird das
passende Gerät in der Darstellung als eigener Wurzelknoten angezeigt.
"""
try:
parsed_root_id = int(root_id) if root_id not in (None, "") else None
except (TypeError, ValueError):
parsed_root_id = None
selected_category_ids: list[int] = []
for value in request.query_params.getlist("category_id"):
try:
category_id = int(value)
except (TypeError, ValueError):
continue
if category_id not in selected_category_ids:
selected_category_ids.append(category_id)
# Zuerst alle für den angemeldeten Benutzer sichtbaren Assets laden. Die
# Filterung erfolgt anschließend in Python, damit Treffer in beliebiger
# Hierarchietiefe nicht durch eine SQL-Abfrage von ihren Eltern getrennt
# oder versehentlich nur auf Wurzelknoten beschränkt werden.
all_assets = (
_apply_asset_access(db.query(Asset), request)
.options(joinedload(Asset.category))
.order_by(Asset.name)
.all()
)
all_by_id = {asset.id: asset for asset in all_assets}
all_children: dict[int | None, list[Asset]] = {}
for asset in all_assets:
parent_key = asset.parent_asset_id if asset.parent_asset_id in all_by_id else None
all_children.setdefault(parent_key, []).append(asset)
# Bei ausgewähltem obersten Knoten nur dessen Teilbaum berücksichtigen.
scope_ids: set[int]
selected_root = all_by_id.get(parsed_root_id) if parsed_root_id else None
if selected_root:
scope_ids = set()
pending = [selected_root.id]
while pending:
current_id = pending.pop()
if current_id in scope_ids:
continue
scope_ids.add(current_id)
pending.extend(child.id for child in all_children.get(current_id, []))
else:
scope_ids = set(all_by_id)
needle = name_filter.strip().casefold()
selected_categories = set(selected_category_ids)
visible_assets: list[Asset] = []
for asset in all_assets:
if asset.id not in scope_ids:
continue
category_matches = not selected_categories or asset.category_id in selected_categories
name_matches = not needle or needle in (asset.name or "").casefold()
if category_matches and name_matches:
visible_assets.append(asset)
# Der explizit gewählte oberste Knoten bleibt als Orientierung sichtbar,
# auch wenn er selbst den Namens- oder Kategorienfilter nicht erfüllt.
if selected_root and selected_root not in visible_assets:
visible_assets.insert(0, selected_root)
visible_ids = {asset.id for asset in visible_assets}
by_parent: dict[int | None, list[Asset]] = {}
for asset in visible_assets:
parent_key = asset.parent_asset_id if asset.parent_asset_id in visible_ids else None
by_parent.setdefault(parent_key, []).append(asset)
for children in by_parent.values():
children.sort(key=lambda item: (item.name or "").casefold())
hide_unassigned = str(hide_without_parent or "").lower() in {"1", "true", "on", "yes"}
if selected_root:
roots = [selected_root]
# Treffer, deren Zwischenknoten weggefiltert wurden, zusätzlich als
# eigene Wurzeln darstellen, damit kein Treffer verloren geht.
roots.extend(
asset for asset in by_parent.get(None, [])
if asset.id != selected_root.id
)
else:
roots = list(by_parent.get(None, []))
if hide_unassigned:
# Nur echte Einzelgeräte ohne hinterlegtes übergeordnetes Gerät
# ausblenden. Geräte, die lediglich wegen eines Filters als Wurzel
# erscheinen, bleiben sichtbar. Ein gewählter oberster Knoten bleibt.
roots = [
asset for asset in roots
if asset.id == parsed_root_id
or asset.parent_asset_id is not None
or bool(by_parent.get(asset.id))
]
categories = db.query(Category).order_by(Category.name).all()
actual_roots = [
asset for asset in all_assets
if asset.parent_asset_id is None or asset.parent_asset_id not in all_by_id
]
actual_roots.sort(key=lambda item: (item.name or "").casefold())
return templates.TemplateResponse(
"asset_tree.html",
{
"request": request,
"roots": roots,
"by_parent": by_parent,
"categories": categories,
"all_roots": actual_roots,
"selected_root_id": parsed_root_id,
"selected_category_ids": selected_category_ids,
"name_filter": name_filter,
"hide_without_parent": hide_unassigned,
},
)
def _inventory_fallback(asset: Asset) -> tuple[dict, list | dict, dict]:
raw = asset.mesh_source_data if isinstance(asset.mesh_source_data, dict) else {}
hardware = asset.mesh_hardware_data if isinstance(getattr(asset, "mesh_hardware_data", None), dict) else {}
software = getattr(asset, "mesh_software_data", None)
# Rebuild from the raw snapshot whenever possible. This upgrades older snapshots
# with language-neutral i18n keys immediately, without waiting for another sync.
if raw:
try:
from app.meshcentral import _build_hardware_inventory
hardware = _build_hardware_inventory(raw)
except Exception:
if not hardware:
hardware = (((raw.get("sys") or {}).get("hardware")) if isinstance(raw.get("sys"), dict) else {}) or {}
if not isinstance(software, (list, dict)) or not software:
candidates = []
sys_data = raw.get("sys") if isinstance(raw.get("sys"), dict) else {}
windows = (((sys_data.get("hardware") or {}).get("windows")) if isinstance(sys_data.get("hardware"), dict) else {}) or {}
for candidate in (sys_data.get("software"), windows.get("software") if isinstance(windows, dict) else None, raw.get("software"), raw.get("installedSoftware"), raw.get("installedApplications")):
if isinstance(candidate, (list, dict)) and candidate:
software = candidate
break
else:
software = []
return hardware, software, raw
def _software_table(software: list | dict) -> tuple[list[str], list[dict]]:
if isinstance(software, dict):
if all(isinstance(value, dict) for value in software.values()):
rows = [{"Name": key, **value} for key, value in software.items()]
else:
rows = [{"Eigenschaft": key, "Wert": value} for key, value in software.items()]
elif isinstance(software, list):
rows = [item if isinstance(item, dict) else {"Wert": item} for item in software]
else:
rows = []
preferred = ["Name", "DisplayName", "name", "Version", "DisplayVersion", "version", "Publisher", "publisher", "InstallDate", "installDate"]
keys = []
for key in preferred:
if any(key in row for row in rows) and key not in keys:
keys.append(key)
for row in rows:
for key in row.keys():
if key not in keys:
keys.append(key)
return keys, rows
@app.get("/assets/{asset_id:int}")
def asset_detail(asset_id: int, request: Request, db: Session = Depends(get_db)):
asset = _apply_asset_access(db.query(Asset), request).options(joinedload(Asset.category).joinedload(Category.visible_fields)).filter(Asset.id == asset_id).first()
if not asset:
raise HTTPException(404, "Asset nicht gefunden")
active_fields = [f for f in asset.category.visible_fields if f.active]
history_count = db.query(AssetHistory).filter(AssetHistory.asset_id == asset.id).count()
hardware_inventory, software_inventory, raw_inventory = _inventory_fallback(asset)
# The software tab on the asset detail page uses the dedicated callback
# inventory. MeshCentral raw details remain available independently.
stored_software = (
db.query(InstalledSoftware)
.filter(InstalledSoftware.asset_id == asset.id)
.order_by(InstalledSoftware.name.asc(), InstalledSoftware.version.asc())
.all()
)
latest_software_job = (
db.query(SoftwareJob)
.options(joinedload(SoftwareJob.package))
.filter(SoftwareJob.asset_id == asset.id)
.order_by(SoftwareJob.created_at.desc())
.first()
)
return templates.TemplateResponse("asset_detail.html", {
"request": request,
"asset": asset,
"active_fields": active_fields,
"history_count": history_count,
"is_admin": _is_admin(request),
"issue_status_name": (_action_status(db, "issue").name if _action_status(db, "issue") else None),
"today": date.today().isoformat(),
"parent_asset": asset.parent_asset,
"child_assets": sorted(asset.child_assets, key=lambda item: item.name.casefold()),
"custom_values": {
f.field_name: _custom_value_get(asset, f.definition, db)
for f in active_fields
if f.definition and not f.definition.is_system
},
"hardware_inventory": hardware_inventory,
"software_inventory": software_inventory,
"installed_software": stored_software,
"latest_software_job": latest_software_job,
"raw_inventory": raw_inventory,
"asset_platform": detect_platform(asset),
"job_definitions": (
db.query(JobDefinition).options(selectinload(JobDefinition.parameters))
.filter(JobDefinition.enabled.is_(True), JobDefinition.platform.in_([detect_platform(asset), "all"]))
.order_by(JobDefinition.name).all()
if _is_admin(request) and asset.mesh_node_id and detect_platform(asset) != "unknown" else []
),
"job_parameter_overrides": ({x.definition_parameter_id:x.value_text for x in db.query(JobDefinitionAssetParameter).filter(JobDefinitionAssetParameter.asset_id==asset.id).all()} if _is_admin(request) else {}),
})
def _reset_mesh_link_runtime_state(asset: Asset) -> None:
"""Clear only connection-derived runtime state after a manual relink.
Inventory and normal asset fields intentionally remain untouched until the
next regular MeshCentral synchronization applies the configured mappings.
"""
asset.mesh_online = None
asset.mesh_online_state = "unknown"
asset.mesh_presence_updated_at = None
asset.mesh_last_online_at = None
asset.mesh_last_seen = None
asset.mesh_presence_source = None
asset.mesh_last_sync = None
def _mesh_link_asset(db: Session, asset_id: int) -> Asset:
asset = db.query(Asset).filter(Asset.id == asset_id).first()
if not asset:
raise HTTPException(404, "Asset nicht gefunden")
return asset
@app.get("/assets/{asset_id}/mesh-link")
def asset_mesh_link_page(asset_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
asset = _mesh_link_asset(db, asset_id)
try:
devices = fetch_device_summaries_for_linking()
except RuntimeError as exc:
return RedirectResponse(
f"/assets/{asset.id}?toast_error=" + quote(
_translate_request(
request,
"mesh.link.fetch_failed",
"The MeshCentral device list could not be loaded: {error}",
error=str(exc),
)
),
status_code=303,
)
node_ids = [item["node_id"] for item in devices if item.get("node_id")]
assigned_rows = (
db.query(Asset.id, Asset.name, Asset.mesh_node_id)
.filter(Asset.mesh_node_id.in_(node_ids))
.all()
if node_ids
else []
)
assigned = {str(row.mesh_node_id): {"asset_id": row.id, "asset_name": row.name} for row in assigned_rows if row.mesh_node_id}
available_devices: list[dict[str, Any]] = []
current_device: dict[str, Any] | None = None
for item in devices:
row = dict(item)
owner = assigned.get(row["node_id"])
row["assigned_asset_id"] = owner["asset_id"] if owner else None
row["assigned_asset_name"] = owner["asset_name"] if owner else None
row["is_current"] = row["node_id"] == (asset.mesh_node_id or "")
if row["is_current"]:
current_device = row
continue
if owner is None:
available_devices.append(row)
return templates.TemplateResponse(
"asset_mesh_link.html",
{
"request": request,
"asset": asset,
"current_device": current_device,
"available_devices": available_devices,
},
)
@app.post("/assets/{asset_id}/mesh-link")
def asset_mesh_link_save(
asset_id: int,
request: Request,
node_id: str = Form(...),
db: Session = Depends(get_db),
):
_require_admin(request)
asset = _mesh_link_asset(db, asset_id)
requested_node_id = str(node_id or "").strip()
if not requested_node_id:
return RedirectResponse(
f"/assets/{asset.id}/mesh-link?toast_error="
+ quote(_translate_request(request, "mesh.link.select_required", "Select a MeshCentral device.")),
status_code=303,
)
try:
device_ids = {item["node_id"] for item in fetch_device_summaries_for_linking() if item.get("node_id")}
except RuntimeError as exc:
return RedirectResponse(
f"/assets/{asset.id}?toast_error=" + quote(
_translate_request(
request,
"mesh.link.fetch_failed",
"The MeshCentral device list could not be loaded: {error}",
error=str(exc),
)
),
status_code=303,
)
if requested_node_id not in device_ids:
return RedirectResponse(
f"/assets/{asset.id}/mesh-link?toast_error="
+ quote(_translate_request(request, "mesh.link.not_found", "The selected MeshCentral device no longer exists.")),
status_code=303,
)
owner = db.query(Asset).filter(Asset.mesh_node_id == requested_node_id, Asset.id != asset.id).first()
if owner:
message = _translate_request(
request,
"mesh.link.already_assigned",
"This MeshCentral device is already linked to asset {asset}.",
asset=owner.name,
)
return RedirectResponse(
f"/assets/{asset.id}/mesh-link?toast_error=" + quote(message),
status_code=303,
)
old_node_id = asset.mesh_node_id or ""
if old_node_id == requested_node_id:
return RedirectResponse(
f"/assets/{asset.id}?toast_warning="
+ quote(_translate_request(request, "mesh.link.unchanged", "The MeshCentral link is already set to this device.")),
status_code=303,
)
asset.mesh_node_id = requested_node_id
_reset_mesh_link_runtime_state(asset)
asset.mesh_sync_status = "pending"
asset.mesh_sync_message = _translate_request(
request,
"mesh.link.pending_sync",
"MeshCentral link changed manually; synchronization is pending.",
)
_record_asset_history(
db,
asset,
{
_translate_request(request, "field.mesh_node_id", "MeshCentral node ID"): {
"old": old_node_id,
"new": requested_node_id,
}
},
source="mesh-link-manual",
changed_by=_changed_by(request),
)
try:
db.commit()
except IntegrityError:
db.rollback()
return RedirectResponse(
f"/assets/{asset.id}/mesh-link?toast_error="
+ quote(_translate_request(request, "mesh.link.conflict", "The MeshCentral node ID was assigned by another operation.")),
status_code=303,
)
return RedirectResponse(
f"/assets/{asset.id}?toast_success="
+ quote(_translate_request(request, "mesh.link.changed", "MeshCentral link changed.")),
status_code=303,
)
@app.post("/assets/{asset_id}/mesh-unlink")
def asset_mesh_unlink(asset_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
asset = _mesh_link_asset(db, asset_id)
old_node_id = asset.mesh_node_id or ""
if not old_node_id:
return RedirectResponse(
f"/assets/{asset.id}?toast_warning="
+ quote(_translate_request(request, "mesh.link.already_unlinked", "This asset is not linked to MeshCentral.")),
status_code=303,
)
asset.mesh_node_id = None
_reset_mesh_link_runtime_state(asset)
asset.mesh_sync_status = "unlinked"
asset.mesh_sync_message = _translate_request(request, "mesh.link.unlinked_message", "MeshCentral link removed manually.")
_record_asset_history(
db,
asset,
{
_translate_request(request, "field.mesh_node_id", "MeshCentral node ID"): {
"old": old_node_id,
"new": "",
}
},
source="mesh-unlink-manual",
changed_by=_changed_by(request),
)
db.commit()
return RedirectResponse(
f"/assets/{asset.id}?toast_success="
+ quote(_translate_request(request, "mesh.link.unlinked", "MeshCentral link removed.")),
status_code=303,
)
def _presence_payload(asset: Asset) -> dict[str, Any]:
return {
"asset_id": asset.id,
"has_mesh_node_id": bool(asset.mesh_node_id),
"state": asset.mesh_online_state or "unknown",
"online": asset.mesh_online,
"updated_at": _utc_iso(asset.mesh_presence_updated_at) or None,
"last_online_at": _utc_iso(asset.mesh_last_online_at) or None,
"last_seen_at": _utc_iso(asset.mesh_presence_updated_at if (asset.mesh_online_state or "unknown") == "online" else asset.mesh_last_seen) or None,
"historical_last_seen_at": _utc_iso(asset.mesh_last_seen) or None,
"source": asset.mesh_presence_source,
}
@app.get("/api/assets/presence")
def assets_presence(request: Request, db: Session = Depends(get_db)):
assets = (
_apply_asset_access(db.query(Asset), request)
.filter(Asset.mesh_node_id.isnot(None), Asset.mesh_node_id != "")
.all()
)
return {"items": [_presence_payload(asset) for asset in assets]}
@app.get("/api/assets/{asset_id}/presence")
def asset_presence(asset_id: int, request: Request, db: Session = Depends(get_db)):
asset = _get_visible_asset(db, request, asset_id)
return _presence_payload(asset)
@app.get("/assets/{asset_id}/mesh-inventory.json")
def asset_mesh_inventory_json(asset_id: int, request: Request, db: Session = Depends(get_db)):
asset = _apply_asset_access(db.query(Asset), request).filter(Asset.id == asset_id).first()
if not asset:
raise HTTPException(404, "Asset nicht gefunden")
hardware, software, raw = _inventory_fallback(asset)
payload = {
"asset_id": asset.id,
"asset_name": asset.name,
"mesh_node_id": asset.mesh_node_id,
"updated_at": _utc_iso(asset.mesh_inventory_updated_at) or None,
"hardware": hardware,
"software": software,
"raw": raw,
}
filename = re.sub(r"[^A-Za-z0-9._-]+", "_", asset.name or f"asset-{asset.id}") + "-mesh-inventory.json"
return Response(content=json.dumps(payload, ensure_ascii=False, indent=2, default=str), media_type="application/json", headers={"Content-Disposition": f'attachment; filename="{filename}"'})
@app.get("/assets/{asset_id}/edit")
def asset_edit(asset_id: int, request: Request, category_id: int | None = None, db: Session = Depends(get_db)):
_require_admin(request)
asset = db.query(Asset).options(joinedload(Asset.category).joinedload(Category.visible_fields)).filter(Asset.id == asset_id).first()
categories = db.query(Category).options(joinedload(Category.visible_fields)).all()
if not asset:
raise HTTPException(404, "Asset nicht gefunden")
selected = next((category for category in categories if category.id == category_id), None) if category_id else asset.category
if selected is None:
selected = asset.category
return templates.TemplateResponse("asset_form.html", {"request": request, "asset": asset, "categories": categories, "selected": selected, "fields": ASSET_FIELDS, "status_options": _active_status_options(db, request), "duplicate_mode": False, "parent_candidates": db.query(Asset).filter(Asset.id != asset.id).order_by(Asset.name).all(), "custom_values": {f.field_name: _custom_value_get(asset, f.definition, db) for f in selected.visible_fields if f.definition and not f.definition.is_system}})
@app.post("/assets/{asset_id}/edit")
async def asset_update(asset_id: int, request: Request, category_id: int = Form(...), image: UploadFile | None = File(None), db: Session = Depends(get_db)):
_require_admin(request)
asset = db.query(Asset).filter(Asset.id == asset_id).first()
category = db.query(Category).options(joinedload(Category.visible_fields)).filter(Category.id == category_id).first()
if not asset or not category:
raise HTTPException(404, "Asset oder Kategorie nicht gefunden")
form = await request.form()
changes: dict[str, dict[str, str]] = {}
# mesh_node_id is a protected system field. On the administrator-only
# edit page it may nevertheless be changed directly or cleared. Do this
# separately from the generic field loop because readonly system fields
# are intentionally skipped there.
if "mesh_node_id" in form:
old_node_id = _normalize_optional_text(asset.mesh_node_id)
requested_node_id = _normalize_optional_text(form.get("mesh_node_id"))
if old_node_id != requested_node_id:
if requested_node_id:
owner = (
db.query(Asset.id, Asset.name)
.filter(Asset.mesh_node_id == requested_node_id, Asset.id != asset.id)
.first()
)
if owner:
message = _translate_request(
request,
"mesh.link.already_assigned",
"This MeshCentral device is already linked to asset {asset}.",
).replace("{asset}", owner.name or str(owner.id))
return RedirectResponse(
f"/assets/{asset.id}/edit?toast_error=" + quote(message),
status_code=303,
)
changes[_translate_request(request, "field.mesh_node_id", "MeshCentral node ID")] = {
"old": old_node_id or "",
"new": requested_node_id or "",
}
asset.mesh_node_id = requested_node_id
_reset_mesh_link_runtime_state(asset)
if requested_node_id:
asset.mesh_sync_status = "pending"
asset.mesh_sync_message = _translate_request(
request,
"mesh.link.pending_sync",
"MeshCentral link changed manually; synchronization is pending.",
)
else:
asset.mesh_sync_status = "unlinked"
asset.mesh_sync_message = _translate_request(
request,
"mesh.link.unlinked_message",
"MeshCentral link removed manually.",
)
old_category_id = asset.category_id
if old_category_id != category_id:
old_category = db.query(Category).filter(Category.id == old_category_id).first()
changes["Kategorie"] = {"old": old_category.name if old_category else str(old_category_id), "new": category.name}
asset.category_id = category_id
for field in category.visible_fields:
if (
field.active
and hasattr(asset, field.field_name)
and field.field_name not in {"room_number", "parent_asset_id"}
and not _definition_readonly(field)
):
old_value = getattr(asset, field.field_name)
new_value = form.get(field.field_name) or None
if field.definition and field.definition.data_type == "date":
new_value = _normalize_date_value(new_value)
elif field.definition and field.definition.data_type == "datetime":
new_value = _normalize_datetime_value(new_value)
if _display_value(old_value) != _display_value(new_value):
changes[_definition_label(field)] = {"old": _display_value(old_value), "new": _display_value(new_value)}
setattr(asset, field.field_name, new_value)
new_room = _normalize_optional_text(form.get("room_number"))
if _display_value(asset.room_number) != _display_value(new_room):
changes["Raumnummer"] = {"old": _display_value(asset.room_number), "new": _display_value(new_room)}
asset.room_number = new_room
parent_raw = _normalize_optional_text(form.get("parent_asset_id"))
new_parent_id = int(parent_raw) if parent_raw else None
if new_parent_id and _would_create_parent_cycle(db, asset.id, new_parent_id):
raise HTTPException(400, "Diese Zuordnung würde einen Kreis in der Baustruktur erzeugen.")
current_parent_id = int(asset.parent_asset_id) if asset.parent_asset_id not in (None, "") else None
if current_parent_id != new_parent_id:
old_parent = db.query(Asset).filter(Asset.id == current_parent_id).first() if current_parent_id else None
new_parent = db.query(Asset).filter(Asset.id == new_parent_id).first() if new_parent_id else None
changes["Übergeordnetes Gerät"] = {"old": old_parent.name if old_parent else "", "new": new_parent.name if new_parent else ""}
asset.parent_asset_id = new_parent_id
for field in category.visible_fields:
if field.active and field.definition and not field.definition.is_system and not field.definition.readonly:
old_custom = _custom_value_get(asset, field.definition, db)
new_custom = form.get(field.field_name)
if _display_value(old_custom) != _display_value(new_custom):
changes[_definition_label(field)] = {"old": _display_value(old_custom), "new": _display_value(new_custom)}
_custom_value_set(asset, field.definition, new_custom, db)
_validate_asset_identifiers(db, asset_tag=asset.asset_tag, manufacturer=asset.manufacturer, serial_number=asset.serial_number, exclude_id=asset.id)
new_image = save_upload(image)
if new_image:
if asset.image_path != new_image:
changes["Bild"] = {"old": asset.image_path or "", "new": new_image}
asset.image_path = new_image
try:
_record_asset_history(db, asset, changes, source="manual-update", changed_by=_changed_by(request))
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Asset konnte nicht gespeichert werden. Inventarnummer, Seriennummer oder MeshCentral Node-ID ist möglicherweise bereits vergeben.") from exc
return RedirectResponse(f"/assets/{asset.id}", status_code=303)
@app.get("/assets/{asset_id}/history")
def asset_history(asset_id: int, request: Request, db: Session = Depends(get_db)):
asset = _get_visible_asset(db, request, asset_id)
entries = db.query(AssetHistory).filter(AssetHistory.asset_id == asset_id).order_by(AssetHistory.changed_at.desc()).all()
return templates.TemplateResponse("asset_history.html", {"request": request, "asset": asset, "entries": entries})
@app.get("/assets/{asset_id}/history/export.xlsx")
def asset_history_export(asset_id: int, request: Request, db: Session = Depends(get_db)):
asset = _get_visible_asset(db, request, asset_id)
if not asset:
raise HTTPException(404, "Asset nicht gefunden")
entries = db.query(AssetHistory).filter(AssetHistory.asset_id == asset_id).order_by(AssetHistory.changed_at.desc()).all()
rows = []
for entry in entries:
for label, change in (entry.changes or {}).items():
rows.append([entry.changed_at, entry.source, entry.changed_by or "", label, change.get("old", ""), change.get("new", "")])
return _xlsx_response(f"asset-{asset.id}-historie.xlsx", ["Zeitpunkt", "Quelle", "Benutzer", "Feld", "Alter Wert", "Neuer Wert"], rows)
@app.get("/assets/export.xlsx")
def assets_export(request: Request, category_id: int | None = None, db: Session = Depends(get_db)):
categories = db.query(Category).options(joinedload(Category.visible_fields)).all()
selected = next((category for category in categories if category.id == category_id), None) if category_id else None
query = _apply_asset_access(db.query(Asset), request)
if selected:
query = query.filter(Asset.category_id == selected.id)
assets = query.order_by(Asset.name).all()
fields = ([field for field in sorted(selected.visible_fields, key=lambda item: item.sort_order) if field.active and field.show_in_list] if selected else _list_fields_for_categories(categories))
headers = ["Asset-ID"] + ((["Kategorie"] if not selected else []) + [field.label for field in fields])
if not any(field.field_name == "name" for field in fields):
headers.insert(1, "Bezeichnung")
rows = []
for asset in assets:
row = [asset.id] + (([asset.category.name] if not selected else []) + [getattr(asset, field.field_name, "") for field in fields])
if not any(field.field_name == "name" for field in fields):
row.insert(1, asset.name)
rows.append(row)
suffix = selected.name if selected else "alle"
return _xlsx_response(f"assets-{suffix}.xlsx", headers, rows)
@app.get("/export/assets.xlsx")
def assets_export_static(request: Request, category_id: int | None = None, db: Session = Depends(get_db)):
return assets_export(request=request, category_id=category_id, db=db)
@app.post("/assets/{asset_id}/issue")
def asset_issue(asset_id: int, request: Request, return_to: str = Form("detail"), assigned_to: str = Form(""), assigned_on: str = Form(""), db: Session = Depends(get_db)):
_require_admin(request)
asset = _get_visible_asset(db, request, asset_id)
status = _action_status(db, "issue")
if not status:
raise HTTPException(400, "Kein Status für die Aktion Ausgeben konfiguriert")
changes = _asset_action_changes(asset, "issue", status.name, assigned_to, assigned_on)
_record_asset_history(db, asset, changes, source="manual-issue", changed_by=_changed_by(request))
db.commit()
return RedirectResponse("/" if return_to == "dashboard" else ("/assets" if return_to == "list" else f"/assets/{asset.id}"), status_code=303)
@app.post("/assets/{asset_id}/return")
def asset_return(asset_id: int, request: Request, return_to: str = Form("detail"), db: Session = Depends(get_db)):
_require_admin(request)
asset = _get_visible_asset(db, request, asset_id)
status = _action_status(db, "return")
if not status:
raise HTTPException(400, "Kein Status für die Aktion Zurücknehmen konfiguriert")
changes = _asset_action_changes(asset, "return", status.name)
_record_asset_history(db, asset, changes, source="manual-return", changed_by=_changed_by(request))
db.commit()
return RedirectResponse("/" if return_to == "dashboard" else ("/assets" if return_to == "list" else f"/assets/{asset.id}"), status_code=303)
MERGE_PREVIEW_FIELDS = [
("id", "Asset-ID"),
("category", "Kategorie"),
("name", "Bezeichnung"),
("model", "Modell"),
("serial_number", "Seriennummer"),
("hostname", "Computername"),
("notes", "Notizen"),
("mesh_node_id", "MeshCentral Node-ID"),
]
MERGE_FILL_FIELDS = ["model", "serial_number", "hostname", "notes", "mesh_node_id"]
def _merge_display_value(asset: Asset, field_name: str) -> str:
if field_name == "id":
return str(asset.id)
if field_name == "category":
return asset.category.name if asset.category else ""
return _display_value(getattr(asset, field_name, None))
def _merge_preview(assets: list[Asset]) -> list[dict]:
result = []
for field_name, label in MERGE_PREVIEW_FIELDS:
source = assets[0]
value = _merge_display_value(source, field_name)
if field_name not in {"id", "category", "name"} and not value:
for candidate in assets[1:]:
candidate_value = _merge_display_value(candidate, field_name)
if candidate_value:
source = candidate
value = candidate_value
break
result.append({
"field_name": field_name,
"label": label,
"value": value,
"source_asset_id": source.id,
})
return result
def _ordered_merge_assets(db: Session, asset_ids: list[str]) -> list[Asset]:
ids = _selected_asset_ids(asset_ids)
if len(ids) < 2:
raise HTTPException(400, "Bitte mindestens zwei Assets auswählen.")
rows = db.query(Asset).options(joinedload(Asset.category)).filter(Asset.id.in_(ids)).all()
by_id = {asset.id: asset for asset in rows}
if len(by_id) != len(ids):
raise HTTPException(404, "Mindestens ein ausgewähltes Asset wurde nicht gefunden.")
return [by_id[asset_id] for asset_id in ids]
@app.post("/assets/merge")
def assets_merge_order(request: Request, asset_ids: list[str] = Form([]), db: Session = Depends(get_db)):
_require_admin(request)
assets = _ordered_merge_assets(db, asset_ids)
return templates.TemplateResponse(
"assets_merge_order.html",
{"request": request, "assets": assets},
)
@app.post("/assets/merge/preview")
def assets_merge_preview(request: Request, asset_ids: list[str] = Form([]), db: Session = Depends(get_db)):
_require_admin(request)
assets = _ordered_merge_assets(db, asset_ids)
return templates.TemplateResponse(
"assets_merge_preview.html",
{"request": request, "assets": assets, "preview": _merge_preview(assets)},
)
@app.post("/assets/merge/apply")
def assets_merge_apply(request: Request, asset_ids: list[str] = Form([]), db: Session = Depends(get_db)):
_require_admin(request)
assets = _ordered_merge_assets(db, asset_ids)
target = assets[0]
sources = assets[1:]
source_ids = [asset.id for asset in sources]
changes = {}
for field_name in MERGE_FILL_FIELDS:
if _display_value(getattr(target, field_name, None)):
continue
for source in sources:
candidate = getattr(source, field_name, None)
if not _display_value(candidate):
continue
if field_name == "mesh_node_id":
# mesh_node_id is unique. When a manually imported asset stays
# and absorbs a MeshCentral asset, the source must release the
# node ID before the target can receive it. Without this flush,
# PostgreSQL may update the target first and reject the merge.
source.mesh_node_id = None
db.flush()
setattr(target, field_name, candidate)
changes[field_name] = {
"old": "",
"new": _display_value(candidate),
"source_asset_id": str(source.id),
}
break
target_values = {
row.field_definition_id: row
for row in db.query(AssetFieldValue).filter(AssetFieldValue.asset_id == target.id).all()
}
for source in sources:
for source_value in db.query(AssetFieldValue).filter(AssetFieldValue.asset_id == source.id).all():
existing = target_values.get(source_value.field_definition_id)
current = ""
if existing:
current = _display_value(
existing.value_text if existing.value_text is not None else
existing.value_number if existing.value_number is not None else
existing.value_boolean if existing.value_boolean is not None else
existing.value_date if existing.value_date is not None else
existing.value_datetime if existing.value_datetime is not None else
existing.value_json
)
if current:
continue
if existing is None:
source_value.asset_id = target.id
target_values[source_value.field_definition_id] = source_value
else:
existing.value_text = source_value.value_text
existing.value_number = source_value.value_number
existing.value_boolean = source_value.value_boolean
existing.value_date = source_value.value_date
existing.value_datetime = source_value.value_datetime
existing.value_json = source_value.value_json
db.delete(source_value)
db.query(AssetHistory).filter(AssetHistory.asset_id.in_(source_ids)).update(
{AssetHistory.asset_id: target.id}, synchronize_session=False
)
db.query(SoftwareJob).filter(SoftwareJob.asset_id.in_(source_ids)).update(
{SoftwareJob.asset_id: target.id}, synchronize_session=False
)
# Compact job states are derived data. Rebuild them for the surviving asset
# after the complete job history has been moved.
db.query(AssetJobState).filter(AssetJobState.asset_id.in_([target.id, *source_ids])).delete(synchronize_session=False)
db.flush()
for merged_job in (
db.query(SoftwareJob)
.filter(SoftwareJob.asset_id == target.id)
.order_by(SoftwareJob.created_at.asc(), SoftwareJob.id.asc())
.all()
):
sync_asset_job_state(
db,
merged_job,
increment_execution=True,
status_at=merged_job.finished_at or merged_job.callback_completed_at or merged_job.started_at or merged_job.sent_at or merged_job.created_at,
)
existing_software = {
(row.name.casefold(), (row.version or "").casefold(), (row.publisher or "").casefold(), (row.architecture or "").casefold())
for row in db.query(InstalledSoftware).filter(InstalledSoftware.asset_id == target.id).all()
}
for row in db.query(InstalledSoftware).filter(InstalledSoftware.asset_id.in_(source_ids)).all():
identity = (row.name.casefold(), (row.version or "").casefold(), (row.publisher or "").casefold(), (row.architecture or "").casefold())
if identity in existing_software:
db.delete(row)
else:
row.asset_id = target.id
existing_software.add(identity)
db.query(Asset).filter(Asset.parent_asset_id.in_(source_ids)).update(
{Asset.parent_asset_id: target.id}, synchronize_session=False
)
if target.parent_asset_id in source_ids:
target.parent_asset_id = None
_record_asset_history(
db,
target,
{
"Duplikate zusammengeführt": {
"old": ", ".join(str(asset_id) for asset_id in source_ids),
"new": str(target.id),
},
**changes,
},
source="duplicate-merge",
changed_by=_changed_by(request),
)
for source in sources:
db.delete(source)
try:
db.commit()
except IntegrityError:
db.rollback()
message = _translate_request(
request,
"duplicates.merge_unique_error",
"The duplicates could not be merged because a unique value is already assigned to another asset.",
)
# Do not redirect back to /assets/merge/preview: that page is POST-only.
# A GET redirect there caused ERR_TOO_MANY_REDIRECTS instead of a toast.
return RedirectResponse(f"/assets?toast_error={quote(message)}", status_code=303)
return RedirectResponse(
f"/assets/{target.id}?toast_success={quote(f'{len(sources)} Duplikate wurden in Asset {target.id} zusammengeführt')}",
status_code=303,
)
@app.post("/assets/bulk-edit")
def assets_bulk_edit_form(request: Request, asset_ids: list[str] = Form([]), category_id: int | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
ids = _selected_asset_ids(asset_ids)
if not ids:
raise HTTPException(400, "Bitte mindestens ein Asset markieren.")
assets = (db.query(Asset).options(joinedload(Asset.category).joinedload(Category.visible_fields)).filter(Asset.id.in_(ids)).order_by(Asset.name).all())
if len(assets) != len(ids):
raise HTTPException(404, "Mindestens ein ausgewähltes Asset wurde nicht gefunden.")
fields = [{"field_name": "category_id", "label": "Kategorie"}] + _bulk_edit_fields_for_assets(assets)
categories = _sort_localized(db.query(Category).all(), "name", request)
return templates.TemplateResponse("assets_bulk_edit.html", {
"request": request,
"assets": assets,
"asset_ids": ids,
"fields": fields,
"status_options": _active_status_options(db, request),
"parent_candidates": _sort_localized(db.query(Asset).filter(~Asset.id.in_(ids)).all(), "name", request),
"categories": categories,
"category_id": category_id,
})
@app.post("/assets/bulk-edit/apply")
def assets_bulk_edit_apply(request: Request, asset_ids: list[str] = Form([]), field_name: str = Form(...), new_value: str = Form(""), category_id: int | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
ids = _selected_asset_ids(asset_ids)
assets = (db.query(Asset).options(joinedload(Asset.category).joinedload(Category.visible_fields)).filter(Asset.id.in_(ids)).all())
if not ids or len(assets) != len(ids):
raise HTTPException(400, "Die Auswahl ist leer oder nicht mehr vollständig vorhanden.")
allowed = {"category_id": "Kategorie", **{item["field_name"]: item["label"] for item in _bulk_edit_fields_for_assets(assets)}}
if field_name not in allowed:
raise HTTPException(400, "Das gewählte Feld ist nicht in allen ausgewählten Gerätekategorien zugelassen.")
value = _normalize_optional_text(new_value)
parent_asset = None
target_category = None
if field_name == "category_id":
try:
target_category_id = int(value) if value else None
except (TypeError, ValueError) as exc:
raise HTTPException(400, "Die Kategorie ist ungültig.") from exc
if target_category_id is None:
raise HTTPException(400, "Bitte eine Kategorie auswählen.")
target_category = db.query(Category).filter(Category.id == target_category_id).first()
if not target_category:
raise HTTPException(404, "Die gewählte Kategorie wurde nicht gefunden.")
value = target_category_id
if field_name == "parent_asset_id":
try:
parent_id = int(value) if value else None
except (TypeError, ValueError) as exc:
raise HTTPException(400, "Das übergeordnete Gerät ist ungültig.") from exc
if parent_id in ids:
raise HTTPException(400, "Ein ausgewähltes Asset kann nicht gleichzeitig als übergeordnetes Gerät für die gesamte Auswahl verwendet werden.")
if parent_id is not None:
parent_asset = db.query(Asset).filter(Asset.id == parent_id).first()
if not parent_asset:
raise HTTPException(404, "Das gewählte übergeordnete Gerät wurde nicht gefunden.")
value = parent_id
for asset in assets:
if parent_id is not None and _would_create_parent_cycle(db, asset.id, parent_id):
raise HTTPException(400, f"Die Zuordnung würde beim Asset {asset.name} einen Kreis in der Baumstruktur erzeugen.")
changed = 0
for asset in assets:
old_value = getattr(asset, field_name, None)
if old_value == value or _display_value(old_value) == _display_value(value):
continue
if field_name == "parent_asset_id":
old_parent_id = int(old_value) if old_value not in (None, "") else None
old_parent = db.query(Asset).filter(Asset.id == old_parent_id).first() if old_parent_id else None
history_old = old_parent.name if old_parent else ""
history_new = parent_asset.name if parent_asset else ""
elif field_name == "category_id":
history_old = asset.category.name if asset.category else str(old_value or "")
history_new = target_category.name
else:
history_old = _display_value(old_value)
history_new = _display_value(value)
setattr(asset, field_name, value)
_record_asset_history(db, asset, {allowed[field_name]: {"old": history_old, "new": history_new}}, source="bulk-update", changed_by=_changed_by(request))
changed += 1
try:
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Die Daten konnten nicht für alle Assets gespeichert werden.") from exc
target = f"/assets?category_id={category_id}" if category_id else "/assets"
return RedirectResponse(f"{target}{'&' if '?' in target else '?'}toast_success={quote(str(changed) + ' Assets aktualisiert')}", status_code=303)
@app.post("/assets/bulk-delete")
def assets_bulk_delete(request: Request, asset_ids: list[str] = Form([]), category_id: int | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
ids = _selected_asset_ids(asset_ids)
if not ids:
raise HTTPException(400, "Bitte mindestens ein Asset markieren.")
assets = db.query(Asset).filter(Asset.id.in_(ids)).all()
if len(assets) != len(ids):
raise HTTPException(404, "Mindestens ein ausgewähltes Asset wurde nicht gefunden.")
image_paths = [asset.image_path for asset in assets if asset.image_path]
for asset in assets:
db.delete(asset)
db.commit()
for image_path in image_paths:
delete_asset_image(image_path)
target = f"/assets?category_id={category_id}" if category_id else "/assets"
return RedirectResponse(f"{target}{'&' if '?' in target else '?'}toast_success={quote(str(len(assets)) + ' Assets gelöscht')}", status_code=303)
@app.post("/assets/{asset_id}/delete")
def asset_delete(asset_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
asset = _get_visible_asset(db, request, asset_id)
image_path = asset.image_path
db.delete(asset)
db.commit()
delete_asset_image(image_path)
return RedirectResponse("/assets", status_code=303)
@app.get("/status-options")
def status_options_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
options = db.query(StatusOption).order_by(StatusOption.sort_order, StatusOption.name).all()
return templates.TemplateResponse("status_options.html", {"request": request, "options": options})
@app.post("/status-options")
async def status_options_save(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
existing = db.query(StatusOption).all()
for option in existing:
name = (form.get(f"name_{option.id}") or "").strip()
if not name:
continue
option.name = name
option.active = form.get(f"active_{option.id}") == "on"
try:
option.sort_order = int(form.get(f"sort_{option.id}") or option.sort_order)
except ValueError:
pass
issue_id = int(form.get("issue_status_id") or 0)
return_id = int(form.get("return_status_id") or 0)
for option in existing:
option.use_for_issue = option.id == issue_id
option.use_for_return = option.id == return_id
new_name = (form.get("new_name") or "").strip()
if new_name:
max_order = max((item.sort_order for item in existing), default=-1)
db.add(StatusOption(name=new_name, sort_order=max_order + 1, active=True))
try:
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Statuswerte müssen eindeutig sein.") from exc
return RedirectResponse("/status-options", status_code=303)
@app.post("/status-options/{option_id}/delete")
def status_option_delete(option_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
option = db.query(StatusOption).filter(StatusOption.id == option_id).first()
if not option:
raise HTTPException(404, "Statuswert nicht gefunden")
db.delete(option)
db.commit()
return RedirectResponse("/status-options", status_code=303)
@app.get("/categories")
def categories_list(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
categories = (
db.query(Category)
.options(selectinload(Category.visible_fields))
.order_by(Category.name)
.all()
)
asset_counts = dict(
db.query(Asset.category_id, func.count(Asset.id))
.group_by(Asset.category_id)
.all()
)
return templates.TemplateResponse(
"categories.html",
{"request": request, "categories": categories, "asset_counts": asset_counts,
"definitions": db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).order_by(FieldDefinition.sort_order, FieldDefinition.label).all()},
)
@app.post("/categories/bulk-edit")
async def categories_bulk_edit(
request: Request,
category_ids: list[str] = Form([]),
field_name: str = Form(...),
field_value: str = Form(""),
clear_value: str | None = Form(None),
image: UploadFile | None = File(None),
db: Session = Depends(get_db),
):
_require_admin(request)
ids = _selected_asset_ids(category_ids)
if not ids:
message = _translate_request(request, "categories.bulk.no_selection", "No categories selected")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
categories = db.query(Category).filter(Category.id.in_(ids)).all()
if not categories:
message = _translate_request(request, "categories.bulk.no_selection", "No categories selected")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
clear = clear_value == "on"
if field_name == "name":
if len(categories) != 1:
message = _translate_request(request, "categories.bulk.name_single_only", "The name can only be changed for one category at a time")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
value = field_value.strip()
if not value:
message = _translate_request(request, "categories.bulk.name_required", "The category name must not be empty")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
categories[0].name = value
elif field_name == "description":
value = None if clear else (field_value.strip() or None)
for category in categories:
category.description = value
elif field_name == "image":
if clear:
for category in categories:
category.image_path = None
else:
new_image = save_upload(image)
if not new_image:
message = _translate_request(request, "categories.bulk.image_required", "Select an image or choose clear value")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
for category in categories:
category.image_path = new_image
else:
raise HTTPException(400, _translate_request(request, "categories.bulk.unsupported_field", "Unsupported category field"))
try:
db.commit()
except IntegrityError:
db.rollback()
message = _translate_request(request, "categories.bulk.name_duplicate", "The category name already exists")
return RedirectResponse(f"/categories?toast_error={quote(message)}", status_code=303)
message = _translate_request(request, "categories.bulk.updated", "Updated {count} categories", count=len(categories))
return RedirectResponse(f"/categories?toast_success={quote(message)}", status_code=303)
@app.post("/categories/bulk-fields")
async def categories_bulk_fields(
request: Request,
category_ids: list[str] = Form([]),
operation: str = Form(...),
field_ids: list[str] = Form([]),
field_definition_id: str = Form(""),
apply_required: str | None = Form(None),
required_value: str | None = Form(None),
apply_show_in_list: str | None = Form(None),
show_in_list_value: str | None = Form(None),
source_category_id: str = Form(""),
db: Session = Depends(get_db),
):
"""Manage category field assignments without touching generic table behavior."""
_require_admin(request)
ids = _selected_asset_ids(category_ids)
if not ids:
message = _translate_request(request, "categories.fields.no_selection", "No categories selected")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
categories = (
db.query(Category)
.options(selectinload(Category.visible_fields).joinedload(CategoryField.definition))
.filter(Category.id.in_(ids))
.all()
)
if not categories:
message = _translate_request(request, "categories.fields.no_selection", "No categories selected")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
definitions = {item.id: item for item in db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).all()}
selected_definition_ids = _selected_asset_ids(field_ids)
changed = 0
if operation in {"add", "remove"}:
selected_definitions = [definitions[item_id] for item_id in selected_definition_ids if item_id in definitions]
if not selected_definitions:
message = _translate_request(request, "categories.fields.choose_fields", "Select at least one field")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
for category in categories:
by_definition = {link.field_definition_id: link for link in category.visible_fields if link.field_definition_id}
by_name = {link.field_name: link for link in category.visible_fields}
next_order = max((link.sort_order for link in category.visible_fields), default=-1) + 1
for definition in selected_definitions:
link = by_definition.get(definition.id) or by_name.get(definition.field_name)
if operation == "add":
if link is None:
link = CategoryField(
category_id=category.id,
field_name=definition.field_name,
label=definition.label,
field_definition_id=definition.id,
active=True,
required=definition.field_name == "name",
show_in_list=definition.field_name == "name",
readonly=definition.readonly,
sort_order=next_order,
)
next_order += 1
db.add(link)
category.visible_fields.append(link)
changed += 1
elif not link.active:
link.active = True
if definition.field_name == "name":
link.required = True
link.show_in_list = True
changed += 1
else:
if definition.field_name == "name":
continue
if link is not None and (link.active or link.required or link.show_in_list):
link.active = False
link.required = False
link.show_in_list = False
changed += 1
elif operation == "settings":
try:
definition_id = int(field_definition_id)
except (TypeError, ValueError):
definition_id = 0
definition = definitions.get(definition_id)
if not definition:
message = _translate_request(request, "categories.fields.choose_field", "Select a field")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
change_required = apply_required == "on"
change_list = apply_show_in_list == "on"
if not change_required and not change_list:
message = _translate_request(request, "categories.fields.choose_setting", "Select at least one setting to change")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
required = required_value == "on"
show_in_list = show_in_list_value == "on"
for category in categories:
link = next((item for item in category.visible_fields if item.field_definition_id == definition.id or item.field_name == definition.field_name), None)
if link is None or not link.active:
continue
before = (link.required, link.show_in_list)
if change_required and definition.field_name != "name" and not definition.readonly:
link.required = required
if change_list:
link.show_in_list = True if definition.field_name == "name" else show_in_list
if before != (link.required, link.show_in_list):
changed += 1
elif operation == "copy":
try:
source_id = int(source_category_id)
except (TypeError, ValueError):
source_id = 0
source = (
db.query(Category)
.options(selectinload(Category.visible_fields).joinedload(CategoryField.definition))
.filter(Category.id == source_id)
.first()
)
if not source:
message = _translate_request(request, "categories.fields.choose_source", "Select a source category")
return RedirectResponse(f"/categories?toast_warning={quote(message)}", status_code=303)
source_links = sorted(source.visible_fields, key=lambda item: item.sort_order)
for category in categories:
if category.id == source.id:
continue
target_by_name = {item.field_name: item for item in category.visible_fields}
source_names = {item.field_name for item in source_links}
# A template is an exact field-configuration copy. Existing links are
# retained but deactivated when absent from the source, preserving data.
for target in category.visible_fields:
if target.field_name not in source_names and target.field_name != "name":
target.active = False
target.required = False
target.show_in_list = False
for source_link in source_links:
target = target_by_name.get(source_link.field_name)
if target is None:
target = CategoryField(
category_id=category.id,
field_name=source_link.field_name,
label=source_link.label,
field_definition_id=source_link.field_definition_id,
readonly=source_link.readonly,
)
db.add(target)
category.visible_fields.append(target)
target.active = source_link.active or source_link.field_name == "name"
target.required = True if source_link.field_name == "name" else source_link.required
target.show_in_list = True if source_link.field_name == "name" else source_link.show_in_list
target.sort_order = source_link.sort_order
changed += 1
else:
raise HTTPException(400, _translate_request(request, "categories.fields.unsupported_operation", "Unsupported field operation"))
db.commit()
message = _translate_request(request, "categories.fields.updated", "Field configuration updated for {count} categories", count=changed)
return RedirectResponse(f"/categories?toast_success={quote(message)}", status_code=303)
@app.post("/categories/activate-default-fields")
def categories_activate_default_fields(
request: Request,
category_ids: list[str] = Form([]),
db: Session = Depends(get_db),
):
_require_admin(request)
ids = _selected_asset_ids(category_ids)
if not ids:
return RedirectResponse("/categories?toast_warning=Keine Kategorien ausgewählt", status_code=303)
categories = (
db.query(Category)
.options(selectinload(Category.visible_fields).joinedload(CategoryField.definition))
.filter(Category.id.in_(ids))
.all()
)
default_definitions = (
db.query(FieldDefinition)
.filter(FieldDefinition.is_active.is_(True), FieldDefinition.is_default.is_(True))
.all()
)
if not default_definitions:
return RedirectResponse("/categories?toast_warning=Es sind keine Standardfelder markiert", status_code=303)
activated = 0
created_links = 0
for category in categories:
by_definition = {link.field_definition_id: link for link in category.visible_fields if link.field_definition_id}
by_name = {link.field_name: link for link in category.visible_fields}
next_order = max((link.sort_order for link in category.visible_fields), default=-1) + 1
for definition in default_definitions:
link = by_definition.get(definition.id) or by_name.get(definition.field_name)
if link is None:
link = CategoryField(
category_id=category.id,
field_name=definition.field_name,
label=definition.label,
field_definition_id=definition.id,
active=True,
required=definition.field_name == "name",
show_in_list=True,
readonly=definition.readonly,
sort_order=next_order,
)
next_order += 1
db.add(link)
created_links += 1
else:
if not link.active:
link.active = True
activated += 1
link.show_in_list = True
if definition.field_name == "name":
link.required = True
db.commit()
message = f"Standardfelder aktiviert: {activated} vorhandene und {created_links} neue Zuordnungen"
return RedirectResponse(f"/categories?toast_success={quote(message)}", status_code=303)
@app.get("/categories/new")
def category_new(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
definitions=db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).order_by(FieldDefinition.sort_order, FieldDefinition.label).all()
return templates.TemplateResponse("category_form.html", {"request": request, "category": None, "definitions": definitions, "selected_fields": {}})
@app.post("/categories/new")
async def category_create(request: Request, name: str = Form(...), description: str = Form(""), image: UploadFile | None = File(None), db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
category = Category(name=name.strip(), description=description or None, image_path=save_upload(image))
db.add(category)
db.flush()
definitions = db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).order_by(FieldDefinition.sort_order).all()
for order, definition in enumerate(definitions):
active = form.get(f"active_{definition.id}") == "on" or definition.field_name == "name"
db.add(CategoryField(category_id=category.id, field_name=definition.field_name, label=definition.label, field_definition_id=definition.id, active=active, required=True if definition.field_name == "name" else form.get(f"required_{definition.id}") == "on", show_in_list=True if definition.field_name == "name" else form.get(f"list_{definition.id}") == "on", readonly=definition.readonly, sort_order=order))
try:
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Kategorie konnte nicht gespeichert werden. Der Kategoriename ist möglicherweise bereits vorhanden.") from exc
return RedirectResponse("/categories", status_code=303)
@app.get("/categories/{category_id}/edit")
def category_edit(category_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
category = db.query(Category).options(joinedload(Category.visible_fields)).filter(Category.id == category_id).first()
if not category:
raise HTTPException(404, "Kategorie nicht gefunden")
selected_fields = {f.field_name: f for f in category.visible_fields}
definitions=db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).order_by(FieldDefinition.sort_order, FieldDefinition.label).all()
return templates.TemplateResponse("category_form.html", {"request": request, "category": category, "definitions": definitions, "selected_fields": selected_fields})
@app.post("/categories/{category_id}/edit")
async def category_update(category_id: int, request: Request, name: str = Form(...), description: str = Form(""), image: UploadFile | None = File(None), db: Session = Depends(get_db)):
_require_admin(request)
category = db.query(Category).options(joinedload(Category.visible_fields)).filter(Category.id == category_id).first()
if not category:
raise HTTPException(404, "Kategorie nicht gefunden")
form = await request.form()
category.name = name.strip()
category.description = description.strip() or None
new_image = save_upload(image)
if new_image:
category.image_path = new_image
# Vorhandene Definitionen explizit löschen und flushen. Das verhindert
# UniqueConstraint-Verletzungen bei anschließend identischen Feldnamen.
db.query(CategoryField).filter(CategoryField.category_id == category.id).delete(synchronize_session=False)
db.flush()
definitions = db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).order_by(FieldDefinition.sort_order).all()
for order, definition in enumerate(definitions):
active = form.get(f"active_{definition.id}") == "on" or definition.field_name == "name"
db.add(CategoryField(
category_id=category.id, field_name=definition.field_name, label=definition.label, field_definition_id=definition.id, active=active,
required=True if definition.field_name == "name" else form.get(f"required_{definition.id}") == "on",
show_in_list=True if definition.field_name == "name" else form.get(f"list_{definition.id}") == "on",
readonly=definition.readonly, sort_order=order,
))
try:
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Kategorie konnte nicht gespeichert werden. Der Kategoriename ist möglicherweise bereits vorhanden.") from exc
return RedirectResponse("/categories", status_code=303)
@app.post("/categories/{category_id}/duplicate")
def category_duplicate(category_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
category = (
db.query(Category)
.options(joinedload(Category.visible_fields))
.filter(Category.id == category_id)
.first()
)
if not category:
raise HTTPException(404, "Kategorie nicht gefunden")
base_name = f"{category.name} (Kopie)"
new_name = base_name
number = 2
while db.query(Category.id).filter(Category.name == new_name).first():
new_name = f"{category.name} (Kopie {number})"
number += 1
duplicate = Category(
name=new_name,
description=category.description,
image_path=category.image_path,
)
db.add(duplicate)
db.flush()
for field in sorted(category.visible_fields, key=lambda item: item.sort_order):
db.add(CategoryField(
category_id=duplicate.id,
field_name=field.field_name,
label=field.label,
field_definition_id=field.field_definition_id,
active=field.active,
required=field.required,
show_in_list=field.show_in_list,
readonly=field.readonly,
sort_order=field.sort_order,
))
db.commit()
return RedirectResponse(
f"/categories/{duplicate.id}/edit?toast_success=Kategorie wurde dupliziert",
status_code=303,
)
@app.post("/categories/{category_id}/delete")
def category_delete(category_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
category = db.query(Category).filter(Category.id == category_id).first()
if not category:
raise HTTPException(404, "Kategorie nicht gefunden")
asset_count = db.query(Asset).filter(Asset.category_id == category.id).count()
if asset_count > 0:
return RedirectResponse(
f"/categories?toast_error=Kategorie kann nicht gelöscht werden: Sie wird noch von {asset_count} Asset(s) verwendet",
status_code=303,
)
db.delete(category)
db.commit()
return RedirectResponse(
"/categories?toast_success=Kategorie wurde gelöscht",
status_code=303,
)
@app.get("/api/software-callback/health")
def software_callback_health():
return {
"ok": True,
"service": "assetmanager-software-callback",
"version": application_version(),
"time": datetime.utcnow().isoformat() + "Z",
}
def _software_overview_statistics(db: Session) -> dict[str, int]:
"""Return only the counters required by the software overview.
Do not load every InstalledSoftware ORM row merely to display three numbers.
PostgreSQL performs the counting and grouping and returns only scalar values.
"""
total_entries, initialized_assets = (
db.query(
func.count(InstalledSoftware.id),
func.count(func.distinct(InstalledSoftware.asset_id)),
)
.one()
)
# Fetch only the three short grouping columns. The final Python casefold
# intentionally preserves the exact grouping semantics used by the existing
# detailed summary (Unicode casefold is stronger than SQL LOWER).
summary_rows = (
db.query(
InstalledSoftware.name,
InstalledSoftware.publisher,
InstalledSoftware.platform,
)
.distinct()
.all()
)
summary_count = len({
(
(name or "").strip().casefold(),
(publisher or "").strip().casefold(),
(platform or "unknown").strip().casefold(),
)
for name, publisher, platform in summary_rows
})
return {
"software_entry_count": int(total_entries or 0),
"software_summary_count": int(summary_count or 0),
"initialized_assets": int(initialized_assets or 0),
}
def _software_group_expressions():
"""Return normalized SQL expressions used by all software groupings."""
name_key = func.lower(func.trim(func.coalesce(InstalledSoftware.name, "")))
publisher_key = func.lower(func.trim(func.coalesce(InstalledSoftware.publisher, "")))
platform_key = func.lower(func.trim(func.coalesce(InstalledSoftware.platform, "unknown")))
version_key = func.lower(func.trim(func.coalesce(InstalledSoftware.version, "")))
return name_key, publisher_key, platform_key, version_key
def _software_summary_rows(db: Session, grouping: str) -> list[dict[str, Any]]:
"""Build the aggregated software list inside PostgreSQL.
The query returns only already aggregated result rows. ``grouping=name_platform``
combines changing publisher names and all versions by software name and
platform. ``grouping=software`` combines all versions by name, publisher and
platform. ``grouping=version`` keeps one row per version.
"""
name_key, publisher_key, platform_key, version_key = _software_group_expressions()
common_columns = (
name_key.label("name_key"),
publisher_key.label("publisher_key"),
platform_key.label("platform_key"),
func.min(func.coalesce(InstalledSoftware.name, "")).label("name"),
func.min(func.coalesce(InstalledSoftware.publisher, "")).label("publisher"),
func.min(func.coalesce(InstalledSoftware.platform, "unknown")).label("platform"),
)
if grouping == "name_platform":
publishers_array = func.array_agg(
distinct(func.nullif(func.trim(func.coalesce(InstalledSoftware.publisher, "")), ""))
)
versions_array = func.array_agg(
distinct(func.nullif(func.trim(func.coalesce(InstalledSoftware.version, "")), ""))
)
query_rows = (
db.query(
name_key.label("name_key"),
platform_key.label("platform_key"),
func.min(func.coalesce(InstalledSoftware.name, "")).label("name"),
func.min(func.coalesce(InstalledSoftware.platform, "unknown")).label("platform"),
publishers_array.label("publishers"),
versions_array.label("versions"),
func.count(func.distinct(InstalledSoftware.asset_id)).label("asset_count"),
func.count(InstalledSoftware.id).label("entries"),
func.max(InstalledSoftware.last_seen_at).label("last_seen_at"),
)
.group_by(name_key, platform_key)
.order_by(platform_key, name_key)
.all()
)
result: list[dict[str, Any]] = []
for row in query_rows:
publishers = sorted(
{str(value).strip() for value in (row.publishers or []) if value and str(value).strip()},
key=str.casefold,
)
versions = sorted(
{str(value).strip() for value in (row.versions or []) if value and str(value).strip()},
key=str.casefold,
)
result.append({
"name": row.name or "",
"publisher": ", ".join(publishers),
"platform": row.platform or "unknown",
"versions": ", ".join(versions),
"version_count": len(versions),
"asset_count": int(row.asset_count or 0),
"entries": int(row.entries or 0),
"last_seen_at": row.last_seen_at,
})
return result
if grouping == "version":
query_rows = (
db.query(
*common_columns,
version_key.label("version_key"),
func.min(func.coalesce(InstalledSoftware.version, "")).label("version"),
func.count(func.distinct(InstalledSoftware.asset_id)).label("asset_count"),
func.count(InstalledSoftware.id).label("entries"),
func.max(InstalledSoftware.last_seen_at).label("last_seen_at"),
)
.group_by(name_key, publisher_key, platform_key, version_key)
.order_by(platform_key, name_key, publisher_key, version_key)
.all()
)
return [
{
"name": row.name or "",
"publisher": row.publisher or "",
"platform": row.platform or "unknown",
"versions": row.version or "",
"version_count": 1 if (row.version or "").strip() else 0,
"asset_count": int(row.asset_count or 0),
"entries": int(row.entries or 0),
"last_seen_at": row.last_seen_at,
}
for row in query_rows
]
# PostgreSQL performs the three-column grouping and returns the distinct
# version values as a compact array. Assets are counted distinctly across
# all versions, so the value remains correct even if an asset contains more
# than one version record for the same software.
versions_array = func.array_agg(
distinct(func.nullif(func.trim(func.coalesce(InstalledSoftware.version, "")), ""))
)
query_rows = (
db.query(
*common_columns,
versions_array.label("versions"),
func.count(func.distinct(InstalledSoftware.asset_id)).label("asset_count"),
func.count(InstalledSoftware.id).label("entries"),
func.max(InstalledSoftware.last_seen_at).label("last_seen_at"),
)
.group_by(name_key, publisher_key, platform_key)
.order_by(platform_key, name_key, publisher_key)
.all()
)
result: list[dict[str, Any]] = []
for row in query_rows:
versions = sorted(
{str(value).strip() for value in (row.versions or []) if value and str(value).strip()},
key=str.casefold,
)
result.append({
"name": row.name or "",
"publisher": row.publisher or "",
"platform": row.platform or "unknown",
"versions": ", ".join(versions),
"version_count": len(versions),
"asset_count": int(row.asset_count or 0),
"entries": int(row.entries or 0),
"last_seen_at": row.last_seen_at,
})
return result
def _software_installation_rows(db: Session):
"""Return only columns rendered by the detailed installation table."""
return (
db.query(
InstalledSoftware.id.label("id"),
InstalledSoftware.asset_id.label("asset_id"),
Asset.name.label("asset_name"),
InstalledSoftware.name.label("name"),
InstalledSoftware.version.label("version"),
InstalledSoftware.publisher.label("publisher"),
InstalledSoftware.architecture.label("architecture"),
InstalledSoftware.platform.label("platform"),
InstalledSoftware.install_date.label("install_date"),
InstalledSoftware.last_seen_at.label("last_seen_at"),
)
.join(Asset, InstalledSoftware.asset_id == Asset.id)
.order_by(
InstalledSoftware.name.asc(),
Asset.name.asc(),
InstalledSoftware.version.asc(),
InstalledSoftware.id.asc(),
)
.all()
)
_SOFTWARE_JOB_FILTER_KEYS = (
"job_id",
"created_at",
"asset",
"job_type",
"package",
"action",
"priority",
"platform",
"status",
)
def _software_job_filter_values(request: Request) -> dict[str, str]:
"""Read and normalize the server-side filters of the recent-job table."""
return {
key: str(request.query_params.get(f"job_filter_{key}") or "").strip()[:200]
for key in _SOFTWARE_JOB_FILTER_KEYS
}
def _sql_text_contains(column: Any, value: str):
"""Case-insensitive literal substring search for SQL text expressions."""
escaped = (
str(value or "")
.casefold()
.replace("\\", "\\\\")
.replace("%", "\\%")
.replace("_", "\\_")
)
return func.lower(func.coalesce(cast(column, String), "")).like(
f"%{escaped}%",
escape="\\",
)
def _translated_filter_codes(
request: Request,
prefix: str,
codes: tuple[str, ...],
search_value: str,
) -> list[str]:
"""Map translated, partially entered labels back to their internal codes.
The translation dictionaries are cached on the request so one filter
request does not open a database session for every possible status value.
"""
needle = str(search_value or "").strip().casefold()
if not needle:
return []
cache = getattr(request.state, "software_job_filter_translations", None)
if cache is None:
user = _session_user(request) or {}
language_code = (
user.get("language_code")
or request.session.get("language_code")
or load_config().get("general", {}).get("default_language", "en")
)
translation_db = SessionLocal()
try:
cache = {
"own": translation_dictionary(translation_db, language_code),
"en": translation_dictionary(translation_db, "en"),
}
finally:
translation_db.close()
request.state.software_job_filter_translations = cache
matches: list[str] = []
for code in codes:
key = f"{prefix}.{code}"
label = cache["own"].get(key) or cache["en"].get(key) or code
if needle in code.casefold() or needle in str(label).casefold():
matches.append(code)
return matches
def _software_job_date_filter(column: Any, search_value: str):
"""Support ISO/database text and the common German date formats."""
value = str(search_value or "").strip()
conditions = [_sql_text_contains(column, value)]
for date_format in ("%d.%m.%Y", "%d.%m.%y", "%Y-%m-%d"):
try:
parsed = datetime.strptime(value, date_format)
except ValueError:
continue
conditions.append(and_(column >= parsed, column < parsed + timedelta(days=1)))
break
return or_(*conditions)
def _apply_software_job_filters(query: Any, request: Request, filters: dict[str, str]):
"""Apply all recent-job filters before LIMIT is evaluated.
This keeps an old job searchable even when the visible result size is only
25 rows. The large script, result and log columns are still excluded from
the final SELECT by ``load_only`` below.
"""
parameters_text = cast(SoftwareJob.parameters, String)
if filters["job_id"]:
value = filters["job_id"]
if value.isdigit():
query = query.filter(SoftwareJob.id == int(value))
else:
query = query.filter(_sql_text_contains(SoftwareJob.id, value))
if filters["created_at"]:
query = query.filter(
_software_job_date_filter(SoftwareJob.created_at, filters["created_at"])
)
if filters["asset"]:
query = query.filter(_sql_text_contains(Asset.name, filters["asset"]))
if filters["job_type"]:
value = filters["job_type"]
conditions = [
_sql_text_contains(SoftwareJob.job_type, value),
_sql_text_contains(parameters_text, value),
]
job_type_codes = _translated_filter_codes(
request,
"jobs.type",
("software_inventory", "job_definition"),
value,
)
if job_type_codes:
conditions.append(SoftwareJob.job_type.in_(job_type_codes))
job_kind_codes = _translated_filter_codes(
request,
"jobdefs.kind",
("script", "registry"),
value,
)
for kind in job_kind_codes:
conditions.append(
and_(
SoftwareJob.job_type == "job_definition",
_sql_text_contains(parameters_text, kind),
)
)
query = query.filter(or_(*conditions))
if filters["package"]:
value = filters["package"]
query = query.filter(
or_(
_sql_text_contains(SoftwarePackage.name, value),
_sql_text_contains(parameters_text, value),
)
)
if filters["action"]:
value = filters["action"]
conditions = [_sql_text_contains(SoftwareJob.action, value)]
action_codes = _translated_filter_codes(
request,
"jobs.action",
("run", "install", "uninstall", "update"),
value,
)
if action_codes:
conditions.append(SoftwareJob.action.in_(action_codes))
query = query.filter(or_(*conditions))
if filters["priority"]:
value = filters["priority"]
if re.fullmatch(r"[-+]?\d+", value):
query = query.filter(SoftwareJob.priority == int(value))
else:
query = query.filter(_sql_text_contains(SoftwareJob.priority, value))
if filters["platform"]:
query = query.filter(
_sql_text_contains(SoftwareJob.platform, filters["platform"])
)
if filters["status"]:
value = filters["status"]
conditions = [_sql_text_contains(SoftwareJob.status, value)]
status_codes = _translated_filter_codes(
request,
"software.status",
(
"created",
"sending",
"sent",
"running",
"success",
"partial",
"failed",
"timeout",
"cancelled",
),
value,
)
if status_codes:
conditions.append(SoftwareJob.status.in_(status_codes))
query = query.filter(or_(*conditions))
return query
@app.get("/software")
def software_overview(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
page_started = time.perf_counter()
timeout_started = time.perf_counter()
_expire_stale_software_jobs(db)
timeout_seconds = time.perf_counter() - timeout_started
allowed_job_limits = (25, 50, 100, 500, 1000)
raw_limit = request.query_params.get("job_limit")
if raw_limit is None:
raw_limit = request.session.get("software_job_limit", 25)
try:
requested_limit = int(raw_limit)
except (TypeError, ValueError):
requested_limit = 25
job_limit = requested_limit if requested_limit in allowed_job_limits else 25
request.session["software_job_limit"] = job_limit
stats_started = time.perf_counter()
statistics = _software_overview_statistics(db)
stats_seconds = time.perf_counter() - stats_started
jobs_started = time.perf_counter()
job_filters = _software_job_filter_values(request)
filters_active = any(job_filters.values())
jobs_query = (
db.query(SoftwareJob)
.join(Asset, Asset.id == SoftwareJob.asset_id)
.outerjoin(SoftwarePackage, SoftwarePackage.id == SoftwareJob.package_id)
)
jobs_query = _apply_software_job_filters(jobs_query, request, job_filters)
# Count after filtering but before LIMIT. This value is only needed while a
# filter is active and therefore adds no overhead to the normal overview.
jobs_filtered_total = (
int(jobs_query.with_entities(func.count(SoftwareJob.id)).scalar() or 0)
if filters_active
else 0
)
jobs = (
jobs_query
.options(
load_only(
SoftwareJob.id,
SoftwareJob.asset_id,
SoftwareJob.package_id,
SoftwareJob.status,
SoftwareJob.job_type,
SoftwareJob.action,
SoftwareJob.priority,
SoftwareJob.platform,
SoftwareJob.parameters,
SoftwareJob.created_at,
),
contains_eager(SoftwareJob.asset).load_only(Asset.id, Asset.name, Asset.mesh_node_id),
contains_eager(SoftwareJob.package).load_only(SoftwarePackage.id, SoftwarePackage.name),
)
.order_by(SoftwareJob.created_at.desc(), SoftwareJob.id.desc())
.limit(job_limit)
.all()
)
_decorate_jobs_for_display(db, jobs)
jobs_seconds = time.perf_counter() - jobs_started
response = templates.TemplateResponse(
"software.html",
{
"request": request,
"jobs": jobs,
"job_limit": job_limit,
"job_limit_options": allowed_job_limits,
"job_filters": job_filters,
"jobs_filters_active": filters_active,
"jobs_filtered_total": jobs_filtered_total,
**statistics,
},
)
total_seconds = time.perf_counter() - page_started
if total_seconds >= 0.5:
logger.info(
"Software overview timing: total=%.3fs timeout=%.3fs statistics=%.3fs jobs=%.3fs rows=%s",
total_seconds,
timeout_seconds,
stats_seconds,
jobs_seconds,
len(jobs),
)
return response
@app.get("/software/job-states")
def software_job_states_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
_expire_stale_software_jobs(db)
rows = (
db.query(AssetJobState, Asset, JobDefinition)
.join(Asset, Asset.id == AssetJobState.asset_id)
.outerjoin(JobDefinition, JobDefinition.id == AssetJobState.job_definition_id)
.order_by(AssetJobState.updated_at.desc(), func.lower(Asset.name), AssetJobState.state_key)
.all()
)
job_states: list[dict[str, Any]] = []
for state, asset, definition in rows:
if state.job_type == "job_definition":
job_label = definition.name if definition is not None else state.state_key
job_type_key = "jobs.type.job_definition"
job_type_fallback = "Jobdefinition"
elif state.job_type == "software_inventory":
job_label = _translate_request(request, "jobs.type.software_inventory", "Software inventory")
job_type_key = "jobs.type.software_inventory"
job_type_fallback = "Software inventory"
else:
job_label = state.state_key
job_type_key = f"jobs.type.{state.job_type}"
job_type_fallback = state.job_type
job_states.append({
"state": state,
"asset": asset,
"job_label": job_label,
"job_type_key": job_type_key,
"job_type_fallback": job_type_fallback,
"last_execution": state.last_started_at or state.last_sent_at or state.last_status_at or state.last_created_at,
})
return templates.TemplateResponse(
"software_job_states.html",
{
"request": request,
"job_states": job_states,
},
)
@app.post("/software/inventory/delete")
async def delete_software_inventory_entries(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
redirect_to = str(form.get("redirect_to") or "/software/installations")
if not redirect_to.startswith("/") or redirect_to.startswith("//"):
redirect_to = "/software/installations"
entry_ids = []
for value in form.getlist("entry_ids"):
try:
entry_ids.append(int(value))
except (TypeError, ValueError):
continue
entry_ids = list(dict.fromkeys(entry_ids))
if not entry_ids:
return RedirectResponse(redirect_to + ("&" if "?" in redirect_to else "?") + "toast_error=" + quote(_translate_request(request, "software.delete.none_selected", "No software inventory entries were selected.")), status_code=303)
deleted = db.query(InstalledSoftware).filter(InstalledSoftware.id.in_(entry_ids)).delete(synchronize_session=False)
db.commit()
message = _translate_request(request, "software.delete.result", "Deleted {count} software inventory entries.").format(count=deleted)
return RedirectResponse(redirect_to + ("&" if "?" in redirect_to else "?") + "toast_success=" + quote(message), status_code=303)
@app.get("/software/installations")
def software_installations_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
started = time.perf_counter()
software_entries = _software_installation_rows(db)
query_seconds = time.perf_counter() - started
response = templates.TemplateResponse(
"software_installations.html",
{
"request": request,
"software_entries": software_entries,
},
)
total_seconds = time.perf_counter() - started
if total_seconds >= 0.5:
logger.info(
"Software installations timing: total=%.3fs query=%.3fs rows=%s",
total_seconds,
query_seconds,
len(software_entries),
)
return response
@app.get("/software/summary")
def software_summary_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
started = time.perf_counter()
allowed_groupings = {"name_platform", "software", "version"}
grouping = str(
request.query_params.get("grouping")
or request.session.get("software_summary_grouping")
or "software"
).strip().lower()
if grouping not in allowed_groupings:
grouping = "software"
request.session["software_summary_grouping"] = grouping
software_summary = _software_summary_rows(db, grouping)
query_seconds = time.perf_counter() - started
response = templates.TemplateResponse(
"software_summary.html",
{
"request": request,
"software_summary": software_summary,
"grouping": grouping,
},
)
total_seconds = time.perf_counter() - started
if total_seconds >= 0.5:
logger.info(
"Software summary timing: total=%.3fs query=%.3fs rows=%s grouping=%s",
total_seconds,
query_seconds,
len(software_summary),
grouping,
)
return response
@app.get("/assets/{asset_id}/software")
def asset_software(asset_id: int, request: Request, db: Session = Depends(get_db)):
asset = _get_visible_asset(db, request, asset_id)
software = db.query(InstalledSoftware).filter(InstalledSoftware.asset_id == asset.id).order_by(InstalledSoftware.name, InstalledSoftware.version).all()
jobs = db.query(SoftwareJob).options(joinedload(SoftwareJob.package)).filter(SoftwareJob.asset_id == asset.id).order_by(SoftwareJob.created_at.desc()).limit(50).all()
_decorate_jobs_for_display(db, jobs)
return templates.TemplateResponse("asset_software.html", {"request": request, "asset": asset, "software": software, "jobs": jobs, "is_admin": _is_admin(request)})
@app.post("/assets/{asset_id}/job-parameters/{definition_id}")
async def asset_job_parameter_overrides(asset_id:int, definition_id:int, request:Request, db:Session=Depends(get_db)):
_require_admin(request)
asset=_get_visible_asset(db,request,asset_id); definition=db.get(JobDefinition,definition_id)
if not definition: raise HTTPException(404)
form=await request.form(); valid_ids={x.id:x for x in definition.parameters}
existing={x.definition_parameter_id:x for x in db.query(JobDefinitionAssetParameter).filter(JobDefinitionAssetParameter.asset_id==asset.id, JobDefinitionAssetParameter.definition_parameter_id.in_(list(valid_ids) or [-1])).all()}
for parameter_id,parameter in valid_ids.items():
use_override=str(form.get(f"override_{parameter_id}") or "") == "on"
row=existing.get(parameter_id)
if use_override:
if row is None:
row=JobDefinitionAssetParameter(asset_id=asset.id,definition_parameter_id=parameter_id);db.add(row)
row.value_text=str(form.get(f"value_{parameter_id}") or "");row.updated_by=_changed_by(request)
elif row is not None:
db.delete(row)
db.commit()
return RedirectResponse(f"/assets/{asset.id}?toast_success="+quote(_translate_request(request,"jobparams.saved","Asset-specific job parameters were saved.")),303)
def _job_parameter_dict(job: SoftwareJob) -> dict[str, Any]:
return job.parameters if isinstance(job.parameters, dict) else {}
def _decorate_jobs_for_display(db: Session, jobs: list[SoftwareJob]) -> None:
"""Attach user-facing job metadata without changing technical job fields.
Job definitions intentionally use the technical database values
``job_type=job_definition`` and the shared system package
``Jobdefinition``. Those values are required by the dispatcher but are
not meaningful in the UI. The selected definition and its real kind are
therefore exposed as transient display attributes.
"""
definition_ids: set[int] = set()
definition_revisions: set[int] = set()
for job in jobs:
if job.job_type != "job_definition":
continue
parameters = _job_parameter_dict(job)
try:
definition_ids.add(int(parameters.get("definition_id")))
except (TypeError, ValueError):
continue
try:
definition_revisions.add(int(parameters.get("definition_revision")))
except (TypeError, ValueError):
pass
definitions_by_id: dict[int, JobDefinition] = {}
revision_snapshots: dict[tuple[int, int], dict[str, Any]] = {}
if definition_ids:
definitions_by_id = {
definition.id: definition
for definition in (
db.query(JobDefinition)
.filter(JobDefinition.id.in_(definition_ids))
.all()
)
}
if definition_revisions:
revision_snapshots = {
(row.definition_id, row.revision): (row.snapshot if isinstance(row.snapshot, dict) else {})
for row in (
db.query(JobDefinitionRevision)
.filter(
JobDefinitionRevision.definition_id.in_(definition_ids),
JobDefinitionRevision.revision.in_(definition_revisions),
)
.all()
)
}
for job in jobs:
parameters = _job_parameter_dict(job)
package_name = job.package.name if job.package else ""
job.display_job_type_key = f"jobs.type.{job.job_type}"
job.display_job_type_fallback = job.job_type
job.display_package_name = package_name
job.display_definition_revision = None
job.display_creation_mode = str(parameters.get("_creation_mode") or "single")
job.display_bulk_batch_id = str(parameters.get("_bulk_batch_id") or "")
if job.job_type != "job_definition":
continue
definition = None
definition_id = None
definition_revision = None
try:
definition_id = int(parameters.get("definition_id"))
definition = definitions_by_id.get(definition_id)
except (TypeError, ValueError):
pass
try:
definition_revision = int(parameters.get("definition_revision"))
except (TypeError, ValueError):
pass
revision_snapshot = (
revision_snapshots.get((definition_id, definition_revision), {})
if definition_id is not None and definition_revision is not None
else {}
)
definition_name = str(
parameters.get("definition_name")
or revision_snapshot.get("name")
or ""
).strip()
if not definition_name and definition is not None:
definition_name = definition.name
definition_kind = str(
parameters.get("definition_job_kind")
or revision_snapshot.get("job_kind")
or ""
).strip().lower()
if not definition_kind and definition is not None:
definition_kind = str(definition.job_kind or "script").strip().lower()
if definition_kind not in {"script", "registry"}:
definition_kind = "script"
job.display_job_type_key = f"jobdefs.kind.{definition_kind}"
job.display_job_type_fallback = definition_kind
job.display_package_name = definition_name or package_name
job.display_definition_revision = (
parameters.get("definition_revision")
if parameters.get("definition_revision") is not None
else (definition.revision if definition is not None else None)
)
def _job_creation_metadata(
*,
creation_mode: str = "single",
bulk_batch_id: str | None = None,
bulk_position: int | None = None,
bulk_total: int | None = None,
) -> dict[str, Any]:
metadata: dict[str, Any] = {"_creation_mode": "bulk" if creation_mode == "bulk" else "single"}
if bulk_batch_id:
metadata["_bulk_batch_id"] = bulk_batch_id
if bulk_position is not None:
metadata["_bulk_position"] = int(bulk_position)
if bulk_total is not None:
metadata["_bulk_total"] = int(bulk_total)
return metadata
def _create_definition_job(
db: Session,
request: Request,
asset: Asset,
definition: JobDefinition,
*,
creation_mode: str = "single",
bulk_batch_id: str | None = None,
bulk_position: int | None = None,
bulk_total: int | None = None,
) -> SoftwareJob:
if not asset.mesh_node_id:
raise ValueError("missing_node_id")
platform = detect_platform(asset)
if platform == "unknown":
raise ValueError("unknown_platform")
if definition.platform not in {platform, "all"}:
raise ValueError("platform_mismatch")
if not definition.enabled:
raise ValueError("definition_disabled")
package = _seed_generic_job_package(db)
token = secrets.token_urlsafe(32)
job = SoftwareJob(
asset_id=asset.id, package_id=package.id, status="created",
job_type="job_definition", action="run", priority=100, platform=platform,
parameters={
"definition_id": definition.id,
"definition_name": definition.name,
"definition_revision": definition.revision,
"definition_job_kind": definition.job_kind,
"definition_interpreter": definition.interpreter,
"definition_source_type": definition.source_type,
"effective_parameters": _job_effective_parameters(db, definition, asset),
**_job_creation_metadata(
creation_mode=creation_mode,
bulk_batch_id=bulk_batch_id,
bulk_position=bulk_position,
bulk_total=bulk_total,
),
},
result_data={}, attempt_count=0, max_attempts=1,
callback_token_hash=token_hash(token),
callback_expires_at=datetime.utcnow() + timedelta(seconds=max(60, min(int(definition.timeout_seconds or 300) + 120, 86400))),
created_by=_changed_by(request),
)
db.add(job)
db.flush()
sync_asset_job_state(db, job, increment_execution=True)
_record_job_event(
db,
job,
"created",
"created",
_translate_request(
request,
"jobs.event.created_named",
'Job "{name}" created',
name=definition.name,
),
_changed_by(request),
)
db.commit()
db.refresh(job)
configured = str(_software_settings().get("callback_base_url", "") or "").strip().rstrip("/")
callback_base = configured or str(request.base_url).rstrip("/")
job_parameters = dict(job.parameters or {})
job_parameters["_callback_base_url"] = callback_base
job.parameters = job_parameters
db.commit()
threading.Thread(target=execute_job, args=(job.id, token, callback_base), daemon=True).start()
return job
@app.post("/assets/{asset_id}/jobs/run")
def asset_run_job_definition(asset_id: int, request: Request, definition_id: int = Form(...), db: Session = Depends(get_db)):
_require_admin(request)
asset = _get_visible_asset(db, request, asset_id)
definition = db.get(JobDefinition, definition_id)
if not definition or not definition.enabled:
return RedirectResponse(f"/assets/{asset.id}?toast_error=" + quote(_translate_request(request, "jobs.definition_unavailable", "The selected job definition is not available.")), status_code=303)
try:
job = _create_definition_job(db, request, asset, definition)
except ValueError as exc:
key = str(exc)
messages = {
"missing_node_id": _translate_request(request, "jobs.asset_missing_node", "The asset has no MeshCentral node ID."),
"unknown_platform": _translate_request(request, "jobs.asset_unknown_platform", "The operating system platform could not be detected."),
"platform_mismatch": _translate_request(request, "jobs.platform_mismatch", "The job definition does not match the asset platform."),
"definition_disabled": _translate_request(request, "jobs.definition_unavailable", "The selected job definition is not available."),
}
return RedirectResponse(f"/assets/{asset.id}?toast_error=" + quote(messages.get(key, key)), status_code=303)
return RedirectResponse(f"/software/jobs/{job.id}?toast_success=" + quote(_translate_request(request, "jobs.started", "Job started.")), status_code=303)
def _create_software_inventory_job(
db: Session,
request: Request,
asset: Asset,
*,
creation_mode: str = "single",
bulk_batch_id: str | None = None,
bulk_position: int | None = None,
bulk_total: int | None = None,
) -> SoftwareJob:
if not asset.mesh_node_id:
raise ValueError("missing_node_id")
platform = detect_platform(asset)
if platform == "android":
raise ValueError("unsupported_android")
if platform == "unknown":
raise ValueError("unknown_platform")
package = _seed_software_inventory_package(db)
token = secrets.token_urlsafe(32)
job = SoftwareJob(
asset_id=asset.id, package_id=package.id, status="created",
job_type="software_inventory", action="run", priority=100, platform=platform,
parameters=_job_creation_metadata(
creation_mode=creation_mode,
bulk_batch_id=bulk_batch_id,
bulk_position=bulk_position,
bulk_total=bulk_total,
),
result_data={}, attempt_count=0, max_attempts=1,
callback_token_hash=token_hash(token),
callback_expires_at=datetime.utcnow() + timedelta(
minutes=max(5, min(int(_software_settings().get("callback_timeout_minutes", package.callback_timeout_minutes or 30) or 30), 240))
),
created_by=_changed_by(request),
)
db.add(job)
db.flush()
sync_asset_job_state(db, job, increment_execution=True)
_record_job_event(db, job, "created", "created", _translate_request(request, "jobs.event.created", "Job created"), _changed_by(request))
db.commit()
db.refresh(job)
configured = str(_software_settings().get("callback_base_url", "") or "").strip().rstrip("/")
callback_base = configured or str(request.base_url).rstrip("/")
job_parameters = dict(job.parameters or {})
job_parameters["_callback_base_url"] = callback_base
job.parameters = job_parameters
db.commit()
threading.Thread(target=execute_job, args=(job.id, token, callback_base), daemon=True).start()
return job
@app.post("/assets/{asset_id}/software-inventory")
def asset_software_inventory(
asset_id: int,
request: Request,
return_to: str = Form("software-page"),
db: Session = Depends(get_db),
):
_require_admin(request)
asset = _get_visible_asset(db, request, asset_id)
try:
_create_software_inventory_job(db, request, asset)
except ValueError as exc:
messages = {
"missing_node_id": "Das Asset besitzt keine MeshCentral Node-ID.",
"unsupported_android": "Der MeshCentral-Android-Agent stellt derzeit keinen vergleichbaren RunCommand-Kanal für diese Inventur bereit.",
"unknown_platform": "Das Betriebssystem konnte nicht sicher erkannt werden.",
}
raise HTTPException(400, messages.get(str(exc), "Die Softwareinventur konnte nicht gestartet werden."))
target = (
f"/assets/{asset.id}?inventory_tab=software"
if return_to == "asset-detail"
else f"/assets/{asset.id}/software"
)
return RedirectResponse(
target + f"&toast_success={quote(_translate_request(request, 'software.inventory_queued', 'Software inventory was submitted to MeshCentral.'))}"
if "?" in target
else target + f"?toast_success={quote(_translate_request(request, 'software.inventory_queued', 'Software inventory was submitted to MeshCentral.'))}",
status_code=303,
)
@app.post("/jobs/bulk-create")
async def bulk_create_jobs(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
job_type = str(form.get("job_type") or "").strip()
definition_id_raw = str(form.get("definition_id") or "").strip()
definition = None
if definition_id_raw:
try:
definition = db.get(JobDefinition, int(definition_id_raw))
except (TypeError, ValueError):
definition = None
if not definition or not definition.enabled:
return RedirectResponse("/assets?toast_error=" + quote(_translate_request(request, "jobs.definition_unavailable", "The selected job definition is not available.")), status_code=303)
elif job_type != "software_inventory":
return RedirectResponse("/assets?toast_error=" + quote(_translate_request(request, "jobs.bulk.invalid_type", "The selected job type is not available.")), status_code=303)
asset_ids = []
for value in form.getlist("asset_ids"):
try: asset_ids.append(int(value))
except (TypeError, ValueError): continue
asset_ids = list(dict.fromkeys(asset_ids))
if not asset_ids:
return RedirectResponse("/assets?toast_error=" + quote(_translate_request(request, "jobs.bulk.none_selected", "No assets were selected.")), status_code=303)
visible_assets = _apply_asset_access(db.query(Asset), request).filter(Asset.id.in_(asset_ids)).all()
visible_by_id = {asset.id: asset for asset in visible_assets}
created = 0
skipped = 0
bulk_batch_id = uuid.uuid4().hex
bulk_total = len(asset_ids)
for bulk_position, asset_id in enumerate(asset_ids, start=1):
asset = visible_by_id.get(asset_id)
if not asset:
skipped += 1
continue
try:
if definition is not None:
_create_definition_job(
db,
request,
asset,
definition,
creation_mode="bulk",
bulk_batch_id=bulk_batch_id,
bulk_position=bulk_position,
bulk_total=bulk_total,
)
else:
_create_software_inventory_job(
db,
request,
asset,
creation_mode="bulk",
bulk_batch_id=bulk_batch_id,
bulk_position=bulk_position,
bulk_total=bulk_total,
)
created += 1
except ValueError:
skipped += 1
if created:
message = _translate_request(request, "jobs.bulk.result", "Created {created} jobs; skipped {skipped} assets.").format(created=created, skipped=skipped)
return RedirectResponse("/assets?toast_success=" + quote(message), status_code=303)
message = _translate_request(request, "jobs.bulk.none_created", "No job could be created; skipped {skipped} assets.").format(skipped=skipped)
return RedirectResponse("/assets?toast_error=" + quote(message), status_code=303)
@app.post("/assets/job-filter/start")
async def start_job_for_filtered_assets(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
category_id = _optional_form_int(form.get("category_id"), field_name="category_id", minimum=1)
job_filter = _normalize_asset_job_filter(
db,
selection=form.getlist("job_filter"),
excluded_statuses=form.getlist("job_filter_exclude"),
revision_mode=str(form.get("job_filter_revision") or "all"),
)
if not job_filter:
return RedirectResponse(
"/assets?toast_error=" + quote(_translate_request(request, "assets.job_filter.invalid", "The job filter is invalid.")),
status_code=303,
)
start_selection = str(form.get("job_filter_start") or "").strip()
if not start_selection:
start_selection = job_filter["selections"][0]
if start_selection not in job_filter["selections"]:
return RedirectResponse(
_asset_job_filter_url(
job_filter,
category_id,
toast_error=_translate_request(request, "assets.job_filter.start_invalid", "Select one of the filtered jobs to run."),
),
status_code=303,
)
start_spec = next((spec for spec in job_filter["specs"] if spec["selection"] == start_selection), None)
if start_spec is None:
return RedirectResponse(
_asset_job_filter_url(
job_filter,
category_id,
toast_error=_translate_request(request, "assets.job_filter.start_invalid", "Select one of the filtered jobs to run."),
),
status_code=303,
)
definition = start_spec.get("definition")
if definition is not None and not definition.enabled:
return RedirectResponse(
_asset_job_filter_url(
job_filter,
category_id,
toast_error=_translate_request(request, "jobs.definition_unavailable", "The selected job definition is not available."),
),
status_code=303,
)
query = _apply_asset_access(db.query(Asset), request)
if category_id:
query = query.filter(Asset.category_id == category_id)
candidate_assets = query.order_by(func.lower(Asset.name), Asset.id).all()
assets, _views = _filter_assets_by_job_states(db, candidate_assets, job_filter)
if not assets:
return RedirectResponse(
_asset_job_filter_url(
job_filter,
category_id,
toast_error=_translate_request(request, "assets.job_filter.no_assets", "No assets match this job filter."),
),
status_code=303,
)
created = 0
skipped = 0
bulk_batch_id = uuid.uuid4().hex
bulk_total = len(assets)
for bulk_position, asset in enumerate(assets, start=1):
try:
if definition is not None:
_create_definition_job(
db,
request,
asset,
definition,
creation_mode="bulk",
bulk_batch_id=bulk_batch_id,
bulk_position=bulk_position,
bulk_total=bulk_total,
)
else:
_create_software_inventory_job(
db,
request,
asset,
creation_mode="bulk",
bulk_batch_id=bulk_batch_id,
bulk_position=bulk_position,
bulk_total=bulk_total,
)
created += 1
except ValueError:
skipped += 1
if created:
message = _translate_request(
request,
"assets.job_filter.started",
"Created {created} jobs for the filtered assets; skipped {skipped} assets.",
created=created,
skipped=skipped,
)
return RedirectResponse(
_asset_job_filter_url(job_filter, category_id, toast_success=message),
status_code=303,
)
message = _translate_request(
request,
"jobs.bulk.none_created",
"No job could be created; skipped {skipped} assets.",
skipped=skipped,
)
return RedirectResponse(
_asset_job_filter_url(job_filter, category_id, toast_error=message),
status_code=303,
)
def _validate_software_callback_token(job_id: int, token: str) -> bool:
"""Validate a callback outside the ASGI event loop.
Returns True when the same valid token belongs to a job that was already
completed. Callback delivery is intentionally idempotent so a client may
safely retry after a lost HTTP response.
"""
db = SessionLocal()
try:
job = db.get(SoftwareJob, job_id)
if not job:
raise HTTPException(404, "Job nicht gefunden")
if not token or not hmac.compare_digest(token_hash(token), job.callback_token_hash):
raise HTTPException(401, "Ungültiger Callback-Token")
if datetime.utcnow() > job.callback_expires_at:
raise HTTPException(410, "Callback-Token ist abgelaufen")
return bool(job.callback_completed_at)
finally:
db.close()
def _decode_software_callback_body(raw_body: bytes) -> tuple[str, str, list[str]]:
decoded_body = ""
detected_encoding = "utf-8"
decode_errors: list[str] = []
candidates: list[tuple[str, bytes]] = []
if raw_body.startswith(b"\xef\xbb\xbf"):
candidates.append(("utf-8-sig", raw_body))
elif raw_body.startswith(b"\xff\xfe"):
candidates.append(("utf-16-le", raw_body[2:]))
elif raw_body.startswith(b"\xfe\xff"):
candidates.append(("utf-16-be", raw_body[2:]))
candidates.extend([
("utf-8", raw_body),
("utf-8-sig", raw_body),
("utf-16-le", raw_body),
("utf-16-be", raw_body),
])
seen_encodings: set[str] = set()
for encoding, value in candidates:
if encoding in seen_encodings:
continue
seen_encodings.add(encoding)
try:
candidate = value.decode(encoding)
if candidate.lstrip().startswith(("{", "[")):
decoded_body = candidate
detected_encoding = encoding
break
except UnicodeDecodeError as exc:
decode_errors.append(f"{encoding}: {exc}")
if not decoded_body:
decoded_body = raw_body.decode("utf-8", errors="replace")
detected_encoding = "utf-8-replacement"
return decoded_body, detected_encoding, decode_errors
def _process_software_job_callback(
job_id: int,
token: str,
raw_body: bytes,
content_type: str,
client_host: str,
queued_at: float,
) -> dict[str, Any]:
"""Process one callback in the dedicated callback worker pool.
Every worker creates its own SQLAlchemy session. The job row is locked so
duplicate callbacks cannot import the same result concurrently. Software
inventory replacement additionally locks the asset row, serializing only
callbacks that target the same asset while callbacks for other assets run
in parallel.
"""
started_at = time.monotonic()
_software_debug_log(
f"CALLBACK worker started | job={job_id} | client={client_host} | "
f"queue_wait_ms={(started_at - queued_at) * 1000:.1f} | workers={SOFTWARE_CALLBACK_WORKER_COUNT}"
)
db = SessionLocal()
try:
job = (
db.query(SoftwareJob)
.filter(SoftwareJob.id == job_id)
.with_for_update()
.first()
)
if not job:
raise HTTPException(404, "Job nicht gefunden")
if not token or not hmac.compare_digest(token_hash(token), job.callback_token_hash):
raise HTTPException(401, "Ungültiger Callback-Token")
if datetime.utcnow() > job.callback_expires_at:
raise HTTPException(410, "Callback-Token ist abgelaufen")
if job.callback_completed_at:
_software_debug_log(f"CALLBACK duplicate accepted | job={job.id} | status={job.status}")
db.rollback()
return {"ok": True, "job_id": job.id, "status": job.status, "duplicate": True}
decoded_body, detected_encoding, decode_errors = _decode_software_callback_body(raw_body)
content_length = len(raw_body)
_software_debug_log(
f"CALLBACK raw body | job={job_id} | encoding={detected_encoding} | BEGIN\n"
+ decoded_body
+ f"\nCALLBACK raw body | job={job_id} | END"
)
try:
payload = json.loads(decoded_body)
except Exception as exc:
failure_log = (
f"JSON parsing failed: {type(exc).__name__}: {exc}\n"
f"Detected encoding: {detected_encoding}\n"
f"Content-Type: {content_type}\n"
f"Body bytes: {content_length}\n"
f"\n--- Complete callback body ---\n{decoded_body}"
)
job.log_text = ((job.log_text or "") + "\n\n--- Callback parsing error ---\n" + failure_log)[-2_000_000:]
job.message = "Callback empfangen, aber JSON konnte nicht verarbeitet werden."
db.commit()
_software_debug_log(
f"CALLBACK JSON ERROR | job={job_id} | {type(exc).__name__}: {exc}"
+ (f" | decode_errors={'; '.join(decode_errors)}" if decode_errors else "")
)
raise HTTPException(400, "Ungültiges JSON") from exc
if not isinstance(payload, dict):
job.log_text = (
(job.log_text or "")
+ f"\n\n--- Callback parsing error ---\nJSON root type: {type(payload).__name__}\n"
+ f"\n--- Complete callback body ---\n{decoded_body}"
)[-2_000_000:]
job.message = "Callback empfangen, aber kein JSON-Objekt erhalten."
db.commit()
_software_debug_log(f"CALLBACK JSON TYPE ERROR | job={job_id} | root={type(payload).__name__}")
raise HTTPException(400, "JSON-Objekt erwartet")
_software_debug_log(
f"CALLBACK JSON parsed | job={job_id} | keys={','.join(sorted(str(key) for key in payload.keys()))}"
)
log_parts: list[str] = []
for label, key in (("Log", "log"), ("STDOUT", "stdout"), ("STDERR", "stderr")):
value = payload.get(key)
if value is None or value == "":
continue
text_value = str(value)
if len(text_value) > 200000:
raise HTTPException(413, f"{label} ist zu groß")
log_parts.append(f"--- {label} ---\n{text_value}")
status = str(payload.get("status") or "running").lower()
job.status = status if status in {"running", "success", "failed", "partial"} else "running"
if job.started_at is None:
job.started_at = datetime.utcnow()
try:
job.exit_code = int(payload.get("exit_code")) if payload.get("exit_code") is not None else job.exit_code
except (TypeError, ValueError):
pass
if payload.get("message") is not None:
job.message = str(payload.get("message"))[:4000]
summary_lines = [
f"Callback received: {datetime.utcnow().isoformat()}Z",
f"Status: {job.status}",
f"Exit code: {job.exit_code if job.exit_code is not None else '-'}",
]
if job.message:
summary_lines.append(f"Message: {job.message}")
callback_log = "\n".join(summary_lines + (["", *log_parts] if log_parts else []))
job.log_text = ((job.log_text or "") + "\n\n--- Client callback ---\n" + callback_log).strip()[-2_000_000:]
job.result_data = {
str(key): value
for key, value in payload.items()
if key not in {"software", "log"}
}
software = payload.get("software")
stored_software_count = 0
excluded_software_count = 0
if isinstance(software, list):
_software_debug_log(f"CALLBACK software processing started | job={job.id} | received={len(software)}")
# Serialize inventory replacement only per asset. Other assets are
# still processed concurrently by the callback worker pool.
db.query(Asset.id).filter(Asset.id == job.asset_id).with_for_update().first()
now = datetime.utcnow()
removed_count = (
db.query(InstalledSoftware)
.filter(InstalledSoftware.asset_id == job.asset_id)
.delete(synchronize_session=False)
)
_software_debug_log(f"CALLBACK previous software removed | job={job.id} | rows={removed_count}")
seen: set[tuple[str, str, str, str]] = set()
inventory_rows: list[dict[str, Any]] = []
inventory_platform = (job.platform or "unknown")[:40].lower()
exclusion_rules = _software_settings().get("inventory_exclusion_rules", [])
for item in software[:50000]:
if not isinstance(item, dict):
continue
name = str(item.get("name") or "").strip()
if not name:
continue
if _software_name_is_excluded(name, inventory_platform, exclusion_rules):
excluded_software_count += 1
continue
values = (
name[:300],
str(item.get("version") or "")[:160],
str(item.get("publisher") or "")[:220],
str(item.get("architecture") or "")[:40],
)
identity = tuple(value.casefold() for value in values)
if identity in seen:
continue
seen.add(identity)
inventory_rows.append({
"asset_id": job.asset_id,
"name": values[0],
"version": values[1],
"publisher": values[2],
"architecture": values[3],
"platform": inventory_platform,
"install_date": str(item.get("install_date") or "")[:40] or None,
"source": "meshcentral-inventory",
"first_seen_at": now,
"last_seen_at": now,
})
if inventory_rows:
db.bulk_insert_mappings(InstalledSoftware, inventory_rows)
stored_software_count = len(inventory_rows)
_software_debug_log(
f"CALLBACK software filtered | job={job.id} | received={len(software)} | "
f"excluded={excluded_software_count} | stored={stored_software_count} | "
f"rules={len(exclusion_rules)}"
)
if isinstance(job.result_data, dict):
job.result_data["software_inventory"] = {
"received": len(software),
"excluded": excluded_software_count,
"stored": stored_software_count,
}
if bool(payload.get("completed")) or job.status in {"success", "failed"}:
job.callback_completed_at = datetime.utcnow()
job.finished_at = datetime.utcnow()
language_code = str(load_config().get("general", {}).get("default_language", "de") or "de")
event_message = job.message or translate(
db,
language_code,
"jobs.event.callback",
"Callback received",
)
_record_job_event(db, job, "callback", job.status, event_message)
sync_asset_job_state(db, job)
db.commit()
processing_ms = (time.monotonic() - started_at) * 1000
_software_debug_log(
f"CALLBACK database commit successful | job={job.id} | stored_software={stored_software_count} | excluded_software={excluded_software_count} | processing_ms={processing_ms:.1f}"
)
_software_debug_log(
f"CALLBACK accepted | job={job.id} | status={job.status} | completed={bool(job.callback_completed_at)} | "
f"software_entries={len(software) if isinstance(software, list) else 0}"
)
return {"ok": True, "job_id": job.id, "status": job.status}
except HTTPException:
db.rollback()
raise
except Exception:
db.rollback()
logger.exception("Software callback worker failed for job %s", job_id)
raise
finally:
db.close()
@app.post("/api/software-jobs/{job_id}/callback")
async def software_job_callback(job_id: int, request: Request, token: str = ""):
client_host = request.client.host if request.client else "-"
content_type = request.headers.get("content-type", "-")
_software_debug_log(
f"CALLBACK received | job={job_id} | client={client_host} | content_type={content_type} | "
f"workers={SOFTWARE_CALLBACK_WORKER_COUNT}"
)
# The short authentication query runs outside the ASGI event loop. It is
# repeated under a row lock inside the worker to close race conditions.
already_completed = await asyncio.to_thread(_validate_software_callback_token, job_id, token)
if already_completed:
_software_debug_log(f"CALLBACK duplicate accepted before body read | job={job_id} | client={client_host}")
return {"ok": True, "job_id": job_id, "duplicate": True}
try:
raw_body = await request.body()
except ClientDisconnect:
_software_debug_log(
f"CALLBACK client disconnected while uploading body | job={job_id} | client={client_host}"
)
return Response(status_code=499)
_software_debug_log(f"CALLBACK body received | job={job_id} | bytes={len(raw_body)}")
queued_at = time.monotonic()
loop = asyncio.get_running_loop()
return await loop.run_in_executor(
SOFTWARE_CALLBACK_EXECUTOR,
_process_software_job_callback,
job_id,
token,
raw_body,
content_type,
client_host,
queued_at,
)
@app.get("/api/software-jobs/{job_id}")
def software_job_status(job_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
_expire_stale_software_jobs(db, job_id)
job = db.get(SoftwareJob, job_id)
if not job:
raise HTTPException(404, "Job not found")
return {
"id": job.id,
"job_type": job.job_type,
"action": job.action,
"priority": job.priority,
"status": job.status,
"status_label": _translate_request(request, f"software.status.{job.status}", job.status),
"message": job.message or "",
"exit_code": job.exit_code,
"attempt_count": job.attempt_count,
"max_attempts": job.max_attempts,
"finished_at": _utc_iso(job.finished_at),
"log_text": job.log_text or "",
}
def _prepare_existing_job_restart(
db: Session,
request: Request | None,
job: SoftwareJob,
*,
automatic: bool = False,
) -> str:
"""Reset an existing job for another attempt and return its callback token.
Manual and automatic retries use this same tested reset path. The caller
commits the transaction and starts ``execute_job`` afterwards.
"""
token = secrets.token_urlsafe(32)
package_timeout = getattr(job.package, "callback_timeout_minutes", None) or 30
timeout_minutes = max(
5,
min(
int(_software_settings().get("callback_timeout_minutes", package_timeout) or 30),
240,
),
)
job.callback_token_hash = token_hash(token)
job.callback_expires_at = datetime.utcnow() + timedelta(minutes=timeout_minutes)
job.callback_completed_at = None
job.status = "created"
job.sent_at = None
job.started_at = None
job.finished_at = None
job.exit_code = None
job.message = None
job.result_data = {}
job.meshctrl_stdout = None
job.meshctrl_stderr = None
job.command_preview = None
job.log_text = ""
job.scheduled_at = None
job.max_attempts = max(int(job.max_attempts or 1), int(job.attempt_count or 0) + 1)
if automatic:
language_code = str(load_config().get("general", {}).get("default_language", "de") or "de")
event_type = "automatic_retry_requested"
event_message = translate(
db,
language_code,
"jobs.event.automatic_retry_requested",
"Automatic retry requested",
)
created_by = "system:auto-retry"
else:
event_type = "restart_requested"
event_message = _translate_request(request, "jobs.event.restart_requested", "Manual restart requested")
created_by = _changed_by(request)
_record_job_event(
db,
job,
event_type,
"created",
event_message,
created_by,
)
sync_asset_job_state(db, job, increment_execution=True)
return token
@app.post("/api/software-jobs/statuses")
def software_job_statuses(
request: Request,
job_ids: list[int] = Body(...),
db: Session = Depends(get_db),
):
"""Return compact live status data for the visible software-job rows."""
_require_admin(request)
normalized_ids: list[int] = []
seen_ids: set[int] = set()
for value in job_ids[:1000]:
try:
job_id = int(value)
except (TypeError, ValueError):
continue
if job_id <= 0 or job_id in seen_ids:
continue
seen_ids.add(job_id)
normalized_ids.append(job_id)
if not normalized_ids:
return {"jobs": []}
_expire_stale_software_jobs(db)
rows = (
db.query(SoftwareJob.id, SoftwareJob.status)
.filter(SoftwareJob.id.in_(normalized_ids))
.all()
)
user = _session_user(request) or {}
language_code = (
user.get("language_code")
or request.session.get("language_code")
or load_config().get("general", {}).get("default_language", "en")
)
status_labels: dict[str, str] = {}
result = []
for job_id, status in rows:
normalized_status = str(status or "created")
if normalized_status not in status_labels:
status_labels[normalized_status] = translate(
db,
language_code,
f"software.status.{normalized_status}",
normalized_status,
)
result.append({
"id": int(job_id),
"status": normalized_status,
"status_label": status_labels[normalized_status],
})
return {"jobs": result}
@app.post("/jobs/{job_id}/restart")
def restart_job(job_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
job = db.query(SoftwareJob).options(joinedload(SoftwareJob.asset), joinedload(SoftwareJob.package)).filter(SoftwareJob.id == job_id).first()
if not job:
raise HTTPException(404, _translate_request(request, "jobs.not_found", "Job not found"))
_expire_stale_software_jobs(db, job.id)
db.refresh(job)
if job.status in {"sending", "running"}:
return RedirectResponse(
f"/jobs/{job.id}?toast_error=" + quote(_translate_request(request, "jobs.restart.active", "An active job cannot be restarted.")),
status_code=303,
)
if not job.asset or not job.asset.mesh_node_id:
return RedirectResponse(
f"/jobs/{job.id}?toast_error=" + quote(_translate_request(request, "jobs.restart.no_node", "The asset has no MeshCentral node ID.")),
status_code=303,
)
token = _prepare_existing_job_restart(db, request, job)
db.commit()
configured = str(_software_settings().get("callback_base_url", "") or "").strip().rstrip("/")
callback_base = configured or str(request.base_url).rstrip("/")
threading.Thread(target=execute_job, args=(job.id, token, callback_base), daemon=True).start()
return RedirectResponse(
f"/jobs/{job.id}?toast_success=" + quote(_translate_request(request, "jobs.restart.started", "Job was queued again.")),
status_code=303,
)
@app.post("/jobs/restart-selected")
async def restart_selected_jobs(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
job_ids: list[int] = []
for value in form.getlist("job_ids"):
try:
job_ids.append(int(value))
except (TypeError, ValueError):
continue
job_ids = list(dict.fromkeys(job_ids))
if not job_ids:
return RedirectResponse(
"/software?toast_error=" + quote(_translate_request(request, "jobs.bulk_restart.none_selected", "No jobs were selected.")),
status_code=303,
)
_expire_stale_software_jobs(db)
jobs = (
db.query(SoftwareJob)
.options(joinedload(SoftwareJob.asset), joinedload(SoftwareJob.package))
.filter(SoftwareJob.id.in_(job_ids))
.all()
)
by_id = {job.id: job for job in jobs}
queued: list[tuple[int, str]] = []
skipped = 0
for job_id in job_ids:
job = by_id.get(job_id)
if (
job is None
or job.status in {"success", "sending", "running"}
or not job.asset
or not job.asset.mesh_node_id
):
skipped += 1
continue
token = _prepare_existing_job_restart(db, request, job)
queued.append((job.id, token))
if not queued:
db.rollback()
return RedirectResponse(
"/software?toast_error=" + quote(_translate_request(request, "jobs.bulk_restart.none_eligible", "None of the selected jobs can be restarted.")),
status_code=303,
)
db.commit()
configured = str(_software_settings().get("callback_base_url", "") or "").strip().rstrip("/")
callback_base = configured or str(request.base_url).rstrip("/")
for queued_job_id, token in queued:
threading.Thread(target=execute_job, args=(queued_job_id, token, callback_base), daemon=True).start()
message = _translate_request(
request,
"jobs.bulk_restart.started",
"{count} jobs were queued again; {skipped} were skipped.",
count=len(queued),
skipped=skipped,
)
return RedirectResponse(
"/software?toast_success=" + quote(message),
status_code=303,
)
@app.get("/jobs/{job_id}")
def job_detail(job_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
job = db.query(SoftwareJob).options(joinedload(SoftwareJob.asset), joinedload(SoftwareJob.package)).filter(SoftwareJob.id == job_id).first()
if not job:
raise HTTPException(404, _translate_request(request, "jobs.not_found", "Job not found"))
_expire_stale_software_jobs(db, job.id)
db.refresh(job)
_decorate_jobs_for_display(db, [job])
events = db.query(JobEvent).filter(JobEvent.job_id == job.id).order_by(JobEvent.created_at.desc(), JobEvent.id.desc()).all()
can_restart = job.status not in {"sending", "running"}
return templates.TemplateResponse("software_job.html", {"request": request, "job": job, "events": events, "can_restart": can_restart})
@app.get("/software/jobs/{job_id}")
def software_job_detail_legacy(job_id: int, request: Request):
_require_admin(request)
return RedirectResponse(f"/jobs/{job_id}", status_code=307)
@app.get("/login")
def login_page(request: Request):
if _session_user(request):
return RedirectResponse("/profile", status_code=303)
mode = load_config().get("authentication", {}).get("mode", "none")
return templates.TemplateResponse("login.html", {"request": request, "mode": mode})
@app.post("/login")
def login_submit(request: Request, username: str = Form(...), password: str = Form(...), db: Session = Depends(get_db)):
username = username.strip()
auth_config = load_config().get("authentication", {})
mode = str(auth_config.get("mode", "none")).lower()
user = db.query(User).filter(func.lower(User.username) == username.casefold()).first()
valid = False
# Ein geschütztes lokales Konto wird auch im LDAP-Modus ausschließlich
# gegen sein lokales Passwort geprüft. Bei falschem Passwort gibt es
# bewusst keinen LDAP-Fallback für denselben Benutzernamen.
protected_local = bool(
user
and getattr(user, "is_protected", False)
and user.auth_source == "local"
)
if protected_local and mode in {"local", "ldap"}:
valid = bool(user.is_active and _verify_password(password, user.password_hash))
elif mode == "local":
valid = bool(user and user.is_active and user.auth_source == "local" and _verify_password(password, user.password_hash))
elif mode == "ldap":
ldap_user = _ldap_authenticate(username, password, auth_config)
if ldap_user:
if not user:
user = User(username=username, auth_source="ldap", is_active=True)
db.add(user)
if user.is_active:
user.auth_source = "ldap"
user.display_name = ldap_user.get("display_name") or username
user.email = ldap_user.get("email") or None
valid = True
elif mode == "none":
return RedirectResponse("/?toast_warning=Die Anmeldung ist derzeit deaktiviert", status_code=303)
if not valid or not user:
logger.warning("Fehlgeschlagene Anmeldung für Benutzer %s im Modus %s", username, mode)
return RedirectResponse("/login?toast_error=Benutzername oder Passwort ist ungültig", status_code=303)
user.last_login = datetime.utcnow()
db.commit()
db.refresh(user)
_login_user(request, user)
return RedirectResponse("/?toast_success=Anmeldung erfolgreich", status_code=303)
@app.get("/logout")
def logout(request: Request):
request.session.clear()
return RedirectResponse("/?toast_success=Sie wurden abgemeldet", status_code=303)
@app.get("/profile")
def profile(request: Request, db: Session = Depends(get_db)):
session_user = _session_user(request)
if not session_user:
return RedirectResponse("/login?toast_warning=Bitte zuerst anmelden", status_code=303)
user = db.query(User).filter(User.id == session_user.get("id")).first()
if not user:
request.session.clear()
return RedirectResponse("/login?toast_error=Benutzerkonto wurde nicht gefunden", status_code=303)
return templates.TemplateResponse("profile.html", {"request": request, "user": user, "languages": i18n_languages(db)})
@app.post("/profile/language")
def profile_language(request: Request, language_code: str = Form(...), next_url: str = Form("/profile"), db: Session = Depends(get_db)):
available = {item.code for item in i18n_languages(db)}
if language_code not in available:
raise HTTPException(400, "Unsupported language")
session_user = _session_user(request)
if session_user:
user = db.query(User).filter(User.id == session_user.get("id")).first()
if user:
user.language_code = language_code
db.commit()
_login_user(request, user)
request.session["language_code"] = language_code
safe_next = next_url if next_url.startswith("/") and not next_url.startswith("//") else "/"
return RedirectResponse(f"{safe_next}?toast_success=Language changed", status_code=303)
@app.post("/profile/browser-state")
def profile_browser_state(
request: Request,
persist_browser_state: str | None = Form(None),
db: Session = Depends(get_db),
):
session_user = _session_user(request)
if not session_user:
return RedirectResponse("/login?toast_warning=Bitte zuerst anmelden", status_code=303)
user = db.query(User).filter(User.id == session_user.get("id")).first()
if not user:
request.session.clear()
return RedirectResponse("/login?toast_error=Benutzerkonto wurde nicht gefunden", status_code=303)
user.persist_browser_state = persist_browser_state == "on"
db.commit()
_login_user(request, user)
message = "Browserspeicherung aktiviert" if user.persist_browser_state else "Browserspeicherung deaktiviert"
return RedirectResponse(f"/profile?toast_success={quote(message)}", status_code=303)
@app.post("/profile/toast-settings")
def profile_toast_settings(
request: Request,
toast_position: str = Form("top-right"),
toast_duration_seconds: int = Form(6),
db: Session = Depends(get_db),
):
session_user = _session_user(request)
if not session_user:
return RedirectResponse("/login?toast_warning=Bitte zuerst anmelden", status_code=303)
user = db.query(User).filter(User.id == session_user.get("id")).first()
if not user:
request.session.clear()
return RedirectResponse("/login?toast_error=Benutzerkonto wurde nicht gefunden", status_code=303)
allowed_positions = {"top-left", "top-center", "top-right", "bottom-left", "bottom-center", "bottom-right"}
user.toast_position = toast_position if toast_position in allowed_positions else "top-right"
user.toast_duration_seconds = max(1, min(int(toast_duration_seconds or 6), 120))
db.commit()
_login_user(request, user)
return RedirectResponse("/profile?toast_success=Toast-Einstellungen gespeichert", status_code=303)
@app.post("/profile/theme")
def profile_theme(
request: Request,
theme_mode: str = Form("light"),
next_url: str = Form("/profile"),
db: Session = Depends(get_db),
):
session_user = _session_user(request)
if not session_user:
return RedirectResponse("/login?toast_warning=Bitte zuerst anmelden", status_code=303)
user = db.query(User).filter(User.id == session_user.get("id")).first()
if not user:
request.session.clear()
return RedirectResponse("/login?toast_error=Benutzerkonto wurde nicht gefunden", status_code=303)
user.theme_mode = theme_mode if theme_mode in {"light", "dark"} else "light"
db.commit()
_login_user(request, user)
safe_next = next_url if next_url.startswith("/") and not next_url.startswith("//") else "/profile"
message = _translate_request(request, "appearance.theme_saved", "Display mode saved")
return RedirectResponse(f"{safe_next}?toast_success={quote(message)}", status_code=303)
@app.post("/profile/password")
def profile_password(request: Request, current_password: str = Form(...), new_password: str = Form(...), db: Session = Depends(get_db)):
session_user = _session_user(request)
if not session_user:
return RedirectResponse("/login", status_code=303)
user = db.query(User).filter(User.id == session_user.get("id")).first()
if not user or user.auth_source != "local":
return RedirectResponse("/profile?toast_warning=Passwortänderungen sind nur für lokale Konten möglich", status_code=303)
if not _verify_password(current_password, user.password_hash):
return RedirectResponse("/profile?toast_error=Das aktuelle Passwort ist falsch", status_code=303)
if len(new_password) < 8:
return RedirectResponse("/profile?toast_error=Das neue Passwort muss mindestens 8 Zeichen lang sein", status_code=303)
user.password_hash = _hash_password(new_password)
db.commit()
return RedirectResponse("/profile?toast_success=Passwort geändert", status_code=303)
@app.get("/users")
def users_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
users = db.query(User).order_by(User.username).all()
mesh_groups = [row[0] for row in db.query(Asset.mesh_group).filter(Asset.mesh_group.isnot(None), Asset.mesh_group != "").distinct().order_by(Asset.mesh_group).all()]
return templates.TemplateResponse("users.html", {"request": request, "users": users, "auth": public_config().get("authentication", {}), "mesh_groups": mesh_groups})
@app.post("/users/new")
def user_create(request: Request, username: str = Form(...), display_name: str = Form(""), email: str = Form(""), password: str = Form(...), is_admin: str | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
username = username.strip()
if not username:
raise HTTPException(400, "Benutzername fehlt")
if len(password) < 8:
raise HTTPException(400, "Das Passwort muss mindestens 8 Zeichen lang sein")
if db.query(User.id).filter(func.lower(User.username) == username.casefold()).first():
raise HTTPException(400, "Der Benutzername ist bereits vorhanden")
user = User(username=username, display_name=display_name.strip() or username, email=email.strip() or None, password_hash=_hash_password(password), auth_source="local", is_active=True, is_admin=is_admin == "on", language_code=load_config().get("general", {}).get("default_language", "en"))
db.add(user)
try:
db.commit()
except IntegrityError as exc:
db.rollback()
raise HTTPException(400, "Der Benutzername ist bereits vorhanden") from exc
return RedirectResponse("/users?toast_success=Benutzer angelegt", status_code=303)
@app.post("/users/{user_id}/update")
def user_update(user_id: int, request: Request, display_name: str = Form(""), email: str = Form(""), profile_department: str = Form(""), profile_location: str = Form(""), is_active: str | None = Form(None), is_admin: str | None = Form(None), new_password: str = Form(""), mesh_groups: list[str] = Form([]), asset_access_scopes: list[str] = Form([]), db: Session = Depends(get_db)):
_require_admin(request)
user = db.query(User).filter(User.id == user_id).first()
if not user:
raise HTTPException(404, "Benutzer nicht gefunden")
user.display_name = display_name.strip() or user.username
user.email = email.strip() or None
if getattr(user, "is_protected", False):
user.is_active = True
user.is_admin = True
user.auth_source = "local"
else:
user.is_active = is_active == "on"
user.is_admin = is_admin == "on"
user.allowed_mesh_groups = sorted({value.strip() for value in mesh_groups if value.strip()})
user.asset_access_scopes = sorted({value for value in asset_access_scopes if value in ASSET_ACCESS_SCOPES})
user.profile_department = profile_department.strip() or None
user.profile_location = profile_location.strip() or None
if new_password:
if user.auth_source != "local":
raise HTTPException(400, "LDAP-Passwörter können hier nicht geändert werden")
if len(new_password) < 8:
raise HTTPException(400, "Das neue Passwort muss mindestens 8 Zeichen lang sein")
user.password_hash = _hash_password(new_password)
db.commit()
return RedirectResponse("/users?toast_success=Benutzer gespeichert", status_code=303)
@app.post("/users/{user_id}/delete")
def user_delete(user_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
user = db.get(User, user_id)
if not user:
raise HTTPException(404, "Benutzer nicht gefunden")
if getattr(user, "is_protected", False):
raise HTTPException(400, _translate_request(
request,
"users.protected_delete_forbidden",
"Protected local administrators cannot be deleted.",
))
session_user = _session_user(request) or {}
if int(session_user.get("id") or 0) == user.id:
raise HTTPException(400, "Der aktuell angemeldete Benutzer kann nicht gelöscht werden.")
if user.is_admin and user.is_active:
other_admins = db.query(User).filter(User.id != user.id, User.is_admin.is_(True), User.is_active.is_(True)).count()
if other_admins == 0:
raise HTTPException(400, "Der letzte aktive Administrator kann nicht gelöscht werden.")
db.delete(user)
db.commit()
return RedirectResponse("/users?toast_success=Benutzer gelöscht", status_code=303)
def _ldap_search_users(query_text: str, auth_config: dict) -> list[dict]:
"""Listet AD-Benutzer über das konfigurierte LDAP-Dienstkonto auf.
Eine leere Suche liefert alle aktiven Benutzerkonten. Bei einem Suchtext
werden Benutzername, Anzeigename und E-Mail durchsucht.
"""
from ldap3 import ALL, Connection, Server, SUBTREE
from ldap3.core.exceptions import LDAPException
from ldap3.utils.conv import escape_filter_chars
ldap = auth_config.get("ldap", {})
server_name = str(ldap.get("server", "")).strip()
use_ssl = bool(ldap.get("use_ssl", False))
start_tls = bool(ldap.get("start_tls", False))
port = int(ldap.get("port") or (636 if use_ssl else 389))
bind_dn = str(ldap.get("bind_dn", "")).strip()
env_name = str(ldap.get("bind_password_env", "LDAP_BIND_PASSWORD")).strip()
bind_password = os.getenv(env_name, "") or str(ldap.get("bind_password", ""))
base_dn = str(ldap.get("user_base_dn", "")).strip()
if not server_name:
raise HTTPException(400, "LDAP-Server ist nicht konfiguriert")
if not base_dn:
raise HTTPException(400, "LDAP Benutzer-Basis-DN ist nicht konfiguriert")
if bind_dn and not bind_password:
raise HTTPException(400, f"LDAP-Bind-Passwort fehlt (Umgebungsvariable {env_name})")
raw_query = query_text.strip()
def ldap_wildcard_pattern(value: str) -> str:
# Benutzer-Wildcards zulassen, sonstige LDAP-Sonderzeichen aber sicher maskieren.
# * = beliebig viele Zeichen.
parts = []
literal = []
for char in value:
if char == "*":
if literal:
parts.append(escape_filter_chars("".join(literal)))
literal = []
parts.append("*")
else:
literal.append(char)
if literal:
parts.append(escape_filter_chars("".join(literal)))
pattern = "".join(parts)
# Ohne explizite Wildcard bleibt die bisherige Teilstringsuche erhalten.
if "*" not in value:
pattern = f"*{pattern}*"
return pattern
base_filter = "(&(objectCategory=person)(objectClass=user)(sAMAccountName=*)(!(userAccountControl:1.2.840.113556.1.4.803:=2)))"
if raw_query:
q = ldap_wildcard_pattern(raw_query)
search_filter = (
"(&(objectCategory=person)(objectClass=user)(sAMAccountName=*)"
"(!(userAccountControl:1.2.840.113556.1.4.803:=2))"
f"(|(sAMAccountName={q})(displayName={q})(mail={q})))"
)
else:
search_filter = base_filter
server = Server(server_name, port=port, use_ssl=use_ssl, get_info=ALL, connect_timeout=10)
conn = Connection(
server,
user=bind_dn or None,
password=bind_password or None,
auto_bind=False,
receive_timeout=30,
raise_exceptions=False,
)
try:
# ldap3.Connection.open() liefert je nach Version keinen verlässlichen
# booleschen Rückgabewert. Daher nur öffnen und anschließend binden.
if start_tls and not use_ssl:
conn.open()
if not conn.start_tls():
raise HTTPException(400, f"LDAP StartTLS fehlgeschlagen: {conn.last_error}")
if not conn.bind():
raise HTTPException(400, f"LDAP Bind fehlgeschlagen: {conn.last_error}")
attributes = ["sAMAccountName", "displayName", "mail", "distinguishedName"]
entries = conn.extend.standard.paged_search(
search_base=base_dn,
search_filter=search_filter,
search_scope=SUBTREE,
attributes=attributes,
paged_size=500,
generator=False,
)
result = []
for entry in entries:
if entry.get("type") != "searchResEntry":
continue
attrs = entry.get("attributes", {})
username = str(attrs.get("sAMAccountName") or "").strip()
if not username:
continue
result.append({
"username": username,
"display_name": str(attrs.get("displayName") or username),
"email": str(attrs.get("mail") or ""),
"dn": str(entry.get("dn") or attrs.get("distinguishedName") or ""),
})
result.sort(key=lambda item: ((item.get("display_name") or item["username"]).casefold(), item["username"].casefold()))
ldap_logger.info("LDAP-Benutzerabfrage erfolgreich: Filter=%r Treffer=%s", search_filter, len(result))
return result
except HTTPException:
raise
except LDAPException as exc:
ldap_logger.exception("LDAP-Ausnahme bei Benutzerabfrage: %s", exc)
raise HTTPException(400, f"LDAP-Abfrage fehlgeschlagen: {exc}") from exc
except Exception as exc:
ldap_logger.exception("Unerwarteter LDAP-Fehler bei Benutzerabfrage: %s", exc)
raise HTTPException(400, f"LDAP-Abfrage fehlgeschlagen: {exc}") from exc
finally:
try:
conn.unbind()
except Exception:
pass
def _mark_existing_ldap_users(results: list[dict], db: Session) -> list[dict]:
usernames = [item["username"].casefold() for item in results]
existing = {}
if usernames:
existing = {
user.username.casefold(): user
for user in db.query(User).filter(func.lower(User.username).in_(usernames)).all()
}
for item in results:
user = existing.get(item["username"].casefold())
item["existing"] = bool(user)
item["existing_source"] = user.auth_source if user else ""
return results
@app.get("/settings/ldap-search")
def ldap_search_form(request: Request, show_all: int = 0, db: Session = Depends(get_db)):
_require_admin(request)
results = []
loaded = False
if show_all:
cfg = load_config()
results = _mark_existing_ldap_users(_ldap_search_users("", cfg.get("authentication", {})), db)
loaded = True
return templates.TemplateResponse("ldap_search.html", {
"request": request, "query": "", "results": results, "loaded": loaded,
})
@app.post("/settings/ldap-search")
def ldap_search_page(request: Request, ldap_query: str = Form(""), db: Session = Depends(get_db)):
_require_admin(request)
cfg = load_config()
results = _mark_existing_ldap_users(_ldap_search_users(ldap_query, cfg.get("authentication", {})), db)
return templates.TemplateResponse("ldap_search.html", {
"request": request, "query": ldap_query, "results": results, "loaded": True,
})
@app.post("/settings/ldap-search/import")
def ldap_users_import(
request: Request,
usernames: list[str] = Form([]),
db: Session = Depends(get_db),
):
_require_admin(request)
selected = {value.strip().casefold() for value in usernames if value.strip()}
if not selected:
return RedirectResponse("/settings/ldap-search?toast_warning=Keine Benutzer ausgewählt", status_code=303)
cfg = load_config()
ldap_users = _ldap_search_users("", cfg.get("authentication", {}))
ldap_by_name = {item["username"].casefold(): item for item in ldap_users}
created = 0
updated = 0
skipped = 0
for normalized in selected:
item = ldap_by_name.get(normalized)
if not item:
skipped += 1
continue
user = db.query(User).filter(func.lower(User.username) == item["username"].casefold()).first()
if user:
# Lokale Konten werden nicht automatisch in LDAP-Konten umgewandelt.
if user.auth_source != "ldap":
skipped += 1
continue
user.display_name = item["display_name"] or item["username"]
user.email = item["email"] or None
updated += 1
else:
db.add(User(
username=item["username"],
display_name=item["display_name"] or item["username"],
email=item["email"] or None,
password_hash=None,
auth_source="ldap",
is_active=True,
is_admin=False,
allowed_mesh_groups=[],
language_code=load_config().get("general", {}).get("default_language", "en"),
))
created += 1
db.commit()
return RedirectResponse(
f"/users?toast_success={created} LDAP-Benutzer angelegt, {updated} aktualisiert, {skipped} übersprungen",
status_code=303,
)
@app.get("/settings/translations")
def translations_page(request: Request, missing_only: int = 0, db: Session = Depends(get_db)):
_require_admin(request)
langs = i18n_languages(db, active_only=False)
rows = {}
for item in db.query(Translation).order_by(Translation.translation_key, Translation.language_code).all():
rows.setdefault(item.translation_key, {})[item.language_code] = item.text_value
if missing_only:
rows = {key: values for key, values in rows.items() if any(not values.get(lang.code) for lang in langs)}
return templates.TemplateResponse("translations.html", {"request": request, "languages": langs, "rows": rows, "missing_only": bool(missing_only)})
@app.post("/settings/translations/save")
async def translations_save(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
langs = {lang.code for lang in i18n_languages(db, active_only=False)}
for field, value in form.multi_items():
if not field.startswith("tr__"):
continue
_, code, key = field.split("__", 2)
if code not in langs:
continue
row = db.query(Translation).filter(Translation.translation_key == key, Translation.language_code == code).first()
if row:
row.text_value = str(value)
elif str(value).strip():
db.add(Translation(translation_key=key, language_code=code, text_value=str(value)))
db.commit(); clear_translation_cache()
return RedirectResponse("/settings/translations?toast_success=Translations updated", status_code=303)
@app.get("/settings/translations/export")
def translations_export(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
langs = i18n_languages(db, active_only=False)
data = {}
for item in db.query(Translation).all():
data.setdefault(item.translation_key, {})[item.language_code] = item.text_value
headers = ["Key"] + [lang.code for lang in langs]
rows = [[key] + [values.get(lang.code, "") for lang in langs] for key, values in sorted(data.items())]
return _xlsx_response("assetmanager-translations.xlsx", headers, rows)
@app.post("/settings/translations/import")
async def translations_import(request: Request, file: UploadFile = File(...), ignore_empty: str | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
payload = await file.read()
workbook = load_workbook(io.BytesIO(payload), data_only=True)
sheet = workbook.active
headers = [str(cell.value or "").strip() for cell in sheet[1]]
if not headers or headers[0].casefold() != "key":
raise HTTPException(400, "The first Excel column must be Key")
known = {lang.code for lang in i18n_languages(db, active_only=False)}
columns = [(idx, code) for idx, code in enumerate(headers[1:], start=2) if code in known]
changed = 0
for values in sheet.iter_rows(min_row=2, values_only=True):
key = str(values[0] or "").strip()
if not key:
continue
for idx, code in columns:
value = values[idx-1]
text_value = "" if value is None else str(value)
if not text_value and ignore_empty == "on":
continue
row = db.query(Translation).filter(Translation.translation_key == key, Translation.language_code == code).first()
if row:
row.text_value = text_value
else:
db.add(Translation(translation_key=key, language_code=code, text_value=text_value))
changed += 1
db.commit(); clear_translation_cache()
return RedirectResponse(f"/settings/translations?toast_success={changed} translations imported", status_code=303)
@app.get("/settings/appearance")
def settings_appearance_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
categories = db.query(Category).order_by(Category.name.asc()).all()
return templates.TemplateResponse(
"settings.html",
{
"request": request,
"settings": public_config(),
"categories": categories,
"languages": i18n_languages(db),
"settings_section": "appearance",
"system_storage": system_storage_information(),
},
)
@app.get("/settings/authentication")
def settings_authentication_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
categories = db.query(Category).order_by(Category.name.asc()).all()
return templates.TemplateResponse(
"settings.html",
{
"request": request,
"settings": public_config(),
"categories": categories,
"languages": i18n_languages(db),
"settings_section": "authentication",
"system_storage": system_storage_information(),
},
)
@app.get("/settings/system")
def settings_system_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
categories = db.query(Category).order_by(Category.name.asc()).all()
return templates.TemplateResponse(
"settings.html",
{
"request": request,
"settings": public_config(),
"categories": categories,
"languages": i18n_languages(db),
"settings_section": "system",
"system_storage": system_storage_information(),
},
)
@app.get("/settings/import-export")
def settings_import_export_page(request: Request):
_require_admin(request)
return templates.TemplateResponse(
"settings_import_export.html",
{"request": request},
)
@app.get("/settings/software")
def settings_software_page(request: Request):
_require_admin(request)
settings = _software_settings()
callback_base = str(settings.get("callback_base_url", "") or "").strip().rstrip("/")
return templates.TemplateResponse(
"settings_software.html",
{
"request": request,
"software_settings": settings,
"callback_health_url": f"{callback_base}/api/software-callback/health" if callback_base else "",
"debug_log": _software_debug_tail(),
"debug_log_path": str(SOFTWARE_CALLBACK_DEBUG_LOG),
},
)
@app.post("/settings/software")
def settings_software_save(
request: Request,
dispatch_delay_seconds: int = Form(2),
callback_base_url: str = Form(""),
callback_timeout_minutes: int = Form(30),
callback_worker_count: int = Form(4),
callback_test_timeout_seconds: int = Form(10),
callback_test_verify_tls: str | None = Form(None),
debug_logging: str | None = Form(None),
debug_log_max_lines: int = Form(1000),
automatic_retry_enabled: str | None = Form(None),
automatic_retry_max_age_hours: int = Form(12),
automatic_retry_interval_seconds: int = Form(60),
automatic_retry_max_retries: int = Form(3),
inventory_exclusion_platform: list[str] = Form(default=[]),
inventory_exclusion_name: list[str] = Form(default=[]),
):
_require_admin(request)
callback_base_url = callback_base_url.strip().rstrip("/")
if callback_base_url:
parsed = urllib.parse.urlsplit(callback_base_url)
if parsed.scheme not in {"http", "https"} or not parsed.hostname:
return RedirectResponse(
"/settings/software?toast_error=" + quote(_translate_request(request, "software.settings.invalid_url", "Enter a valid HTTP or HTTPS base URL.")),
status_code=303,
)
submitted_exclusion_rules = _normalize_software_exclusion_rules([
{"platform": platform, "name_contains": name}
for platform, name in zip(inventory_exclusion_platform, inventory_exclusion_name)
])
current = load_config()
software = dict(current.get("software", {}))
software.update({
"dispatch_delay_seconds": max(0, min(int(dispatch_delay_seconds), 60)),
"callback_base_url": callback_base_url,
"callback_timeout_minutes": max(5, min(int(callback_timeout_minutes), 240)),
"callback_worker_count": max(1, min(int(callback_worker_count), 8)),
"callback_test_timeout_seconds": max(2, min(int(callback_test_timeout_seconds), 60)),
"callback_test_verify_tls": callback_test_verify_tls == "on",
"debug_logging": debug_logging == "on",
"debug_log_max_lines": max(100, min(int(debug_log_max_lines), 10000)),
"automatic_retry_enabled": automatic_retry_enabled == "on",
"automatic_retry_max_age_hours": max(1, min(int(automatic_retry_max_age_hours), 168)),
"automatic_retry_interval_seconds": max(10, min(int(automatic_retry_interval_seconds), 86400)),
"automatic_retry_max_retries": max(1, min(int(automatic_retry_max_retries), 10)),
"inventory_exclusion_rules": submitted_exclusion_rules,
})
save_config({"software": software})
_software_debug_log(
f"SETTINGS saved | dispatch_delay_seconds={software['dispatch_delay_seconds']} | callback_base_url={callback_base_url or '<automatic>'} | callback_timeout_minutes={software['callback_timeout_minutes']} | callback_worker_count={software['callback_worker_count']} | automatic_retry_enabled={software['automatic_retry_enabled']} | automatic_retry_max_age_hours={software['automatic_retry_max_age_hours']} | automatic_retry_interval_seconds={software['automatic_retry_interval_seconds']} | automatic_retry_max_retries={software['automatic_retry_max_retries']} | inventory_exclusion_rules={len(submitted_exclusion_rules)} | verify_tls={software['callback_test_verify_tls']}",
force=True,
)
return RedirectResponse(
"/settings/software?toast_success=" + quote(_translate_request(request, "software.settings.saved", "Software settings saved.")),
status_code=303,
)
@app.post("/settings/software/test-callback")
def settings_software_test_callback(request: Request):
_require_admin(request)
settings = _software_settings()
callback_base = str(settings.get("callback_base_url", "") or "").strip().rstrip("/")
if not callback_base:
callback_base = str(request.base_url).rstrip("/")
health_url = f"{callback_base}/api/software-callback/health"
timeout = max(2, min(int(settings.get("callback_test_timeout_seconds", 10) or 10), 60))
verify_tls = bool(settings.get("callback_test_verify_tls", True))
parsed = urllib.parse.urlsplit(health_url)
_software_debug_log("=" * 72, force=True)
_software_debug_log("CALLBACK TEST started", force=True)
_software_debug_log(f"Base URL: {callback_base}", force=True)
_software_debug_log(f"Health URL: {health_url}", force=True)
_software_debug_log(f"Scheme: {parsed.scheme} | Host: {parsed.hostname} | Port: {parsed.port or ('443' if parsed.scheme == 'https' else '80')}", force=True)
_software_debug_log(f"Timeout: {timeout}s | Verify TLS: {verify_tls}", force=True)
try:
dns_started = time.monotonic()
addresses = sorted({item[4][0] for item in socket.getaddrinfo(parsed.hostname, parsed.port or (443 if parsed.scheme == "https" else 80), type=socket.SOCK_STREAM)})
_software_debug_log(f"DNS: {', '.join(addresses)} ({(time.monotonic()-dns_started)*1000:.1f} ms)", force=True)
except Exception as exc:
_software_debug_log(f"DNS ERROR: {type(exc).__name__}: {exc}", force=True)
context = None
if parsed.scheme == "https" and not verify_tls:
context = ssl._create_unverified_context()
request_started = time.monotonic()
try:
req = urllib.request.Request(
health_url,
headers={"User-Agent": f"AssetManager/{application_version()} CallbackTest"},
method="GET",
)
with urllib.request.urlopen(req, timeout=timeout, context=context) as response:
body = response.read(65536).decode("utf-8", errors="replace")
elapsed = (time.monotonic() - request_started) * 1000
_software_debug_log(f"HTTP STATUS: {response.status} {response.reason} ({elapsed:.1f} ms)", force=True)
_software_debug_log(f"Content-Type: {response.headers.get('Content-Type', '-')}", force=True)
_software_debug_log(f"Response body: {body}", force=True)
payload = json.loads(body)
if response.status == 200 and payload.get("ok") is True:
_software_debug_log("RESULT: Callback health endpoint is reachable.", force=True)
return RedirectResponse(
"/settings/software?toast_success=" + quote(_translate_request(request, "software.settings.test_success", "Callback endpoint is reachable.")),
status_code=303,
)
raise ValueError("Health response did not contain ok=true")
except urllib.error.HTTPError as exc:
body = exc.read(65536).decode("utf-8", errors="replace")
_software_debug_log(f"HTTP ERROR: {exc.code} {exc.reason}", force=True)
_software_debug_log(f"Response body: {body}", force=True)
message = f"HTTP {exc.code}: {exc.reason}"
except Exception as exc:
_software_debug_log(f"REQUEST ERROR: {type(exc).__name__}: {exc}", force=True)
message = f"{type(exc).__name__}: {exc}"
_software_debug_log("RESULT: Callback health endpoint is not reachable.", force=True)
return RedirectResponse(
"/settings/software?toast_error=" + quote(_translate_request(request, "software.settings.test_failed", "Callback test failed: {error}", error=message)),
status_code=303,
)
@app.get("/settings/software/debug-log")
def settings_software_debug_log(request: Request):
_require_admin(request)
return Response(_software_debug_tail(), media_type="text/plain; charset=utf-8", headers={"Cache-Control": "no-store"})
@app.post("/settings/software/debug-log/clear")
def settings_software_debug_log_clear(request: Request):
_require_admin(request)
try:
SOFTWARE_CALLBACK_DEBUG_LOG.unlink(missing_ok=True)
except OSError:
logger.exception("Software callback debug log could not be cleared")
return RedirectResponse(
"/settings/software?toast_success=" + quote(_translate_request(request, "software.settings.log_cleared", "Debug log cleared.")),
status_code=303,
)
@app.get("/settings/backup")
def settings_backup_page(request: Request):
_require_admin(request)
return templates.TemplateResponse(
"settings_backup.html",
{
"request": request,
"backups": list_backups(),
"backup_directory": str(BACKUP_DIR),
"backup_interval_hours": BACKUP_INTERVAL_HOURS,
"backup_retention_days": BACKUP_RETENTION_DAYS,
},
)
@app.post("/settings/backup/create")
def settings_backup_create(request: Request):
_require_admin(request)
try:
created_by = _changed_by(request) or "admin"
create_backup(created_by=created_by, reason="manual")
return RedirectResponse("/settings/backup?toast_success=" + quote(_translate_request(request, "backup.created_success", "Backup created")), status_code=303)
except Exception as exc:
logger.exception("Backup creation failed")
return RedirectResponse("/settings/backup?toast_error=" + quote(str(exc)), status_code=303)
@app.get("/settings/backup/{backup_name}/download")
def settings_backup_download(request: Request, backup_name: str):
_require_admin(request)
try:
path = backup_path(backup_name)
if not path.is_file():
raise HTTPException(status_code=404)
return FileResponse(path, media_type="application/zip", filename=path.name)
except ValueError:
raise HTTPException(status_code=404)
@app.post("/settings/backup/{backup_name}/delete")
def settings_backup_delete(request: Request, backup_name: str):
_require_admin(request)
try:
delete_backup(backup_name)
return RedirectResponse("/settings/backup?toast_success=" + quote(_translate_request(request, "backup.deleted_success", "Backup deleted")), status_code=303)
except Exception as exc:
logger.exception("Backup deletion failed")
return RedirectResponse("/settings/backup?toast_error=" + quote(str(exc)), status_code=303)
@app.post("/settings/backup/upload")
async def settings_backup_upload(request: Request, backup_file: UploadFile = File(...)):
_require_admin(request)
try:
content = await backup_file.read()
store_uploaded_backup(backup_file.filename or "backup.zip", content)
return RedirectResponse("/settings/backup?toast_success=" + quote(_translate_request(request, "backup.uploaded_success", "Backup uploaded")), status_code=303)
except Exception as exc:
logger.exception("Backup upload failed")
return RedirectResponse("/settings/backup?toast_error=" + quote(str(exc)), status_code=303)
@app.post("/settings/backup/{backup_name}/restore")
def settings_backup_restore(request: Request, backup_name: str, confirmation: str = Form("")):
_require_admin(request)
if confirmation.strip().upper() != "RESTORE":
return RedirectResponse("/settings/backup?toast_warning=" + quote(_translate_request(request, "backup.restore_confirmation_missing", "Restore confirmation is missing")), status_code=303)
try:
# Safety copy immediately before replacing data.
create_backup(created_by="system", reason="before_restore")
restore_backup(backup_name)
clear_translation_cache()
return RedirectResponse("/settings/backup?toast_success=" + quote(_translate_request(request, "backup.restored_success", "Backup restored")), status_code=303)
except Exception as exc:
logger.exception("Backup restore failed")
return RedirectResponse("/settings/backup?toast_error=" + quote(str(exc)), status_code=303)
@app.get("/settings")
def settings_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
categories = db.query(Category).order_by(Category.name.asc()).all()
return templates.TemplateResponse(
"settings.html",
{
"request": request,
"settings": public_config(),
"categories": categories,
"languages": i18n_languages(db),
"settings_section": "appearance",
"system_storage": system_storage_information(),
},
)
THEME_COLOR_DEFAULTS = {
"primary_color": "#0878bd",
"primary_hover_color": "#075f94",
"header_color": "#063c66",
"navigation_color": "#ffffff",
"page_background": "#f2f4f6",
"surface_color": "#ffffff",
"text_color": "#1f2933",
"border_color": "#d7dde3",
"success_color": "#198754",
"warning_color": "#d97706",
"danger_color": "#b42318",
}
def _safe_hex_color(value: str, fallback: str) -> str:
candidate = (value or "").strip()
return candidate.lower() if re.fullmatch(r"#[0-9a-fA-F]{6}", candidate) else fallback
@app.post("/settings")
async def settings_save(
request: Request,
title: str = Form("AssetManager"),
company_name: str = Form(""),
about_author: str = Form(""),
about_organization: str = Form(""),
about_contact: str = Form(""),
about_website: str = Form(""),
about_repository: str = Form(""),
about_license: str = Form(""),
about_copyright: str = Form(""),
callback_base_url: str = Form(""),
default_language: str = Form("en"),
custom_css: str = Form(""),
default_theme: str = Form("light"),
primary_color: str = Form("#0878bd"),
primary_hover_color: str = Form("#075f94"),
header_color: str = Form("#063c66"),
navigation_color: str = Form("#ffffff"),
page_background: str = Form("#f2f4f6"),
surface_color: str = Form("#ffffff"),
text_color: str = Form("#1f2933"),
border_color: str = Form("#d7dde3"),
success_color: str = Form("#198754"),
warning_color: str = Form("#d97706"),
danger_color: str = Form("#b42318"),
mesh_fallback_category_id: int = Form(1),
mesh_category_1: int | None = Form(None),
mesh_category_2: int | None = Form(None),
mesh_category_3: int | None = Form(None),
mesh_category_4: int | None = Form(None),
mesh_category_5: int | None = Form(None),
mesh_category_6: int | None = Form(None),
mesh_category_7: int | None = Form(None),
mesh_category_8: int | None = Form(None),
auth_mode: str = Form("none"),
ldap_server: str = Form(""),
ldap_port: int = Form(389),
ldap_use_ssl: str | None = Form(None),
ldap_start_tls: str | None = Form(None),
ldap_bind_dn: str = Form(""),
ldap_bind_password_env: str = Form("LDAP_BIND_PASSWORD"),
ldap_user_base_dn: str = Form(""),
ldap_user_filter: str = Form("(sAMAccountName={username})"),
ldap_display_name_attribute: str = Form("displayName"),
ldap_email_attribute: str = Form("mail"),
remove_logo: str | None = Form(None),
remove_favicon: str | None = Form(None),
logo: UploadFile | None = File(None),
favicon: UploadFile | None = File(None),
):
_require_admin(request)
current = load_config()
general = dict(current.get("general", {}))
general["title"] = title.strip() or "AssetManager"
general["company_name"] = company_name.strip()
general["about_author"] = about_author.strip()
general["about_organization"] = about_organization.strip()
general["about_contact"] = about_contact.strip()
general["about_website"] = about_website.strip()
general["about_repository"] = about_repository.strip()
general["about_license"] = about_license.strip()
general["about_copyright"] = about_copyright.strip()
general["callback_base_url"] = callback_base_url.strip().rstrip("/")
general["default_language"] = default_language if default_language in {"en", "de"} else "en"
general["custom_css"] = custom_css.replace("\x00", "").strip()
general["default_theme"] = default_theme if default_theme in {"light", "dark"} else "light"
colors = {}
submitted_colors = {
"primary_color": primary_color,
"primary_hover_color": primary_hover_color,
"header_color": header_color,
"navigation_color": navigation_color,
"page_background": page_background,
"surface_color": surface_color,
"text_color": text_color,
"border_color": border_color,
"success_color": success_color,
"warning_color": warning_color,
"danger_color": danger_color,
}
for color_name, fallback in THEME_COLOR_DEFAULTS.items():
colors[color_name] = _safe_hex_color(submitted_colors.get(color_name, ""), fallback)
general["colors"] = colors
if remove_logo == "on":
old_logo = general.get("logo")
if old_logo:
delete_asset_image(old_logo)
general["logo"] = ""
if logo and logo.filename:
old_logo = general.get("logo")
new_logo = save_upload(logo)
if old_logo and old_logo != new_logo:
delete_asset_image(old_logo)
general["logo"] = new_logo or ""
if remove_favicon == "on":
old_favicon = general.get("favicon")
if old_favicon:
delete_asset_image(old_favicon)
general["favicon"] = ""
if favicon and favicon.filename:
old_favicon = general.get("favicon")
new_favicon = save_upload(favicon)
if old_favicon and old_favicon != new_favicon:
delete_asset_image(old_favicon)
general["favicon"] = new_favicon or ""
if auth_mode not in {"none", "local", "ldap"}:
auth_mode = "none"
authentication = dict(current.get("authentication", {}))
authentication["mode"] = auth_mode
ldap = dict(authentication.get("ldap", {}))
ldap.update({
"server": ldap_server.strip(),
"port": ldap_port,
"use_ssl": ldap_use_ssl == "on",
"start_tls": ldap_start_tls == "on",
"bind_dn": ldap_bind_dn.strip(),
"bind_password_env": ldap_bind_password_env.strip() or "LDAP_BIND_PASSWORD",
"user_base_dn": ldap_user_base_dn.strip(),
"user_filter": ldap_user_filter.strip() or "(sAMAccountName={username})",
"display_name_attribute": ldap_display_name_attribute.strip() or "displayName",
"email_attribute": ldap_email_attribute.strip() or "mail",
})
authentication["ldap"] = ldap
meshcentral = dict(current.get("meshcentral", {}))
meshcentral["fallback_category_id"] = int(mesh_fallback_category_id or 1)
meshcentral["category_mapping"] = {
"1": mesh_category_1,
"2": mesh_category_2,
"3": mesh_category_3,
"4": mesh_category_4,
"5": mesh_category_5,
"6": mesh_category_6,
"7": mesh_category_7,
"8": mesh_category_8,
}
save_config({"general": general, "authentication": authentication, "meshcentral": meshcentral})
return RedirectResponse("/settings?toast_success=Einstellungen gespeichert", status_code=303)
@app.get("/sync/meshcentral")
def meshcentral_sync_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
runs = db.query(SyncRun).order_by(SyncRun.started_at.desc()).limit(20).all()
linked = db.query(Asset).filter(Asset.mesh_node_id.isnot(None)).count()
conflicts = db.query(Asset).filter(Asset.mesh_sync_status == "conflict").count()
missing = db.query(Asset).filter(Asset.mesh_sync_status == "missing").count()
definitions = db.query(FieldDefinition).filter(FieldDefinition.is_active.is_(True)).order_by(FieldDefinition.sort_order, FieldDefinition.label).all()
mappings = db.query(MeshFieldMapping).order_by(MeshFieldMapping.priority, MeshFieldMapping.id).all()
return templates.TemplateResponse("meshcentral_sync.html", {
"request": request, "runs": runs, "config": public_config(),
"linked": linked, "conflicts": conflicts, "missing": missing,
"definitions": definitions,
"mappings": mappings,
"transform_options": [
"raw", "string", "integer", "decimal", "boolean", "json",
"normalize_serial", "normalize_mac", "active_ipv4_address", "active_ipv4_mac",
"sum_memory_gb", "format_memory_modules", "sum_storage_gb", "format_drives",
"gpu_names", "active_antivirus",
],
})
@app.post("/sync/meshcentral/start")
def meshcentral_sync_start(request: Request):
_require_admin(request)
with SYNC_JOBS_LOCK:
if any(not job.get("finished") for job in SYNC_JOBS.values()):
return JSONResponse({"error": "Es läuft bereits eine Synchronisierung."}, status_code=409)
job_id = uuid.uuid4().hex
SYNC_JOBS[job_id] = {"status": "running", "finished": False, "run_id": None, "logs": [], "created_at": datetime.now().isoformat()}
_job_log(job_id, "info", "Synchronisierungsauftrag angelegt.")
threading.Thread(target=_run_sync_job, args=(job_id,), daemon=True).start()
return {"job_id": job_id}
@app.get("/sync/meshcentral/jobs/{job_id}")
def meshcentral_sync_job(job_id: str, offset: int = 0):
with SYNC_JOBS_LOCK:
job = SYNC_JOBS.get(job_id)
if not job:
raise HTTPException(404, "Synchronisierungsauftrag nicht gefunden")
logs = list(job["logs"][max(0, offset):])
return {
"status": job["status"],
"finished": job["finished"],
"run_id": job["run_id"],
"logs": logs,
"next_offset": max(0, offset) + len(logs),
}
@app.post("/sync/meshcentral")
def meshcentral_sync_run():
return RedirectResponse("/sync/meshcentral", status_code=303)
# ---------------------------------------------------------------------------
# v0.3.8.0: central field catalog and MeshCentral synchronization settings
# ---------------------------------------------------------------------------
@app.get('/settings/fields')
def field_definitions_page(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
definitions = db.query(FieldDefinition).order_by(FieldDefinition.sort_order, FieldDefinition.label).all()
value_counts = dict(db.query(AssetFieldValue.field_definition_id, __import__('sqlalchemy').func.count(AssetFieldValue.id)).group_by(AssetFieldValue.field_definition_id).all())
return templates.TemplateResponse('field_definitions.html', {'request': request, 'definitions': definitions, 'data_types': FIELD_DATA_TYPES, 'value_counts': value_counts})
@app.get('/settings/fields/new')
def field_definition_new(request: Request):
_require_admin(request)
return templates.TemplateResponse('field_definition_form.html', {'request': request, 'definition': None, 'data_types': FIELD_DATA_TYPES})
@app.post('/settings/fields/new')
def field_definition_create(request: Request, field_name: str = Form(...), label: str = Form(...), description: str = Form(''), data_type: str = Form('short_text'), readonly: str | None = Form(None), searchable: str | None = Form(None), filterable: str | None = Form(None), chart_enabled: str | None = Form(None), excel_import: str | None = Form(None), excel_export: str | None = Form(None), bulk_edit: str | None = Form(None), is_default: str | None = Form(None), list_width: str | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
name = re.sub(r'[^a-z0-9_]', '_', field_name.strip().lower()).strip('_')
if not name or name in ASSET_FIELDS:
raise HTTPException(400, 'The field name is invalid or already reserved.')
if data_type not in FIELD_DATA_TYPES:
raise HTTPException(400, 'Invalid data type.')
definition = FieldDefinition(field_name=name, label=label.strip(), description=description.strip() or None, data_type=data_type, readonly=readonly == 'on', is_system=False, is_active=True, translation_key=f'field.{name}', searchable=searchable == 'on', filterable=filterable == 'on', chart_enabled=chart_enabled == 'on', excel_import=excel_import == 'on', excel_export=excel_export == 'on', bulk_edit=bulk_edit == 'on', is_default=is_default == 'on', list_width=_optional_form_int(list_width, field_name="list_width", minimum=40, maximum=800), sort_order=(db.query(FieldDefinition).count()+1))
db.add(definition)
try:
db.flush()
for category in db.query(Category).all():
db.add(CategoryField(category_id=category.id, field_name=name, label=definition.label, field_definition_id=definition.id, active=False, required=False, show_in_list=False, readonly=definition.readonly, sort_order=definition.sort_order))
# Add English fallback translation; German can be maintained in translation editor.
db.add(Translation(translation_key=definition.translation_key, language_code='en', text_value=definition.label))
db.commit(); clear_translation_cache()
except IntegrityError as exc:
db.rollback(); raise HTTPException(400, 'A field with this name already exists.') from exc
return RedirectResponse('/settings/fields?toast_success=Field created', status_code=303)
@app.get('/settings/fields/{definition_id}/edit')
def field_definition_edit(definition_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
definition = db.get(FieldDefinition, definition_id)
if not definition: raise HTTPException(404, 'Field not found')
return templates.TemplateResponse('field_definition_form.html', {'request': request, 'definition': definition, 'data_types': FIELD_DATA_TYPES})
@app.post('/settings/fields/{definition_id}/edit')
def field_definition_update(definition_id: int, request: Request, label: str = Form(...), description: str = Form(''), data_type: str = Form('short_text'), readonly: str | None = Form(None), is_active: str | None = Form(None), searchable: str | None = Form(None), filterable: str | None = Form(None), chart_enabled: str | None = Form(None), excel_import: str | None = Form(None), excel_export: str | None = Form(None), bulk_edit: str | None = Form(None), is_default: str | None = Form(None), list_width: str | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
definition = db.get(FieldDefinition, definition_id)
if not definition: raise HTTPException(404, 'Field not found')
if data_type not in FIELD_DATA_TYPES: raise HTTPException(400, 'Invalid data type')
if not definition.is_system and definition.data_type != data_type and db.query(AssetFieldValue).filter(AssetFieldValue.field_definition_id == definition.id).count():
raise HTTPException(400, 'The data type can only be changed while the custom field has no stored values.')
definition.label=label.strip(); definition.description=description.strip() or None; definition.data_type=data_type
definition.readonly=True if definition.field_name in SYSTEM_READONLY_FIELDS else readonly == 'on'
definition.is_active=True if definition.is_system else is_active == 'on'
definition.searchable=searchable == 'on'; definition.filterable=filterable == 'on'; definition.chart_enabled=chart_enabled == 'on'
definition.excel_import=False if definition.readonly else excel_import == 'on'; definition.excel_export=excel_export == 'on'; definition.bulk_edit=False if definition.readonly else bulk_edit == 'on'; definition.is_default=is_default == 'on'; definition.list_width=_optional_form_int(list_width, field_name="list_width", minimum=40, maximum=800)
for link in definition.category_links:
link.label=definition.label; link.readonly=definition.readonly
db.commit(); clear_translation_cache()
return RedirectResponse('/settings/fields?toast_success=Field saved', status_code=303)
@app.post('/settings/fields/{definition_id}/delete')
def field_definition_delete(definition_id: int, request: Request, confirm_values: str | None = Form(None), db: Session = Depends(get_db)):
_require_admin(request)
definition = db.get(FieldDefinition, definition_id)
if not definition: raise HTTPException(404, 'Field not found')
if definition.is_system: raise HTTPException(400, 'System fields cannot be deleted.')
count = db.query(AssetFieldValue).filter(AssetFieldValue.field_definition_id == definition.id).count()
if count and confirm_values != 'on': raise HTTPException(400, f'This field contains values for {count} assets. Confirm deletion of all values.')
db.delete(definition); db.commit()
return RedirectResponse('/settings/fields?toast_success=Field deleted', status_code=303)
@app.post('/sync/meshcentral/settings')
async def meshcentral_settings_save(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
current = load_config()
mesh = dict(current.get('meshcentral', {}))
try:
timeout_seconds = max(10, min(int(form.get('timeout_seconds') or 600), 3600))
except (TypeError, ValueError):
timeout_seconds = 600
mesh.update({
'enabled': form.get('enabled') == 'on',
'url': (form.get('url') or '').strip(),
'username': (form.get('username') or '').strip(),
'password_env': (form.get('password_env') or 'MESHCENTRAL_PASSWORD').strip(),
'tenant': (form.get('tenant') or '').strip(),
'meshctrl_path': (form.get('meshctrl_path') or '/opt/meshcentral/node_modules/meshcentral/meshctrl.js').strip(),
'verify_tls': form.get('verify_tls') == 'on',
'include_details': form.get('include_details') == 'on',
'timeout_seconds': timeout_seconds,
'create_missing_assets': form.get('create_missing_assets') == 'on',
'mark_missing_devices': form.get('mark_missing_devices') == 'on',
'store_source_json': form.get('store_source_json') == 'on',
})
mesh['match_order'] = form.getlist('match_order') or ['node_id', 'serial_number']
# Klartextpasswort aus Altversionen nicht weiter speichern.
mesh.pop('password', None)
mesh.pop('field_mappings', None)
mesh.pop('field_rules', None)
save_config({'meshcentral': mesh})
return RedirectResponse('/sync/meshcentral?toast_success=MeshCentral-Einstellungen gespeichert', status_code=303)
@app.post('/sync/meshcentral/mappings/save')
async def meshcentral_mappings_save(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
allowed_transforms = {
'raw','string','integer','decimal','boolean','json','normalize_serial','normalize_mac',
'active_ipv4_address','active_ipv4_mac','sum_memory_gb','format_memory_modules',
'sum_storage_gb','format_drives','gpu_names','active_antivirus'
}
allowed_multi = {'first','join','json','count'}
allowed_rules = {'fill_empty','meshcentral_wins','local_wins'}
for row in db.query(MeshFieldMapping).all():
prefix = f'mapping_{row.id}_'
source = str(form.get(prefix + 'source_path') or '').strip()
target = str(form.get(prefix + 'target_field_name') or '').strip()
if not source or not target:
continue
row.source_path = source
row.target_field_name = target
transform = str(form.get(prefix + 'transform') or 'raw')
row.transform = transform if transform in allowed_transforms else 'raw'
multi = str(form.get(prefix + 'multi_value_mode') or 'first')
row.multi_value_mode = multi if multi in allowed_multi else 'first'
row.separator = str(form.get(prefix + 'separator') or '\\n').replace('\\n', '\n').replace('\\t', '\t')
row.enabled = form.get(prefix + 'enabled') == 'on'
try: row.priority = int(form.get(prefix + 'priority') or 0)
except (TypeError, ValueError): row.priority = 0
rule = str(form.get(prefix + 'update_rule') or 'fill_empty')
row.update_rule = rule if rule in allowed_rules else 'fill_empty'
row.description = str(form.get(prefix + 'description') or '').strip() or None
db.commit()
return RedirectResponse('/sync/meshcentral?toast_success=Mapping-Regeln gespeichert', status_code=303)
@app.post('/sync/meshcentral/mappings/new')
async def meshcentral_mapping_new(request: Request, db: Session = Depends(get_db)):
_require_admin(request)
form = await request.form()
source = str(form.get('source_path') or '').strip()
target = str(form.get('target_field_name') or '').strip()
if not source or not target:
return RedirectResponse('/sync/meshcentral?toast_error=Quellpfad und Zielfeld sind erforderlich', status_code=303)
definition = db.query(FieldDefinition).filter(FieldDefinition.field_name == target, FieldDefinition.is_active.is_(True)).first()
if not definition:
return RedirectResponse('/sync/meshcentral?toast_error=Das Zielfeld existiert nicht oder ist inaktiv', status_code=303)
max_priority = db.query(MeshFieldMapping).order_by(MeshFieldMapping.priority.desc()).first()
row = MeshFieldMapping(
source_path=source,
target_field_name=target,
transform=str(form.get('transform') or 'raw'),
multi_value_mode=str(form.get('multi_value_mode') or 'first'),
separator=str(form.get('separator') or '\\n').replace('\\n', '\n').replace('\\t', '\t'),
enabled=form.get('enabled') == 'on',
priority=(max_priority.priority + 10) if max_priority else 10,
description=str(form.get('description') or '').strip() or None,
update_rule=str(form.get('update_rule') or 'fill_empty'),
is_system=False,
)
db.add(row)
db.commit()
return RedirectResponse('/sync/meshcentral?toast_success=Neue Mapping-Regel angelegt', status_code=303)
@app.post('/sync/meshcentral/mappings/{mapping_id}/delete')
def meshcentral_mapping_delete(mapping_id: int, request: Request, db: Session = Depends(get_db)):
_require_admin(request)
row = db.query(MeshFieldMapping).filter(MeshFieldMapping.id == mapping_id).first()
if not row:
raise HTTPException(status_code=404, detail='Mapping-Regel nicht gefunden')
if row.is_system:
return RedirectResponse('/sync/meshcentral?toast_error=System-Mappings können deaktiviert, aber nicht gelöscht werden', status_code=303)
db.delete(row)
db.commit()
return RedirectResponse('/sync/meshcentral?toast_success=Mapping-Regel gelöscht', status_code=303)