software package profiles implemeted, job deletion optimized
This commit is contained in:
+202
-11
@@ -48,7 +48,7 @@ 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 .software_control import token_hash, detect_platform, execute_job, build_registry_user_script, cleanup_stale_remote_job_directories
|
||||
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
|
||||
@@ -56,19 +56,20 @@ from .privacy import merge_privacy_settings, localized_privacy_settings, normali
|
||||
from .privacy_retention import check_retention_category, delete_retention_category, append_deletion_audit, deletion_audit_tail, IMPLEMENTED_RETENTION_KEYS
|
||||
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 .setup_analyzer import register_setup_analyzer
|
||||
from .software_packages import delete_package_storage, human_size, load_package_manifest, package_execution_timeout_seconds, package_summary, update_package_post_install, update_package_process_control
|
||||
from .software_packages import delete_package_storage, human_size, load_package_manifest, package_execution_timeout_seconds, package_summary, update_package_post_install, update_package_process_control, export_package_bundle, import_package_bundle, PACKAGE_IMPORT_MAX_MB
|
||||
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"
|
||||
DATA_ROOT = Path(os.getenv("ASSETMANAGER_DATA_ROOT", "/assetmanager-data"))
|
||||
UPLOAD_DIR = Path(os.getenv("UPLOAD_DIR", str(DATA_ROOT / "uploads")))
|
||||
UPLOAD_DIR.mkdir(parents=True, exist_ok=True)
|
||||
STANDARD_IMAGE_DIR = Path(os.getenv("STANDARD_IMAGE_DIR", str(BASE_DIR / "static" / "uploads" / "library")))
|
||||
STANDARD_IMAGE_DIR = Path(os.getenv("STANDARD_IMAGE_DIR", str(UPLOAD_DIR / "library")))
|
||||
STANDARD_IMAGE_DIR.mkdir(parents=True, exist_ok=True)
|
||||
STANDARD_IMAGE_EXTENSIONS = {".png", ".jpg", ".jpeg", ".webp", ".gif"}
|
||||
LOG_DIR = Path(os.getenv("SYNC_LOG_DIR", "/app/data/logs/sync"))
|
||||
LOG_DIR = Path(os.getenv("SYNC_LOG_DIR", str(DATA_ROOT / "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 = Path(os.getenv("APP_LOG_DIR", str(DATA_ROOT / "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()
|
||||
@@ -226,6 +227,9 @@ callback_app = FastAPI(
|
||||
_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"
|
||||
# Persistent uploads live outside the application image in Docker deployments.
|
||||
# Mount this route before /static so existing /static/uploads/... URLs remain valid.
|
||||
app.mount("/static/uploads", StaticFiles(directory=UPLOAD_DIR), name="uploads")
|
||||
app.mount("/static", StaticFiles(directory=BASE_DIR / "static"), name="static")
|
||||
templates = Jinja2Templates(directory=BASE_DIR / "templates")
|
||||
templates.env.globals["application_config"] = load_config
|
||||
@@ -247,11 +251,12 @@ def application_version() -> str:
|
||||
|
||||
|
||||
templates.env.globals["application_version"] = application_version
|
||||
templates.env.globals["application_info_path"] = lambda: str(Path(os.getenv("APPINFO_PATH", str(DATA_ROOT / "config" / "APPINFO.json"))))
|
||||
|
||||
def application_info() -> dict[str, str]:
|
||||
"""Read persistent application metadata.
|
||||
|
||||
Primary file: /app/config/APPINFO.json (or APPINFO_PATH).
|
||||
Primary file: APPINFO_PATH below the persistent AssetManager data root.
|
||||
A legacy APPINFO.json is migrated once when possible.
|
||||
"""
|
||||
defaults = {
|
||||
@@ -259,7 +264,7 @@ def application_info() -> dict[str, str]:
|
||||
"contact": "", "website": "", "repository": "",
|
||||
"license": "", "copyright": "", "description": "",
|
||||
}
|
||||
persistent = Path(os.getenv("APPINFO_PATH", "/app/config/APPINFO.json"))
|
||||
persistent = Path(os.getenv("APPINFO_PATH", str(DATA_ROOT / "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
|
||||
@@ -401,6 +406,102 @@ 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"}
|
||||
REMOTE_JOB_CLEANUP_INTERVAL_SECONDS = 3600
|
||||
_REMOTE_JOB_CLEANUP_LAST_RUN_MONOTONIC = 0.0
|
||||
|
||||
|
||||
def _remote_job_cleanup_settings() -> tuple[bool, int]:
|
||||
settings = _software_settings()
|
||||
enabled = bool(settings.get("remote_job_cleanup_enabled", True))
|
||||
try:
|
||||
retention_hours = int(settings.get("remote_job_retention_hours", 24) or 24)
|
||||
except (TypeError, ValueError):
|
||||
retention_hours = 24
|
||||
return enabled, max(1, min(retention_hours, 8760))
|
||||
|
||||
|
||||
def _active_remote_job_directory_names(db: Session, asset_id: int) -> list[str]:
|
||||
rows = (
|
||||
db.query(SoftwareJob.id, SoftwareJob.attempt_count)
|
||||
.filter(
|
||||
SoftwareJob.asset_id == asset_id,
|
||||
SoftwareJob.status.notin_(SOFTWARE_JOB_TERMINAL_STATES),
|
||||
SoftwareJob.attempt_count > 0,
|
||||
)
|
||||
.all()
|
||||
)
|
||||
return [f"{job_id}-{int(attempt_count or 0)}" for job_id, attempt_count in rows if int(attempt_count or 0) > 0]
|
||||
|
||||
|
||||
def _run_remote_job_cleanup_maintenance() -> None:
|
||||
enabled, retention_hours = _remote_job_cleanup_settings()
|
||||
if not enabled:
|
||||
return
|
||||
|
||||
runtime_config = load_config()
|
||||
cfg = runtime_config.get("meshcentral", {})
|
||||
password_env = str(cfg.get("password_env") or "MESHCENTRAL_PASSWORD")
|
||||
password = os.getenv(password_env, "")
|
||||
if not password:
|
||||
logger.warning("Remote client job cleanup skipped because %s is not set", password_env)
|
||||
return
|
||||
|
||||
db = SessionLocal()
|
||||
try:
|
||||
assets = (
|
||||
db.query(Asset)
|
||||
.join(SoftwareJob, SoftwareJob.asset_id == Asset.id)
|
||||
.filter(
|
||||
Asset.mesh_node_id.isnot(None),
|
||||
Asset.mesh_online.is_(True),
|
||||
)
|
||||
.distinct()
|
||||
.order_by(Asset.id)
|
||||
.all()
|
||||
)
|
||||
if not assets:
|
||||
return
|
||||
|
||||
cleaned_clients = 0
|
||||
failed_clients = 0
|
||||
for asset in assets:
|
||||
platform = detect_platform(asset)
|
||||
if platform not in {"windows", "linux"}:
|
||||
continue
|
||||
excluded = _active_remote_job_directory_names(db, asset.id)
|
||||
try:
|
||||
result = cleanup_stale_remote_job_directories(
|
||||
cfg,
|
||||
asset,
|
||||
password,
|
||||
platform,
|
||||
retention_hours,
|
||||
min(max(30, int(cfg.get("timeout_seconds") or 120)), 120),
|
||||
excluded,
|
||||
)
|
||||
output = ((result.stdout or "") + "\n" + (result.stderr or "")).strip()
|
||||
if result.returncode == 0:
|
||||
cleaned_clients += 1
|
||||
if output:
|
||||
logger.info("Remote job cleanup asset=%s (%s): %s", asset.id, asset.name, output.replace("\n", " | "))
|
||||
else:
|
||||
failed_clients += 1
|
||||
logger.warning(
|
||||
"Remote job cleanup failed asset=%s (%s) rc=%s: %s",
|
||||
asset.id, asset.name, result.returncode, output,
|
||||
)
|
||||
except Exception as exc:
|
||||
failed_clients += 1
|
||||
logger.warning("Remote job cleanup exception asset=%s (%s): %s", asset.id, asset.name, exc)
|
||||
|
||||
if cleaned_clients or failed_clients:
|
||||
logger.info(
|
||||
"Remote client job cleanup cycle completed: clients=%s failed=%s retention_hours=%s",
|
||||
cleaned_clients, failed_clients, retention_hours,
|
||||
)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
|
||||
def _configured_callback_worker_count() -> int:
|
||||
@@ -555,6 +656,18 @@ def _software_job_timeout_loop(stop_event: threading.Event) -> None:
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
global _REMOTE_JOB_CLEANUP_LAST_RUN_MONOTONIC
|
||||
now_monotonic = time.monotonic()
|
||||
if (
|
||||
_REMOTE_JOB_CLEANUP_LAST_RUN_MONOTONIC <= 0
|
||||
or now_monotonic - _REMOTE_JOB_CLEANUP_LAST_RUN_MONOTONIC >= REMOTE_JOB_CLEANUP_INTERVAL_SECONDS
|
||||
):
|
||||
_REMOTE_JOB_CLEANUP_LAST_RUN_MONOTONIC = now_monotonic
|
||||
try:
|
||||
_run_remote_job_cleanup_maintenance()
|
||||
except Exception:
|
||||
logger.exception("Remote client job cleanup maintenance failed")
|
||||
|
||||
for job_id, token, callback_base in queued:
|
||||
threading.Thread(
|
||||
target=execute_job,
|
||||
@@ -5929,16 +6042,88 @@ def software_packages_page(request: Request, db: Session = Depends(get_db)):
|
||||
.order_by(func.lower(SoftwarePackage.name), SoftwarePackage.id)
|
||||
.all()
|
||||
)
|
||||
package_ids = [package.id for package in packages]
|
||||
job_counts: dict[int, int] = {}
|
||||
if package_ids:
|
||||
job_counts = {
|
||||
int(package_id): int(count or 0)
|
||||
for package_id, count in (
|
||||
db.query(SoftwareJob.package_id, func.count(SoftwareJob.id))
|
||||
.filter(SoftwareJob.package_id.in_(package_ids))
|
||||
.group_by(SoftwareJob.package_id)
|
||||
.all()
|
||||
)
|
||||
}
|
||||
package_rows = []
|
||||
for package in packages:
|
||||
row = package_summary(package)
|
||||
row["job_count"] = job_counts.get(package.id, 0)
|
||||
package_rows.append(row)
|
||||
return templates.TemplateResponse(
|
||||
"software_packages.html",
|
||||
{
|
||||
"request": request,
|
||||
"package_rows": [package_summary(package) for package in packages],
|
||||
"package_rows": package_rows,
|
||||
"human_size": human_size,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
@app.post("/software/packages/import")
|
||||
async def software_package_import(request: Request, package_file: UploadFile = File(...), confirm_scripts: str | None = Form(None), import_profile: str | None = Form(None), db: Session = Depends(get_db)):
|
||||
_require_admin(request)
|
||||
if str(confirm_scripts or "").strip().lower() not in {"1", "true", "yes", "on"}:
|
||||
message = _translate_request(request, "software_packages.import_confirm_required", "Confirm that imported packages can contain executable scripts.")
|
||||
return RedirectResponse("/software/packages?toast_error=" + quote(message), status_code=303)
|
||||
suffix = Path(package_file.filename or "package.ampkg").suffix.lower()
|
||||
if suffix not in {".ampkg", ".zip"}:
|
||||
await package_file.close()
|
||||
message = _translate_request(request, "software_packages.import_invalid_type", "Select an AssetManager .ampkg or compatible .zip package.")
|
||||
return RedirectResponse("/software/packages?toast_error=" + quote(message), status_code=303)
|
||||
temporary_path = None
|
||||
try:
|
||||
with tempfile.NamedTemporaryFile(prefix="assetmanager-package-import-", suffix=suffix, delete=False) as handle:
|
||||
temporary_path = Path(handle.name)
|
||||
total_bytes = 0
|
||||
limit_bytes = PACKAGE_IMPORT_MAX_MB * 1024 * 1024
|
||||
while True:
|
||||
chunk = await package_file.read(1024 * 1024)
|
||||
if not chunk:
|
||||
break
|
||||
total_bytes += len(chunk)
|
||||
if total_bytes > limit_bytes:
|
||||
raise ValueError(f"Package bundle exceeds import size limit ({PACKAGE_IMPORT_MAX_MB} MB).")
|
||||
handle.write(chunk)
|
||||
package = import_package_bundle(db, temporary_path, source_name=package_file.filename or "", import_profile=str(import_profile or "").strip().lower() in {"1", "true", "yes", "on"})
|
||||
except Exception as exc:
|
||||
logger.exception("Software package import failed")
|
||||
message = _translate_request(request, "software_packages.import_failed", "Software package import failed: {error}", error=str(exc))
|
||||
return RedirectResponse("/software/packages?toast_error=" + quote(message), status_code=303)
|
||||
finally:
|
||||
await package_file.close()
|
||||
if temporary_path:
|
||||
temporary_path.unlink(missing_ok=True)
|
||||
message = _translate_request(request, "software_packages.imported", 'Software package "{package}" was imported.', package=package.name)
|
||||
return RedirectResponse(f"/software/packages/{package.id}?toast_success=" + quote(message), status_code=303)
|
||||
|
||||
|
||||
@app.get("/software/packages/{package_id}/export")
|
||||
def software_package_export(package_id: int, request: Request, db: Session = Depends(get_db)):
|
||||
_require_admin(request)
|
||||
package = db.get(SoftwarePackage, package_id)
|
||||
if not package or package.package_type != "deployment":
|
||||
raise HTTPException(404, _translate_request(request, "software_packages.not_found", "Software package not found."))
|
||||
try:
|
||||
data = export_package_bundle(package.id, APP_VERSION)
|
||||
manifest = load_package_manifest(package.id)
|
||||
except Exception as exc:
|
||||
raise HTTPException(500, _translate_request(request, "software_packages.export_failed", "Software package export failed: {error}", error=str(exc))) from exc
|
||||
base = re.sub(r"[^A-Za-z0-9._+-]+", "-", str(manifest.get("name") or package.name)).strip("-.") or f"package-{package.id}"
|
||||
version = re.sub(r"[^A-Za-z0-9._+-]+", "-", str(manifest.get("version") or "")).strip("-.")
|
||||
filename = f"{base}-{version}.ampkg" if version else f"{base}.ampkg"
|
||||
return StreamingResponse(io.BytesIO(data), media_type="application/zip", headers={"Content-Disposition": f'attachment; filename="{filename}"'})
|
||||
|
||||
|
||||
@app.get("/software/packages/{package_id}")
|
||||
def software_package_page(package_id: int, request: Request, db: Session = Depends(get_db)):
|
||||
_require_admin(request)
|
||||
@@ -6078,6 +6263,7 @@ async def software_package_delete(package_id: int, request: Request, db: Session
|
||||
|
||||
form = await request.form()
|
||||
delete_jobs = str(form.get("delete_jobs") or "").strip().lower() in {"1", "true", "yes", "on"}
|
||||
return_to = str(form.get("return_to") or "detail").strip().lower()
|
||||
job_ids = [row[0] for row in db.query(SoftwareJob.id).filter(SoftwareJob.package_id == package.id).all()]
|
||||
if job_ids and not delete_jobs:
|
||||
message = _translate_request(
|
||||
@@ -6086,8 +6272,9 @@ async def software_package_delete(package_id: int, request: Request, db: Session
|
||||
"This package is referenced by {count} software jobs. Confirm deletion of the associated jobs first.",
|
||||
count=len(job_ids),
|
||||
)
|
||||
redirect_path = "/software/packages" if return_to == "overview" else f"/software/packages/{package_id}"
|
||||
return RedirectResponse(
|
||||
f"/software/packages/{package_id}?toast_error=" + quote(message),
|
||||
redirect_path + "?toast_error=" + quote(message),
|
||||
status_code=303,
|
||||
)
|
||||
|
||||
@@ -8781,6 +8968,8 @@ def settings_software_save(
|
||||
automatic_retry_max_age_hours: int = Form(12),
|
||||
automatic_retry_interval_seconds: int = Form(60),
|
||||
automatic_retry_max_retries: int = Form(3),
|
||||
remote_job_cleanup_enabled: str | None = Form(None),
|
||||
remote_job_retention_hours: int = Form(24),
|
||||
inventory_exclusion_platform: list[str] = Form(default=[]),
|
||||
inventory_exclusion_name: list[str] = Form(default=[]),
|
||||
):
|
||||
@@ -8814,11 +9003,13 @@ def settings_software_save(
|
||||
"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)),
|
||||
"remote_job_cleanup_enabled": remote_job_cleanup_enabled == "on",
|
||||
"remote_job_retention_hours": max(1, min(int(remote_job_retention_hours), 8760)),
|
||||
"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']}",
|
||||
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']} | remote_job_cleanup_enabled={software['remote_job_cleanup_enabled']} | remote_job_retention_hours={software['remote_job_retention_hours']} | inventory_exclusion_rules={len(submitted_exclusion_rules)} | verify_tls={software['callback_test_verify_tls']}",
|
||||
force=True,
|
||||
)
|
||||
return RedirectResponse(
|
||||
|
||||
Reference in New Issue
Block a user