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 _DISPATCH_SLOT_LOCK = threading.Lock() _NEXT_DISPATCH_MONOTONIC = 0.0 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]: 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, 'Für dieses Betriebssystem ist kein sicherer MeshCtrl-Handler verfügbar.') 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, ) -> subprocess.CompletedProcess: """Create the job directory and remove a stale target file. MeshCtrl's Upload action does not overwrite an existing file reliably on all Windows agents. Failed jobs intentionally retain their directories, so a repeated dispatch must explicitly remove a previous run.ps1 before the upload starts. """ 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;" command += "& icacls.exe $p /inheritance:r /grant:r '*S-1-5-18:(OI)(CI)F' '*S-1-5-32-544:(OI)(CI)F' /T /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}'" command += f" && chmod 700 '{remote_dir}'" 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 _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() cfg=load_config().get('meshcentral',{}) 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 {} 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'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 "-"}', ] 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)}'] 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". # Recreate the directory, remove any stale target and retry once. repair=_prepare_remote_directory(cfg,asset,password,job.platform,remote_dir,min(timeout,120),remote_file) 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()}') 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.'] 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()