from pathlib import Path import base64 import hashlib import os import subprocess import tempfile import threading import time from datetime import datetime from urllib.parse import quote from fastapi import HTTPException from .config import load_config from .database import SessionLocal from .models import Asset, SoftwareJob, JobEvent, JobDefinition from .job_state import sync_asset_job_state from .software_packages import build_deployment_user_script, load_package_manifest, package_payload_files _DISPATCH_SLOT_LOCK = threading.Lock() _NEXT_DISPATCH_MONOTONIC = 0.0 _SOFTWARE_INVENTORY_MESHCTRL_LOCK = threading.Lock() def _wait_for_dispatch_slot() -> tuple[float, float]: """Reserve the next global MeshCtrl dispatch slot. Bulk operations currently create one lightweight worker thread per job. Reserving slots here prevents all workers from opening database and MeshCtrl connections at the same time while still allowing already dispatched client scripts to run in parallel. """ global _NEXT_DISPATCH_MONOTONIC settings = load_config().get("software", {}) try: delay_seconds = float(settings.get("dispatch_delay_seconds", 2) or 0) except (TypeError, ValueError): delay_seconds = 2.0 delay_seconds = max(0.0, min(delay_seconds, 60.0)) if delay_seconds <= 0: return 0.0, 0.0 with _DISPATCH_SLOT_LOCK: now = time.monotonic() reserved_at = max(now, _NEXT_DISPATCH_MONOTONIC) wait_seconds = max(0.0, reserved_at - now) _NEXT_DISPATCH_MONOTONIC = reserved_at + delay_seconds if wait_seconds > 0: time.sleep(wait_seconds) return delay_seconds, wait_seconds def token_hash(token: str) -> str: return hashlib.sha256(token.encode("utf-8")).hexdigest() def detect_platform(asset: Asset) -> str: text = f"{asset.operating_system or ''} {asset.mesh_source_data or ''}".casefold() if "windows" in text: return "windows" if "macos" in text or "mac os" in text or "darwin" in text: return "macos" if "android" in text: return "android" if any(x in text for x in ("linux", "ubuntu", "debian", "fedora", "centos", "red hat", "suse", "arch")): return "linux" return "unknown" def powershell_payload(script: str) -> str: """Compatibility helper retained for older callers. Since 0.5.5.0 scripts are uploaded as files and are no longer encoded into the command line. """ return str(script or "") def powershell_encoded(script: str) -> str: return powershell_payload(script) def windows_inventory(callback_url: str, job_id: int) -> str: script = """ $ErrorActionPreference='Stop' $log=Join-Path $env:ProgramData 'AssetManager\\commands.log' New-Item -ItemType Directory -Force -Path (Split-Path $log)|Out-Null $lines=New-Object System.Collections.Generic.List[string] function L([string]$m){$x="$(Get-Date -Format o) job=JOBID $m";Add-Content -Path $log -Value $x -Encoding UTF8;$lines.Add($x)} function Send-Callback([byte[]]$Payload){ $lastError=$null for($attempt=1;$attempt -le 3;$attempt++){ try{ Invoke-RestMethod -Method Post -Uri 'CALLBACK' -ContentType 'application/json; charset=utf-8' -Body $Payload -TimeoutSec 60|Out-Null if($attempt -gt 1){L ("callback accepted on attempt $attempt of 3")} return }catch{ $lastError=$_.Exception L ("callback attempt $attempt of 3 failed: "+$_.Exception.Message) if($attempt -lt 3){Start-Sleep -Seconds (2*$attempt)} } } throw ("callback failed after 3 attempts: "+$lastError.Message) } L 'software inventory started' try { $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 L ("software inventory completed; entries="+$software.Count) $body=@{status='success';exit_code=0;completed=$true;software=$software;log=($lines -join "`n")}|ConvertTo-Json -Depth 6 -Compress $bodyBytes=[System.Text.Encoding]::UTF8.GetBytes($body) L 'sending callback' Send-Callback $bodyBytes L 'callback accepted' }catch{ L ("software inventory failed: "+$_.Exception.Message) $body=@{status='failed';exit_code=1;completed=$true;message=$_.Exception.Message;log=($lines -join "`n")}|ConvertTo-Json -Depth 4 -Compress $bodyBytes=[System.Text.Encoding]::UTF8.GetBytes($body) try{ Send-Callback $bodyBytes }catch{ L ("callback failed: "+$_.Exception.Message) } exit 1 } """.replace("JOBID", str(job_id)).replace("CALLBACK", callback_url) return script def unix_python_payload(platform: str, callback_url: str, job_id: int) -> str: if platform == "macos": collection = """ raw=subprocess.check_output(["system_profiler","SPApplicationsDataType","-json"],text=True,errors="replace") data=json.loads(raw) for x in data.get("SPApplicationsDataType",[]): items.append({"name":x.get("_name",""),"version":x.get("version",""),"publisher":str(x.get("signed_by","")),"architecture":x.get("arch_kind",""),"install_date":None}) """ log_path = "/Library/Logs/AssetManager-commands.log" else: collection = """ if os.path.exists("/usr/bin/dpkg-query"): out=subprocess.check_output(["dpkg-query","-W","-f=${binary:Package}\\t${Version}\\t${Maintainer}\\n"],text=True,errors="replace") for line in out.splitlines(): p=line.split("\\t");items.append({"name":p[0],"version":p[1] if len(p)>1 else "","publisher":p[2] if len(p)>2 else "","architecture":"","install_date":None}) elif os.path.exists("/usr/bin/rpm"): out=subprocess.check_output(["rpm","-qa","--qf","%{NAME}\\t%{VERSION}-%{RELEASE}\\t%{VENDOR}\\t%{ARCH}\\n"],text=True,errors="replace") for line in out.splitlines(): p=line.split("\\t");items.append({"name":p[0],"version":p[1] if len(p)>1 else "","publisher":p[2] if len(p)>2 else "","architecture":p[3] if len(p)>3 else "","install_date":None}) elif os.path.exists("/sbin/apk") or os.path.exists("/usr/sbin/apk"): out=subprocess.check_output(["apk","info","-v"],text=True,errors="replace") for line in out.splitlines():items.append({"name":line,"version":"","publisher":"","architecture":"","install_date":None}) """ log_path = "/var/log/assetmanager-commands.log" py = f"""import json,subprocess,urllib.request,datetime,os,time log={log_path!r} lines=[] def add(m): x=datetime.datetime.now(datetime.timezone.utc).isoformat()+" job={job_id} "+m lines.append(x) try: with open(log,"a",encoding="utf-8") as f:f.write(x+"\\n") except Exception:pass add("software inventory started") items=[] try: {textwrap_indent(collection, 4)} add("software inventory completed; entries="+str(len(items))) payload={{"status":"success","exit_code":0,"completed":True,"software":items,"log":"\\n".join(lines)}} except Exception as e: add("software inventory failed: "+str(e)) payload={{"status":"failed","exit_code":1,"completed":True,"message":str(e),"log":"\\n".join(lines)}} callback_data=json.dumps(payload).encode() last_error=None for attempt in range(1,4): req=urllib.request.Request({callback_url!r},data=callback_data,headers={{"Content-Type":"application/json"}},method="POST") try: urllib.request.urlopen(req,timeout=60).read() break except Exception as exc: last_error=exc add("callback attempt "+str(attempt)+" of 3 failed: "+str(exc)) if attempt < 3: time.sleep(2*attempt) else: raise RuntimeError("callback failed after 3 attempts: "+str(last_error)) """ return py def textwrap_indent(value: str, spaces: int) -> str: prefix = " " * spaces return "\n".join(prefix + line if line.strip() else line for line in value.strip().splitlines()) def _replace_job_placeholders(script: str, asset: Asset, callback_url: str, job_id: int, parameters: dict | None = None) -> str: values={ 'AssetName':asset.name or '', 'Hostname':asset.hostname or '', 'IPAddress':asset.ip_address or '', 'SerialNumber':asset.serial_number or '', 'MeshNodeId':asset.mesh_node_id or '', 'Manufacturer':asset.manufacturer or '', 'Model':asset.model or '', 'Department':asset.department or '', 'Location':asset.location or '', 'CurrentUser':asset.assigned_to or '', 'CallbackUrl':callback_url, 'JobId':str(job_id), } for key,value in values.items(): script=script.replace('{{'+key+'}}',str(value)) for key,value in (parameters or {}).items(): script=script.replace('{{Param.'+str(key)+'}}',str(value)).replace('{{'+str(key)+'}}',str(value)) return script def build_registry_user_script(configuration: dict | None) -> str: """Build PowerShell user code for a declarative Windows Registry job.""" config=configuration or {} entries=config.get("registry_entries") if isinstance(config,dict) else [] backup=bool(config.get("registry_backup")) if isinstance(config,dict) else False hive_map={ "HKLM":"Registry::HKEY_LOCAL_MACHINE", "HKCU":"Registry::HKEY_CURRENT_USER", "HKCR":"Registry::HKEY_CLASSES_ROOT", "HKU":"Registry::HKEY_USERS", "HKCC":"Registry::HKEY_CURRENT_CONFIG", } type_map={"REG_SZ":"String","REG_EXPAND_SZ":"ExpandString","REG_MULTI_SZ":"MultiString","REG_DWORD":"DWord","REG_QWORD":"QWord","REG_BINARY":"Binary"} lines=["Write-JobLog 'Registry job started'", "$registryResults = @()"] if backup: lines += [ "$backupRoot = Join-Path $env:ProgramData ('AssetManager\\RegistryBackups\\Job-' + '{{JobId}}')", "New-Item -ItemType Directory -Force -Path $backupRoot | Out-Null", "Write-JobLog (\"Registry backup directory: $backupRoot\")", ] for index,item in enumerate(entries if isinstance(entries,list) else [],1): if not isinstance(item,dict): continue action=str(item.get('action') or 'set') hive=str(item.get('hive') or 'HKLM').upper() path=str(item.get('path') or '').strip().strip('\\') name=str(item.get('name') or '') value_type=str(item.get('value_type') or 'REG_SZ').upper() value=str(item.get('value') or '') set_value=bool(item.get('set_value', True)) only_if_missing=bool(item.get('only_if_missing')) permissions_enabled=bool(item.get('permissions_enabled')) if hive not in hive_map or not path: continue ps_path=hive_map[hive]+'\\'+path reg_path=hive+'\\'+path q=lambda x: "'"+x.replace("'","''")+"'" lines.append(f"# Registry entry {index}") lines.append(f"$registryPath = {q(ps_path)}") lines.append(f"$registryValueName = {q(name)}") if backup: safe_name=str(index).zfill(3)+'-'+hive+'-'+path.replace('\\','_').replace(':','_')+'.reg' lines += [ f"$backupFile = Join-Path $backupRoot {q(safe_name)}", f"& reg.exe export {q(reg_path)} $backupFile /y 2>$null | Out-Null", ] if action=='delete_key': lines.append("if (Test-Path -LiteralPath $registryPath) { Remove-Item -LiteralPath $registryPath -Recurse -Force -ErrorAction Stop; Write-JobLog (\"Registry key deleted: $registryPath\"); $registryResults += @{path=$registryPath; action='delete_key'; status='changed'} } else { Write-JobLog (\"Registry key already absent: $registryPath\"); $registryResults += @{path=$registryPath; action='delete_key'; status='unchanged'} }") continue if action=='delete_value': lines += [ "$existing = Get-ItemProperty -LiteralPath $registryPath -Name $registryValueName -ErrorAction SilentlyContinue", "if ($null -ne $existing) { Remove-ItemProperty -LiteralPath $registryPath -Name $registryValueName -Force -ErrorAction Stop; Write-JobLog (\"Registry value deleted: $registryPath\\$registryValueName\"); $registryResults += @{path=$registryPath; name=$registryValueName; action='delete_value'; status='changed'} } else { Write-JobLog (\"Registry value already absent: $registryPath\\$registryValueName\"); $registryResults += @{path=$registryPath; name=$registryValueName; action='delete_value'; status='unchanged'} }" ] elif action=='set': lines.append("if (-not (Test-Path -LiteralPath $registryPath)) { New-Item -Path $registryPath -Force | Out-Null; Write-JobLog (\"Registry key created: $registryPath\") }") if set_value: property_type=type_map.get(value_type,'String') lines.append(f"$registryRawValue = {q(value)}") if value_type=='REG_MULTI_SZ': lines.append("$registryTypedValue = @($registryRawValue -split \"`r?`n\")") elif value_type=='REG_BINARY': lines.append("$registryTypedValue = [byte[]](($registryRawValue -replace '[^0-9A-Fa-f]','' -split '(?<=\\G..)(?=.)') | Where-Object { $_ } | ForEach-Object { [Convert]::ToByte($_,16) })") elif value_type=='REG_DWORD': lines.append("$registryTypedValue = if ($registryRawValue -match '^0x') {[Convert]::ToUInt32($registryRawValue.Substring(2),16)} else {[Convert]::ToUInt32($registryRawValue)}") elif value_type=='REG_QWORD': lines.append("$registryTypedValue = if ($registryRawValue -match '^0x') {[Convert]::ToUInt64($registryRawValue.Substring(2),16)} else {[Convert]::ToUInt64($registryRawValue)}") else: lines.append("$registryTypedValue = $registryRawValue") lines.append("$existingProperty = Get-ItemProperty -LiteralPath $registryPath -Name $registryValueName -ErrorAction SilentlyContinue") if only_if_missing: lines.append("if ($null -ne $existingProperty) { Write-JobLog (\"Registry value exists; skipped: $registryPath\\$registryValueName\"); $registryResults += @{path=$registryPath; name=$registryValueName; action='set'; status='skipped'} } else {") indent=' ' else: indent='' lines += [ indent+f"New-ItemProperty -LiteralPath $registryPath -Name $registryValueName -PropertyType {property_type} -Value $registryTypedValue -Force -ErrorAction Stop | Out-Null", indent+"Write-JobLog (\"Registry value set: $registryPath\\$registryValueName\")", indent+"$registryResults += @{path=$registryPath; name=$registryValueName; action='set'; status='changed'; type='"+value_type+"'}", ] if only_if_missing: lines.append("}") else: lines.append("Write-JobLog (\"Registry key ensured without changing a value: $registryPath\")") lines.append("$registryResults += @{path=$registryPath; action='ensure_key'; status='ready'}") if 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)) lines += [ "if (-not (Test-Path -LiteralPath $registryPath)) { New-Item -Path $registryPath -Force | Out-Null; Write-JobLog (\"Registry key created for permissions: $registryPath\") }", f"$registryPrincipal = {q(principal)}", f"$registryRightsName = {q(rights)}", f"$registryAccessTypeName = {q(access_type)}", "try { $registryIdentity = New-Object System.Security.Principal.SecurityIdentifier($registryPrincipal) } catch { $registryIdentity = $registryPrincipal }", "try { $registryRights = [System.Enum]::Parse([System.Security.AccessControl.RegistryRights], $registryRightsName, $true) } catch { throw \"Invalid RegistryRights '$registryRightsName' for $registryPath\" }", "$registryAccessType = [System.Enum]::Parse([System.Security.AccessControl.AccessControlType], $registryAccessTypeName, $true)", "$registryInheritance = " + ("[System.Security.AccessControl.InheritanceFlags]::ContainerInherit" if inherit_subkeys else "[System.Security.AccessControl.InheritanceFlags]::None"), "$registryPropagation = [System.Security.AccessControl.PropagationFlags]::None", "Write-JobLog (\"Registry ACL target: $registryPath; exists=$([bool](Test-Path -Path $registryPath))\")", "$registryAcl = Get-Acl -Path $registryPath -ErrorAction Stop", "$registryRule = New-Object System.Security.AccessControl.RegistryAccessRule($registryIdentity, $registryRights, $registryInheritance, $registryPropagation, $registryAccessType)", ("$registryAcl.SetAccessRule($registryRule)" if replace_existing else "$registryAcl.AddAccessRule($registryRule) | Out-Null"), "Set-Acl -Path $registryPath -AclObject $registryAcl -ErrorAction Stop", "$registryAclAfter = Get-Acl -Path $registryPath -ErrorAction Stop", "$registryMatchedRules = @($registryAclAfter.Access | Where-Object { try { $_.IdentityReference.Translate([System.Security.Principal.SecurityIdentifier]).Value -eq $registryIdentity.Value } catch { $_.IdentityReference.Value -eq $registryPrincipal -or $_.IdentityReference.Value -like ('*\\' + $registryPrincipal) } })", "Write-JobLog (\"Registry permission set: $registryPath; identity=$registryPrincipal; rights=$registryRightsName; type=$registryAccessTypeName; matching_rules=$($registryMatchedRules.Count)\")", "$registryResults += @{path=$registryPath; action='set_permissions'; status='changed'; principal=$registryPrincipal; rights=$registryRightsName; access_type=$registryAccessTypeName; matching_rules=$registryMatchedRules.Count}", ] lines += ["$JobResult['registry_entries'] = $registryResults", "Write-JobLog ('Registry job completed; entries=' + $registryResults.Count)"] return "\n".join(lines) def _runtime_wrap_script(interpreter: str, user_script: str) -> str: # Job definitions contain user code only. The technical execution wrapper # is always generated from the current interpreter template at runtime. source=str(user_script or "") root=Path(__file__).resolve().parent / "templates" / "jobs" files={"powershell":"powershell_wrapper.ps1","python":"python_wrapper.py.txt","bash":"bash_wrapper.sh","cmd":"cmd_wrapper.cmd"} 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 ---")} filename=files.get(interpreter) if not filename: raise RuntimeError(f"Unbekannter Interpreter: {interpreter}") wrapper=(root/filename).read_text(encoding="utf-8").replace("\r\n","\n") start,end=markers[interpreter]; block=start+"\n\n"+end indent={"powershell":" ","python":" ","bash":"","cmd":""}[interpreter] commands="\n".join(indent+line if line.strip() else line for line in str(user_script or "").splitlines()) return wrapper.replace(block,start+"\n"+commands+"\n"+end,1) def _definition_execution_artifacts(db, platform: str, callback_url: str, job_id: int, asset: Asset, definition_id: int | None = None) -> tuple[str, str, bool]: """Return (resolved script/command, interpreter, upload_required).""" if definition_id is not None: definition=db.query(JobDefinition).filter(JobDefinition.id==definition_id,JobDefinition.enabled.is_(True)).first() else: definition=db.query(JobDefinition).filter(JobDefinition.system_key==f'software_inventory_{platform}',JobDefinition.enabled.is_(True)).first() if not definition: if definition_id is not None: raise RuntimeError('Jobdefinition wurde nicht gefunden oder ist deaktiviert.') return '', '', False if definition.platform not in {platform, 'all'}: raise RuntimeError(f'Jobdefinition passt nicht zur Plattform {platform}.') effective_parameters=((db.get(SoftwareJob,job_id).parameters or {}).get('effective_parameters') or {}) if db.get(SoftwareJob,job_id) else {} if definition.job_kind=='registry': script=_runtime_wrap_script('powershell', build_registry_user_script(definition.configuration)) interpreter='powershell' elif definition.source_type=='inline': interpreter=definition.interpreter script=_runtime_wrap_script(interpreter, definition.inline_script or '') elif definition.source_type=='mounted': interpreter=definition.interpreter path=Path(definition.source_path or '') if not path.is_file(): raise RuntimeError(f'Skriptdatei nicht gefunden: {path}') script=_runtime_wrap_script(interpreter, path.read_text(encoding='utf-8',errors='replace')) elif definition.source_type=='unc': path=_replace_job_placeholders(definition.source_path or '',asset,callback_url,job_id,effective_parameters) arguments=_replace_job_placeholders(definition.arguments or '',asset,callback_url,job_id,effective_parameters) if definition.interpreter=='powershell': command=f'powershell.exe -NoProfile -NonInteractive -ExecutionPolicy Bypass -File "{path}" {arguments}'.strip() elif definition.interpreter=='cmd': command=f'cmd.exe /d /s /c ""{path}" {arguments}"'.strip() elif definition.interpreter=='python': command=f'python3 "{path}" {arguments}'.strip() else: command=f'bash "{path}" {arguments}'.strip() return command, definition.interpreter, False else: raise RuntimeError(f'Unbekannte Skriptquelle: {definition.source_type}') script=_replace_job_placeholders(script,asset,callback_url,job_id,effective_parameters) return script, interpreter, True def _job_execution_artifacts(db, platform: str, callback_url: str, job_id: int, asset: Asset, definition_id: int | None = None) -> tuple[str, str, bool]: job = db.get(SoftwareJob, job_id) if job and job.job_type == 'software_deployment': if platform != 'windows': raise RuntimeError('Softwarepakete werden derzeit nur fuer Windows-Assets unterstuetzt.') manifest = load_package_manifest(job.package_id) user_script = build_deployment_user_script(manifest, job.action) script = _runtime_wrap_script('powershell', user_script) script = _replace_job_placeholders( script, asset, callback_url, job_id, job.parameters if isinstance(job.parameters, dict) else {}, ) return script, 'powershell', True configured=_definition_execution_artifacts(db,platform,callback_url,job_id,asset,definition_id) if configured[0]: return configured if platform == 'windows': return windows_inventory(callback_url, job_id), 'powershell', True if platform in {'linux','macos'}: return unix_python_payload(platform, callback_url, job_id), 'python', True raise HTTPException(400, 'Fuer dieses Betriebssystem ist kein sicherer MeshCtrl-Handler verfuegbar.') def _script_extension(interpreter: str) -> str: return {'powershell':'.ps1','cmd':'.cmd','python':'.py','bash':'.sh'}.get(interpreter,'.txt') def _remote_job_paths(platform: str, job_id: int, attempt: int, interpreter: str) -> tuple[str, str]: suffix=f'{job_id}-{attempt}' filename='run'+_script_extension(interpreter) if platform == 'windows': directory=f'C:\\ProgramData\\AssetManager\\Jobs\\{suffix}' return directory, directory+'\\'+filename directory=f'/var/lib/assetmanager/jobs/{suffix}' return directory, directory+'/'+filename def _meshctrl_base_args(cfg: dict, asset: Asset, password: str) -> list[str]: args=['node',str(cfg.get('meshctrl_path') or '/opt/meshcentral/node_modules/meshcentral/meshctrl.js')] args += ['--url',str(cfg.get('url') or ''),'--loginuser',str(cfg.get('username') or ''),'--loginpass',password] tenant=str(cfg.get('tenant') or '').strip() if tenant: args += ['--domain',tenant] return args def _run_meshctrl(cfg: dict, asset: Asset, password: str, action_args: list[str], timeout: int) -> subprocess.CompletedProcess: base=_meshctrl_base_args(cfg,asset,password) args=base[:2]+action_args+base[2:] return subprocess.run(args,capture_output=True,text=True,timeout=timeout,check=False) def _prepare_remote_directory( cfg: dict, asset: Asset, password: str, platform: str, remote_dir: str, timeout: int, remote_file: str | None = None, reset_acl: bool = True, ) -> subprocess.CompletedProcess: """Create the job directory and optionally remove one stale target file. ACLs are applied to the directory only. Never use a recursive icacls /T here: directory inheritance flags such as (OI)(CI) must not be rewritten onto already uploaded files. Doing so can leave files with no effective ACEs and make them unreadable even for the SYSTEM account. """ if platform == 'windows': command=( "$p='"+remote_dir.replace("'","''")+"';" "New-Item -ItemType Directory -Force -Path $p|Out-Null;" ) if remote_file: command += "Remove-Item -LiteralPath '"+remote_file.replace("'","''")+"' -Force -ErrorAction SilentlyContinue;" if reset_acl: command += ( "& icacls.exe $p /inheritance:r " "/grant:r '*S-1-5-18:(OI)(CI)F' '*S-1-5-32-544:(OI)(CI)F' " "/C|Out-Null;" ) return _run_meshctrl(cfg,asset,password,['RunCommand','--id',asset.mesh_node_id,'--run',command,'--powershell','--reply'],timeout) command=f"mkdir -p '{remote_dir}'" if remote_file: command += f" && rm -f -- '{remote_file}'" if reset_acl: command += f" && chmod 700 '{remote_dir}'" return _run_meshctrl(cfg,asset,password,['RunCommand','--id',asset.mesh_node_id,'--run',command,'--reply'],timeout) def _remote_script_preflight( cfg: dict, asset: Asset, password: str, platform: str, remote_file: str, timeout: int, ) -> subprocess.CompletedProcess: """Verify that the uploaded job script exists and is readable. The Windows diagnostic also prints the effective ACL so an upload/ACL problem is visible in the dispatcher log before PowerShell is launched. """ if platform == 'windows': escaped=remote_file.replace("'","''") command=( "$f='"+escaped+"';" "$exists=Test-Path -LiteralPath $f -PathType Leaf;" "Write-Output ('File exists: '+$exists);" "if($exists){" "try{$i=Get-Item -LiteralPath $f -ErrorAction Stop;Write-Output ('Size: '+$i.Length)}" "catch{Write-Output ('Size: ERROR - '+$_.Exception.Message)};" "try{$s=[System.IO.File]::Open($f,[System.IO.FileMode]::Open,[System.IO.FileAccess]::Read,[System.IO.FileShare]::ReadWrite);$s.Close();Write-Output 'Readable: True'}" "catch{Write-Output ('Readable: False - '+$_.Exception.Message)};" "Write-Output 'ACL:'; & icacls.exe $f" "}" ) return _run_meshctrl(cfg,asset,password,['RunCommand','--id',asset.mesh_node_id,'--run',command,'--powershell','--reply'],timeout) escaped=remote_file.replace("'","'\''") command=f"test -f '{escaped}' && test -r '{escaped}' && ls -l '{escaped}'" return _run_meshctrl(cfg,asset,password,['RunCommand','--id',asset.mesh_node_id,'--run',command,'--reply'],timeout) def _upload_script(cfg: dict, asset: Asset, password: str, local_file: str, remote_dir: str, timeout: int) -> subprocess.CompletedProcess: return _run_meshctrl(cfg,asset,password,['Upload','--id',asset.mesh_node_id,'--file',local_file,'--target',remote_dir],timeout) def _uploaded_script_launch_spec(platform: str, interpreter: str, remote_file: str) -> tuple[str, str, list[str]]: """Return the exact remote command, shell label and MeshCtrl shell flags. Windows PowerShell files are deliberately launched through MeshCtrl's Windows Command Prompt mode (type 0). The command then starts powershell itself with ExecutionPolicy Bypass. Adding MeshCtrl's --powershell flag would switch to type 2 and cause the agent to execute the .ps1 inside the PowerShell host, which is precisely what must be avoided here. """ if platform == 'windows' and interpreter == 'powershell': # Important: do not add MeshCtrl's --powershell option here. Without # that option MeshCtrl uses the normal Windows command prompt, which # then starts PowerShell with the same command that was tested manually. escaped_file = remote_file.replace('"', '""') return ( f'powershell -ExecutionPolicy Bypass -File "{escaped_file}"', 'Windows Command Prompt', [], ) if platform == 'windows' and interpreter == 'cmd': return f'cmd.exe /d /s /c ""{remote_file}""', 'Windows Command Prompt', [] if interpreter == 'python': return f"python3 '{remote_file}'", 'System shell', [] return f"chmod 700 '{remote_file}' && bash '{remote_file}'", 'System shell', [] def _launch_uploaded_script(cfg: dict, asset: Asset, password: str, platform: str, interpreter: str, remote_file: str, timeout: int) -> subprocess.CompletedProcess: command,_shell_label,meshctrl_shell_flags=_uploaded_script_launch_spec(platform,interpreter,remote_file) action=['RunCommand','--id',asset.mesh_node_id,'--run',command]+meshctrl_shell_flags+['--reply'] if platform == 'windows' and interpreter == 'powershell' and '--powershell' in action: raise RuntimeError('Interner Dispatcherfehler: Dateibasierte PowerShell-Jobs dürfen MeshCtrl nicht im PowerShell-Modus starten.') return _run_meshctrl(cfg,asset,password,action,timeout) def _cleanup_remote_directory(cfg: dict, asset: Asset, password: str, platform: str, remote_dir: str, timeout: int) -> subprocess.CompletedProcess: if platform == 'windows': command="Remove-Item -LiteralPath '"+remote_dir.replace("'","''")+"' -Recurse -Force -ErrorAction SilentlyContinue" action=['RunCommand','--id',asset.mesh_node_id,'--run',command,'--powershell','--reply'] else: action=['RunCommand','--id',asset.mesh_node_id,'--run',f"rm -rf -- '{remote_dir}'",'--reply'] return _run_meshctrl(cfg,asset,password,action,timeout) def cleanup_stale_remote_job_directories( cfg: dict, asset: Asset, password: str, platform: str, retention_hours: int, timeout: int, exclude_directory_names: list[str] | None = None, ) -> subprocess.CompletedProcess: """Remove stale AssetManager job-attempt directories on the client. Windows age is based on directory CreationTimeUtc, not LastWriteTimeUtc. That represents the age of the job workspace and is not extended when a script or installer later touches files in that directory. Only numeric AssetManager job directories (legacy ```` and current ``-``) are considered. Active/current directory names supplied by the caller are excluded. Deletion is best-effort. """ try: retention = max(1, min(int(retention_hours), 8760)) except (TypeError, ValueError): retention = 24 excluded = sorted({str(name or '').strip() for name in (exclude_directory_names or []) if str(name or '').strip()}) if platform == 'windows': root = r'C:\ProgramData\AssetManager\Jobs' escaped_root = root.replace("'", "''") excluded_ps = ','.join("'" + name.replace("'", "''") + "'" for name in excluded) command = ( f"$root='{escaped_root}';" f"$cutoff=[DateTime]::UtcNow.AddHours(-{retention});" f"$excluded=@({excluded_ps});" "if(Test-Path -LiteralPath $root){" "$found=0;$eligible=0;$removed=0;$failed=0;$young=0;$active=0;$ignored=0;" "Get-ChildItem -LiteralPath $root -Directory -Force -ErrorAction SilentlyContinue|ForEach-Object{" "$item=$_;$found++;" "if($item.Name -notmatch '^[0-9]+(?:-[0-9]+)?$'){$ignored++;return};" "if($excluded -contains $item.Name){$active++;Write-Output ('Kept active job directory: '+$item.FullName);return};" "$created=$item.CreationTimeUtc;" "if($created -ge $cutoff){$young++;return};" "$eligible++;$dir=$item.FullName;" "try{Remove-Item -LiteralPath $dir -Recurse -Force -ErrorAction Stop;$removed++;Write-Output ('Removed stale job directory: '+$dir+' created_utc='+$created.ToString('o'))}" "catch{" "$firstError=$_.Exception.Message;" "try{" "& icacls.exe $dir /inheritance:e /grant:r '*S-1-5-18:(OI)(CI)F' '*S-1-5-32-544:(OI)(CI)F' /T /C /Q | Out-Null;" "Remove-Item -LiteralPath $dir -Recurse -Force -ErrorAction Stop;" "$removed++;Write-Output ('Removed stale job directory after ACL repair: '+$dir+' first_error='+$firstError)" "}catch{$failed++;Write-Output ('Failed stale job directory: '+$dir+' - '+$_.Exception.Message+' first_error='+$firstError)}" "}" "};" f"Write-Output ('Stale cleanup summary: found='+$found+' eligible='+$eligible+' removed='+$removed+' failed='+$failed+' young='+$young+' active='+$active+' ignored='+$ignored+' retention_hours={retention}')" "}else{Write-Output 'Stale cleanup summary: job root not present'}" ) action = ['RunCommand','--id',asset.mesh_node_id,'--run',command,'--powershell','--reply'] return _run_meshctrl(cfg,asset,password,action,timeout) root = '/var/lib/assetmanager/jobs' minutes = retention * 60 exclude_tests = ' '.join(f"! -name '{name}'" for name in excluded) command = ( f"root='{root}'; " "if [ -d \"$root\" ]; then " f"find \"$root\" -mindepth 1 -maxdepth 1 -type d -mmin +{minutes} \\( -name '[0-9]*' -o -name '[0-9]*-[0-9]*' \\) {exclude_tests} -print -exec rm -rf -- {{}} \\;; " f"echo 'Stale cleanup completed; retention_hours={retention}'; " "else echo 'Stale cleanup summary: job root not present'; fi" ) action = ['RunCommand','--id',asset.mesh_node_id,'--run',command,'--reply'] return _run_meshctrl(cfg,asset,password,action,timeout) def _cleanup_stale_remote_job_directories( cfg: dict, asset: Asset, password: str, platform: str, current_remote_dir: str, retention_hours: int, timeout: int, ) -> subprocess.CompletedProcess: """Compatibility wrapper for dispatcher-side cleanup.""" current_name = Path(current_remote_dir.replace('\\', '/')).name return cleanup_stale_remote_job_directories( cfg, asset, password, platform, retention_hours, timeout, [current_name] if current_name else [] ) def _add_job_event(db, job: SoftwareJob, event_type: str, status: str | None = None, message: str | None = None) -> None: db.add(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, )) def execute_job(job_id: int, token: str, callback_base: str) -> None: # Reserve the dispatch slot before opening a database session. With large # bulk selections this keeps waiting workers from occupying the connection # pool while they wait for their configured MeshCtrl start interval. dispatch_delay_seconds, dispatch_wait_seconds = _wait_for_dispatch_slot() db = SessionLocal() temp_path: str | None = None temp_dir: str | None = None try: job = db.get(SoftwareJob, job_id) asset = db.get(Asset, job.asset_id) if job else None if not job or not asset or not asset.mesh_node_id: if job: job.status='failed' job.message='Asset besitzt keine MeshCentral Node-ID.' job.finished_at=datetime.utcnow() sync_asset_job_state(db, job) db.commit() return callback_url=f"{callback_base}/api/software-jobs/{job.id}/callback?token={quote(token)}" definition_id=None if job.job_type=='job_definition': try: definition_id=int((job.parameters or {}).get('definition_id')) except (TypeError,ValueError): raise RuntimeError('Jobdefinition fehlt oder ist ungültig.') payload,interpreter,upload_required=_job_execution_artifacts(db,job.platform,callback_url,job.id,asset,definition_id) job.attempt_count=(job.attempt_count or 0)+1 attempt=job.attempt_count remote_dir,remote_file=_remote_job_paths(job.platform,job.id,attempt,interpreter) job.resolved_script=payload launch_command=None launch_shell=None if upload_required: launch_command,launch_shell,_launch_flags=_uploaded_script_launch_spec(job.platform,interpreter,remote_file) job.command_preview=f'UPLOAD -> {remote_file}\nSHELL -> {launch_shell}\nEXECUTE -> {launch_command}' else: job.command_preview=payload job.status='sending' _add_job_event(db,job,'sending','sending','Jobdatei wird über MeshCentral übertragen.') sync_asset_job_state(db, job) db.commit() runtime_config=load_config() cfg=runtime_config.get('meshcentral',{}) software_settings=runtime_config.get('software',{}) remote_cleanup_enabled=bool(software_settings.get('remote_job_cleanup_enabled', True)) try: remote_job_retention_hours=max(1,min(int(software_settings.get('remote_job_retention_hours',24) or 24),8760)) except (TypeError,ValueError): remote_job_retention_hours=24 password_env=str(cfg.get('password_env') or 'MESHCENTRAL_PASSWORD') password=os.getenv(password_env,'') if not password: raise RuntimeError(f'MeshCentral-Passwortvariable {password_env} ist nicht gesetzt.') timeout=max(30,int(cfg.get('timeout_seconds') or 600)) job_parameters = job.parameters if isinstance(job.parameters, dict) else {} try: execution_timeout_seconds = int(job_parameters.get('_execution_timeout_seconds') or 0) except (TypeError, ValueError): execution_timeout_seconds = 0 if execution_timeout_seconds > 0: timeout = max(timeout, min(execution_timeout_seconds + 180, 86400)) diagnostics=[ '--- File dispatcher diagnostics ---', f'Dispatch interval configured: {dispatch_delay_seconds:g} seconds', f'Dispatch queue wait: {dispatch_wait_seconds:.3f} seconds', f'Job type (technical): {job.job_type}', f'Package (technical): {job.package.name if job.package else job.package_id}', f'Creation mode: {job_parameters.get("_creation_mode") or "single"}', f'Bulk batch id: {job_parameters.get("_bulk_batch_id") or "-"}', f'Bulk position: {job_parameters.get("_bulk_position") or "-"}/{job_parameters.get("_bulk_total") or "-"}', f'Mesh node id: {asset.mesh_node_id}', f'Interpreter: {interpreter}', f'Upload required: {upload_required}', f'Remote stale cleanup enabled: {remote_cleanup_enabled}', f'Remote job retention: {remote_job_retention_hours} hours', f'Resolved script characters: {len(payload)}', f'Resolved script SHA-256: {hashlib.sha256(payload.encode("utf-8")).hexdigest()}', ] if job.job_type == 'job_definition': diagnostics += [ f'Definition id: {job_parameters.get("definition_id") or "-"}', f'Definition name: {job_parameters.get("definition_name") or "-"}', f'Definition kind: {job_parameters.get("definition_job_kind") or "-"}', f'Definition revision: {job_parameters.get("definition_revision") or "-"}', f'Definition interpreter: {job_parameters.get("definition_interpreter") or interpreter}', f'Definition source type: {job_parameters.get("definition_source_type") or "-"}', ] elif job.job_type == 'software_deployment': diagnostics += [ f'Deployment action: {job.action}', f'Deployment product: {job_parameters.get("product_name") or "-"}', f'Deployment version: {job_parameters.get("product_version") or "-"}', f'Installer type: {job_parameters.get("installer_type") or "-"}', f'Execution timeout: {execution_timeout_seconds or timeout} seconds', ] inventory_meshctrl_lock_acquired = False if job.job_type == 'software_inventory': lock_started = time.monotonic() _SOFTWARE_INVENTORY_MESHCTRL_LOCK.acquire() inventory_meshctrl_lock_acquired = True diagnostics += [ 'MeshCtrl isolation: software inventory serialized', f'MeshCtrl isolation wait: {time.monotonic() - lock_started:.3f} seconds', ] try: if upload_required: suffix=_script_extension(interpreter) temp_dir=tempfile.mkdtemp(prefix=f'assetmanager-job-{job.id}-') temp_path=str(Path(temp_dir)/('run'+suffix)) Path(temp_path).write_text(payload,encoding='utf-8',newline='\n') diagnostics += [f'Remote directory: {remote_dir}',f'Remote file: {remote_file}',f'Local staging bytes: {os.path.getsize(temp_path)}'] if remote_cleanup_enabled: stale_cleanup=_cleanup_stale_remote_job_directories( cfg,asset,password,job.platform,remote_dir,remote_job_retention_hours,min(timeout,120) ) diagnostics += [ '--- Stale remote job cleanup stdout ---',stale_cleanup.stdout or '', '--- Stale remote job cleanup stderr ---',stale_cleanup.stderr or '', f'Stale remote job cleanup return code: {stale_cleanup.returncode}', ] prepared=_prepare_remote_directory(cfg,asset,password,job.platform,remote_dir,timeout,remote_file) diagnostics += ['--- Prepare directory stdout ---',prepared.stdout or '','--- Prepare directory stderr ---',prepared.stderr or '',f'Prepare return code: {prepared.returncode}'] if prepared.returncode!=0: raise RuntimeError(f'Remote Jobverzeichnis konnte nicht erstellt werden: {prepared.stderr or prepared.stdout}') upload_results=[] uploaded=_upload_script(cfg,asset,password,temp_path,remote_dir,timeout) upload_results.append(uploaded) diagnostics += ['--- Upload attempt 1 stdout ---',uploaded.stdout or '','--- Upload attempt 1 stderr ---',uploaded.stderr or '',f'Upload attempt 1 return code: {uploaded.returncode}'] upload_text=(uploaded.stdout or '')+'\n'+(uploaded.stderr or '') upload_ok=uploaded.returncode==0 and 'Upload done' in upload_text and 'Upload error' not in upload_text if not upload_ok: # MeshCtrl can return code 0 together with "Upload error". # Remove only the stale target and retry once. Do not touch ACLs of already uploaded files. repair=_prepare_remote_directory(cfg,asset,password,job.platform,remote_dir,min(timeout,120),remote_file,reset_acl=False) diagnostics += ['--- Upload repair stdout ---',repair.stdout or '','--- Upload repair stderr ---',repair.stderr or '',f'Upload repair return code: {repair.returncode}'] uploaded=_upload_script(cfg,asset,password,temp_path,remote_dir,timeout) upload_results.append(uploaded) diagnostics += ['--- Upload attempt 2 stdout ---',uploaded.stdout or '','--- Upload attempt 2 stderr ---',uploaded.stderr or '',f'Upload attempt 2 return code: {uploaded.returncode}'] upload_text=(uploaded.stdout or '')+'\n'+(uploaded.stderr or '') upload_ok=uploaded.returncode==0 and 'Upload done' in upload_text and 'Upload error' not in upload_text if not upload_ok: combined='\n\n'.join(((item.stdout or '')+'\n'+(item.stderr or '')).strip() for item in upload_results if ((item.stdout or '')+(item.stderr or '')).strip()) raise RuntimeError(f'MeshCentral-Dateiupload fehlgeschlagen nach 2 Versuchen: {combined.strip()}') if job.job_type == 'software_deployment': payload_files = package_payload_files(job.package_id, job.action) diagnostics += [f'Package payload files: {len(payload_files)}'] for package_file in payload_files: remote_package_file = remote_dir + '\\' + package_file.name package_uploads = [] package_uploaded = _upload_script(cfg, asset, password, str(package_file), remote_dir, timeout) package_uploads.append(package_uploaded) diagnostics += [ f'--- Package upload {package_file.name} attempt 1 stdout ---', package_uploaded.stdout or '', f'--- Package upload {package_file.name} attempt 1 stderr ---', package_uploaded.stderr or '', f'Package upload {package_file.name} attempt 1 return code: {package_uploaded.returncode}', ] package_text = (package_uploaded.stdout or '') + '\n' + (package_uploaded.stderr or '') package_ok = package_uploaded.returncode == 0 and 'Upload done' in package_text and 'Upload error' not in package_text if not package_ok: package_repair = _prepare_remote_directory( cfg, asset, password, job.platform, remote_dir, min(timeout, 120), remote_package_file, reset_acl=False, ) diagnostics += [ f'--- Package upload {package_file.name} repair stdout ---', package_repair.stdout or '', f'--- Package upload {package_file.name} repair stderr ---', package_repair.stderr or '', f'Package upload {package_file.name} repair return code: {package_repair.returncode}', ] package_uploaded = _upload_script(cfg, asset, password, str(package_file), remote_dir, timeout) package_uploads.append(package_uploaded) diagnostics += [ f'--- Package upload {package_file.name} attempt 2 stdout ---', package_uploaded.stdout or '', f'--- Package upload {package_file.name} attempt 2 stderr ---', package_uploaded.stderr or '', f'Package upload {package_file.name} attempt 2 return code: {package_uploaded.returncode}', ] package_text = (package_uploaded.stdout or '') + '\n' + (package_uploaded.stderr or '') package_ok = package_uploaded.returncode == 0 and 'Upload done' in package_text and 'Upload error' not in package_text if not package_ok: combined = '\n\n'.join( ((item.stdout or '') + '\n' + (item.stderr or '')).strip() for item in package_uploads if ((item.stdout or '') + (item.stderr or '')).strip() ) raise RuntimeError( f'MeshCentral-Paketdateiupload fehlgeschlagen fuer {package_file.name}: {combined.strip()}' ) preflight_ok = False preflight_attempts = 3 for preflight_attempt in range(1, preflight_attempts + 1): preflight = _remote_script_preflight( cfg, asset, password, job.platform, remote_file, min(timeout, 120), ) diagnostics += [ f'--- Remote script preflight attempt {preflight_attempt}/{preflight_attempts} stdout ---', preflight.stdout or '', f'--- Remote script preflight attempt {preflight_attempt}/{preflight_attempts} stderr ---', preflight.stderr or '', f'Remote script preflight attempt {preflight_attempt}/{preflight_attempts} return code: {preflight.returncode}', ] preflight_text = (preflight.stdout or '') + '\n' + (preflight.stderr or '') preflight_ok = preflight.returncode == 0 and ( job.platform != 'windows' or ('File exists: True' in preflight_text and 'Readable: True' in preflight_text) ) if preflight_ok: break if preflight_attempt < preflight_attempts: diagnostics += [ f'Remote script preflight attempt {preflight_attempt}/{preflight_attempts}: not readable', 'Retrying remote script preflight in 1 second...', ] time.sleep(1) if not preflight_ok: raise RuntimeError('Remote Jobskript ist nach dem Upload nicht lesbar. Siehe Remote script preflight im Dispatcherlog.') diagnostics += [f'Remote execution shell: {launch_shell}',f'Remote execution command: {launch_command}',f'MeshCtrl PowerShell mode: False'] started=datetime.utcnow() result=_launch_uploaded_script(cfg,asset,password,job.platform,interpreter,remote_file,timeout) finished=datetime.utcnow() diagnostics += [f'Execution started UTC: {started.isoformat()}Z',f'Execution finished UTC: {finished.isoformat()}Z','--- Execution stdout ---',result.stdout or '','--- Execution stderr ---',result.stderr or '',f'Execution return code: {result.returncode}'] else: started=datetime.utcnow() action=['RunCommand','--id',asset.mesh_node_id,'--run',payload,'--reply'] if job.platform=='windows' and interpreter=='powershell': action += ['--powershell'] result=_run_meshctrl(cfg,asset,password,action,timeout) finished=datetime.utcnow() diagnostics += ['--- Direct execution stdout ---',result.stdout or '','--- Direct execution stderr ---',result.stderr or '',f'Execution return code: {result.returncode}'] now=datetime.utcnow() # MeshCtrl itself returns code 0 even for some remote PowerShell parser # errors. Only delete the remote files after a completed successful # callback, never merely because the MeshCtrl process exited cleanly. db.expire_all() callback_state=db.get(SoftwareJob,job.id) cleanup_result=None callback_success=bool(callback_state and callback_state.callback_completed_at is not None and callback_state.status=='success') diagnostics += [f'Callback completed before cleanup: {bool(callback_state and callback_state.callback_completed_at)}',f'Callback status before cleanup: {callback_state.status if callback_state else "unknown"}'] if upload_required and callback_success: cleanup_result=_cleanup_remote_directory(cfg,asset,password,job.platform,remote_dir,min(timeout,120)) diagnostics += ['--- Cleanup stdout ---',cleanup_result.stdout or '','--- Cleanup stderr ---',cleanup_result.stderr or '',f'Cleanup return code: {cleanup_result.returncode}'] elif upload_required: diagnostics += ['Remote cleanup skipped: no successful completed callback; files retained for diagnosis.'] finally: if inventory_meshctrl_lock_acquired: _SOFTWARE_INVENTORY_MESHCTRL_LOCK.release() db.query(SoftwareJob).filter(SoftwareJob.id==job.id).update({SoftwareJob.meshctrl_stdout:'\n'.join(diagnostics)[-30000:],SoftwareJob.meshctrl_stderr:(result.stderr or '')[-20000:],SoftwareJob.sent_at:now},synchronize_session=False) db.commit() terminal_states={'success','failed','partial','timeout','cancelled'} pending=db.query(SoftwareJob).filter(SoftwareJob.id==job.id,SoftwareJob.callback_completed_at.is_(None),SoftwareJob.status.notin_(terminal_states)) if result.returncode==0: message='Skriptdatei ausgeführt; Callback wird erwartet.' if upload_required and not callback_success: message+=' Jobdateien bleiben bis zu einem bestätigten erfolgreichen Callback erhalten.' if upload_required and cleanup_result is not None and cleanup_result.returncode!=0: message+=' Die lokale Jobdatei konnte nicht automatisch gelöscht werden.' pending.update({SoftwareJob.status:'sent',SoftwareJob.message:message},synchronize_session=False) else: pending.update({SoftwareJob.status:'failed',SoftwareJob.message:f'MeshCtrl-Ausführungsfehler {result.returncode}; Jobdateien bleiben zur Diagnose erhalten.',SoftwareJob.finished_at:now},synchronize_session=False) db.expire_all() refreshed=db.get(SoftwareJob,job.id) if refreshed: _add_job_event(db,refreshed,'meshctrl_result',refreshed.status,refreshed.message) sync_asset_job_state(db, refreshed) db.commit() except Exception as exc: job=db.get(SoftwareJob,job_id) if job: job.status='failed' job.message=str(exc) job.finished_at=datetime.utcnow() if 'diagnostics' in locals() and diagnostics: job.meshctrl_stdout='\n'.join(diagnostics)[-30000:] _add_job_event(db,job,'error','failed',str(exc)) sync_asset_job_state(db, job) db.commit() finally: if temp_path: try: os.unlink(temp_path) except OSError: pass if temp_dir: try: os.rmdir(temp_dir) except OSError: pass db.close()