diff --git a/README.md b/README.md index 7315d56..29ed1f8 100644 --- a/README.md +++ b/README.md @@ -68,6 +68,23 @@ In XTMF2, open **Settings**, open **RunServers**, and add a remote endpoint. Ent The GUI pins the server certificate to this fingerprint and authenticates with the token before creating the RunServer bus. A mismatched certificate or token is rejected. Do not expose the TCP port to untrusted networks; use firewall rules or a private network as appropriate. +## Distributed estimation + +Estimation runs can evaluate candidate parameter vectors concurrently across multiple connected RunServers. Normal model runs and calibration runs continue to use one RunServer. + +To start a distributed estimation run: + +1. Connect the required RunServers from **Settings** > **RunServers**. +2. Open the model system and choose **Run Estimation**. +3. Select two or more RunServers in the run configuration dialog and choose the orchestrator RunServer. The other selected servers are workers. +4. If the model has an estimation `InputDirectory`, review its worker-local value for each selected worker. These values are saved in the GUI settings by model node ID and are reused on later estimation runs. + +The orchestrator RunServer owns the estimation algorithm and connects directly to the workers. Each worker constructs and validates its own local copy of the model. Input-directory overrides are applied only to that worker's copy; the saved model system is not modified. The orchestrator must be able to reach every worker's configured TCP endpoint. + +After submission, the GUI is not required for the optimization to continue. The orchestrator persists `estimation-completion.json` in the run directory when the job finishes, including the best parameter values and completion status. + +RunServers that receive path overrides must be running the current shared-estimation worker build. Runs without path overrides retain the version-1 shared-estimation protocol and remain compatible with workers that support the original distributed-estimation protocol. + ## Main Branches There are 4 major branches for XTMF 2 intended for different purposes: diff --git a/src/XTMF2.Client/Program.cs b/src/XTMF2.Client/Program.cs index ac9ec8c..8807715 100644 --- a/src/XTMF2.Client/Program.cs +++ b/src/XTMF2.Client/Program.cs @@ -19,6 +19,8 @@ You should have received a copy of the GNU General Public License using System; using System.Collections.Generic; using System.Diagnostics; +using System.IO.Compression; +using System.Linq; using System.IO; using System.Threading; using System.Threading.Tasks; @@ -27,15 +29,25 @@ You should have received a copy of the GNU General Public License using System.Security.Cryptography; using System.Security.Cryptography.X509Certificates; using XTMF2.Bus; +using XTMF2.Bus.Optimization; using XTMF2.Configuration; namespace XTMF2.Client { public class Program { + private static string[] _originalArguments = Array.Empty(); + [MTAThread] static void Main(string[] args) { + _originalArguments = args; + if (string.Equals(Environment.GetEnvironmentVariable("XTMF2_DEPLOYMENT_WATCHDOG"), "1", StringComparison.Ordinal)) + { + RunDeploymentWatchdog(args); + return; + } + LogDeploymentStartupSummary(); if (args.Length == 0) { Console.WriteLine("Usage: XTMF.Run [-setup-security DIRECTORY] [-loadDLL dllPath] [-tcp ADDRESS PORT -security DIRECTORY] [-namedPipe PIPE_NAME]"); @@ -177,8 +189,11 @@ private static void RunTcpServer(string address, int port, string? securityDirec if (certificate is not null) Console.WriteLine($"Certificate fingerprint: {RunServerSecurity.GetFingerprint(certificate)}"); Console.Out.Flush(); + SignalDeploymentReady(); using (tcpListener) using (var shutdown = new CancellationTokenSource()) + using (var remoteEstimationRegistry = new RemoteSharedEstimationRegistry()) + using (var remoteRunRegistry = new RemoteRunRegistry()) { Console.CancelKeyPress += (_, eventArgs) => { @@ -212,13 +227,16 @@ private static void RunTcpServer(string address, int port, string? securityDirec Console.Out.Flush(); var acceptedClient = client; _ = securityDirectory is null - ? Task.Run(() => RunTcpClientUnsecured(acceptedClient, extraDlls)) - : Task.Run(() => RunTcpClient(acceptedClient, certificate!, token!, extraDlls)); + ? Task.Run(() => RunTcpClientUnsecured(acceptedClient, extraDlls, + remoteEstimationRegistry, remoteRunRegistry)) + : Task.Run(() => RunTcpClient(acceptedClient, certificate!, token!, extraDlls, + remoteEstimationRegistry, remoteRunRegistry)); } } } - private static void RunTcpClientUnsecured(TcpClient client, List extraDlls) + private static void RunTcpClientUnsecured(TcpClient client, List extraDlls, + RemoteSharedEstimationRegistry remoteEstimationRegistry, RemoteRunRegistry remoteRunRegistry) { var remoteEndpoint = client.Client.RemoteEndPoint?.ToString() ?? "unknown endpoint"; Console.WriteLine($"Local RunServer connection accepted from {remoteEndpoint}"); @@ -228,7 +246,10 @@ private static void RunTcpClientUnsecured(TcpClient client, List extraDl { try { - RunClient(stream, extraDlls, usePrivateWorkspace: true); + RunClient(stream, extraDlls, usePrivateWorkspace: true, + allowDeployment: false, + remoteEstimationRegistry: remoteEstimationRegistry, + remoteRunRegistry: remoteRunRegistry); } catch (Exception ex) { @@ -238,7 +259,8 @@ private static void RunTcpClientUnsecured(TcpClient client, List extraDl } } - private static void RunTcpClient(TcpClient client, X509Certificate2 certificate, string token, List extraDlls) + private static void RunTcpClient(TcpClient client, X509Certificate2 certificate, string token, List extraDlls, + RemoteSharedEstimationRegistry remoteEstimationRegistry, RemoteRunRegistry remoteRunRegistry) { var remoteEndpoint = client.Client.RemoteEndPoint?.ToString() ?? "unknown endpoint"; using (client) @@ -255,7 +277,10 @@ private static void RunTcpClient(TcpClient client, X509Certificate2 certificate, { try { - RunClient(stream, extraDlls, usePrivateWorkspace: true); + RunClient(stream, extraDlls, usePrivateWorkspace: true, + allowDeployment: true, + remoteEstimationRegistry: remoteEstimationRegistry, + remoteRunRegistry: remoteRunRegistry); } catch (Exception ex) { @@ -268,7 +293,9 @@ private static void RunTcpClient(TcpClient client, X509Certificate2 certificate, } } - private static void RunClient(Stream serverStream, List extraDlls, SystemConfiguration? config = null, bool usePrivateWorkspace = false) + private static void RunClient(Stream serverStream, List extraDlls, SystemConfiguration? config = null, + bool usePrivateWorkspace = false, RemoteSharedEstimationRegistry? remoteEstimationRegistry = null, + RemoteRunRegistry? remoteRunRegistry = null, bool allowDeployment = false) { var runtime = XTMFRuntime.CreateRuntime(config); var loadedConfig = runtime.SystemConfiguration; @@ -276,8 +303,350 @@ private static void RunClient(Stream serverStream, List extraDlls, Syste { loadedConfig.LoadAssembly(dll); } - using var clientBus = new RunServerBus(serverStream, true, runtime, extraDlls, System.Diagnostics.Debugger.IsAttached, usePrivateWorkspace); - clientBus.ProcessRequests(); + using var ownedRemoteEstimationRegistry = remoteEstimationRegistry is null + ? new RemoteSharedEstimationRegistry() + : null; + var registry = remoteEstimationRegistry ?? ownedRemoteEstimationRegistry!; + using var clientBus = new RunServerBus(serverStream, true, runtime, extraDlls, + System.Diagnostics.Debugger.IsAttached, usePrivateWorkspace, remoteRunRegistry, allowDeployment); + using var sharedEstimationWorker = clientBus.AttachSharedEstimationWorker(); + using var remoteCoordinator = new RemoteSharedEstimationCoordinatorSession(clientBus, registry); + clientBus.SetDeploymentGate( + () => + { + remoteRunRegistry?.BeginDrain(); + registry.BeginDrain(); + }, + () => + { + remoteRunRegistry?.EndDrain(); + registry.EndDrain(); + }, + () => + (remoteRunRegistry?.IsIdle ?? true) && registry.IsIdle, + RestartWithOriginalArguments); + clientBus.SetSharedActivityProviders(registry.GetActiveActivities, + sharedEstimationWorker.GetActiveActivities); + try + { + clientBus.ProcessRequests(); + } + finally + { + remoteRunRegistry?.Detach(clientBus); + } + } + + private static bool RestartWithOriginalArguments(string stagingRoot) + { + if (!Directory.Exists(stagingRoot) || !File.Exists(Path.Combine(stagingRoot, "deployment.zip"))) + return false; + + var processPath = Environment.ProcessPath; + if (string.IsNullOrWhiteSpace(processPath)) + return false; + + string? deploymentDirectory = null; + try + { + var processDirectory = GetInstallationDirectory(processPath); + if (string.IsNullOrWhiteSpace(processDirectory)) + return false; + + deploymentDirectory = Path.Combine(processDirectory, ".deployments", + DateTime.UtcNow.ToString("yyyyMMddHHmmssfff") + "-" + Guid.NewGuid().ToString("N")); + Directory.CreateDirectory(deploymentDirectory); + var archivePath = Path.Combine(deploymentDirectory, "deployment.zip"); + File.Copy(Path.Combine(stagingRoot, "deployment.zip"), archivePath); + ExtractDeploymentArchive(archivePath, deploymentDirectory); + EnsureDeploymentPayload(deploymentDirectory); + + var readyFile = Path.Combine(deploymentDirectory, ".deployment-ready"); + var watchdogInfo = CreateRunServerStartInfo(processPath, processDirectory, _originalArguments); + watchdogInfo.Environment["XTMF2_DEPLOYMENT_WATCHDOG"] = "1"; + watchdogInfo.Environment["XTMF2_DEPLOYMENT_WATCHDOG_OLD_PID"] = + Environment.ProcessId.ToString(); + watchdogInfo.Environment["XTMF2_DEPLOYMENT_WATCHDOG_READY_FILE"] = readyFile; + watchdogInfo.Environment["XTMF2_DEPLOYMENT_WATCHDOG_PROCESS_PATH"] = processPath; + watchdogInfo.Environment["XTMF2_DEPLOYMENT_WATCHDOG_DEPLOYMENT_DIRECTORY"] = deploymentDirectory; + foreach (var argument in _originalArguments) + watchdogInfo.ArgumentList.Add(argument); + if (Process.Start(watchdogInfo) is null) + throw new InvalidOperationException("Unable to start the deployment watchdog."); + Environment.Exit(0); + return true; + } + catch (Exception exception) when (exception is IOException or InvalidDataException or UnauthorizedAccessException or + InvalidOperationException or System.ComponentModel.Win32Exception) + { + Console.Error.WriteLine($"Unable to restart RunServer: {exception.Message}"); + if (deploymentDirectory is not null) + { + try + { + Directory.Delete(deploymentDirectory, recursive: true); + } + catch (IOException cleanupException) + { + Console.Error.WriteLine($"Unable to clean up failed deployment: {cleanupException.Message}"); + } + } + return false; + } + } + + private static void SignalDeploymentReady() + { + var readyFile = Environment.GetEnvironmentVariable("XTMF2_DEPLOYMENT_READY_FILE"); + if (string.IsNullOrWhiteSpace(readyFile)) + return; + + try + { + File.WriteAllText(readyFile, "ready"); + } + catch (IOException exception) + { + Console.Error.WriteLine($"Unable to signal RunServer readiness: {exception.Message}"); + } + } + + private static void RunDeploymentWatchdog(string[] arguments) + { + var readyFile = Environment.GetEnvironmentVariable("XTMF2_DEPLOYMENT_WATCHDOG_READY_FILE"); + var processPath = Environment.GetEnvironmentVariable("XTMF2_DEPLOYMENT_WATCHDOG_PROCESS_PATH"); + var deploymentDirectory = Environment.GetEnvironmentVariable("XTMF2_DEPLOYMENT_WATCHDOG_DEPLOYMENT_DIRECTORY"); + if (string.IsNullOrWhiteSpace(readyFile) || string.IsNullOrWhiteSpace(processPath) || + string.IsNullOrWhiteSpace(deploymentDirectory)) + return; + + var oldProcessIdText = Environment.GetEnvironmentVariable("XTMF2_DEPLOYMENT_WATCHDOG_OLD_PID"); + _ = int.TryParse(oldProcessIdText, out var oldProcessId); + var installationDirectory = GetInstallationDirectory(processPath); + var deadline = DateTime.UtcNow.AddSeconds(30); + while (DateTime.UtcNow < deadline) + { + try + { + if (oldProcessId <= 0 || Process.GetProcessById(oldProcessId).HasExited) + break; + } + catch (ArgumentException) + { + break; + } + Thread.Sleep(250); + } + + string? backupDirectory = null; + Process? deployedProcess = null; + try + { + if (oldProcessId > 0) + { + try + { + if (!Process.GetProcessById(oldProcessId).HasExited) + throw new InvalidOperationException("The previous RunServer process did not exit before activation."); + } + catch (ArgumentException) + { + } + } + + backupDirectory = Path.Combine(deploymentDirectory, ".backup"); + Directory.CreateDirectory(backupDirectory); + EnsureDeploymentPayload(deploymentDirectory); + ReplaceDeploymentFile(Path.Combine(installationDirectory, "XTMF2.dll"), + Path.Combine(deploymentDirectory, "XTMF2.dll"), Path.Combine(backupDirectory, "XTMF2.dll")); + ReplaceDeploymentDirectory(Path.Combine(installationDirectory, "Modules"), + Path.Combine(deploymentDirectory, "Modules"), Path.Combine(backupDirectory, "Modules")); + + var startInfo = CreateRunServerStartInfo(processPath, installationDirectory, arguments); + ClearDeploymentWatchdogEnvironment(startInfo); + startInfo.Environment["XTMF2_DEPLOYMENT_MANIFEST"] = + Path.Combine(deploymentDirectory, "deployment-manifest.txt"); + startInfo.Environment["XTMF2_DEPLOYMENT_READY_FILE"] = readyFile; + foreach (var argument in arguments) + startInfo.ArgumentList.Add(argument); + deployedProcess = Process.Start(startInfo) + ?? throw new InvalidOperationException("Unable to start the deployed RunServer."); + + var readyDeadline = DateTime.UtcNow.AddSeconds(30); + while (!File.Exists(readyFile) && DateTime.UtcNow < readyDeadline) + { + if (deployedProcess.HasExited) + throw new InvalidOperationException("The deployed RunServer exited before becoming ready."); + Thread.Sleep(250); + } + if (!File.Exists(readyFile)) + throw new InvalidOperationException("The deployed RunServer did not become ready."); + + Directory.Delete(backupDirectory, recursive: true); + Directory.Delete(deploymentDirectory, recursive: true); + } + catch (Exception exception) when (exception is IOException or InvalidDataException or UnauthorizedAccessException or InvalidOperationException or System.ComponentModel.Win32Exception) + { + Console.Error.WriteLine($"Unable to activate the deployed RunServer: {exception.Message}"); + if (backupDirectory is null) + return; + try + { + if (deployedProcess is not null && !deployedProcess.HasExited) + deployedProcess.Kill(entireProcessTree: true); + RestoreDeploymentFile(Path.Combine(installationDirectory, "XTMF2.dll"), + Path.Combine(backupDirectory ?? string.Empty, "XTMF2.dll")); + RestoreDeploymentDirectory(Path.Combine(installationDirectory, "Modules"), + Path.Combine(backupDirectory ?? string.Empty, "Modules")); + var rollbackInfo = CreateRunServerStartInfo(processPath, installationDirectory, arguments); + ClearDeploymentWatchdogEnvironment(rollbackInfo); + Process.Start(rollbackInfo); + } + catch (Exception rollbackException) when (rollbackException is IOException or UnauthorizedAccessException or InvalidOperationException or System.ComponentModel.Win32Exception) + { + Console.Error.WriteLine($"Unable to restore the previous RunServer version: {rollbackException.Message}"); + } + } + } + + private static void ClearDeploymentWatchdogEnvironment(ProcessStartInfo startInfo) + { + startInfo.Environment.Remove("XTMF2_DEPLOYMENT_WATCHDOG"); + startInfo.Environment.Remove("XTMF2_DEPLOYMENT_WATCHDOG_OLD_PID"); + startInfo.Environment.Remove("XTMF2_DEPLOYMENT_WATCHDOG_READY_FILE"); + startInfo.Environment.Remove("XTMF2_DEPLOYMENT_WATCHDOG_PROCESS_PATH"); + startInfo.Environment.Remove("XTMF2_DEPLOYMENT_WATCHDOG_DEPLOYMENT_DIRECTORY"); + } + + private static ProcessStartInfo CreateRunServerStartInfo(string processPath, string workingDirectory, + IEnumerable arguments) + { + var startInfo = new ProcessStartInfo(processPath) + { + UseShellExecute = false, + WorkingDirectory = workingDirectory, + RedirectStandardInput = false, + RedirectStandardOutput = false, + RedirectStandardError = false, + CreateNoWindow = false + }; + if (IsDotnetHost(processPath)) + startInfo.ArgumentList.Add(typeof(Program).Assembly.Location); + foreach (var argument in arguments) + startInfo.ArgumentList.Add(argument); + return startInfo; + } + + private static bool IsDotnetHost(string processPath) + { + var fileName = Path.GetFileNameWithoutExtension(processPath); + return string.Equals(fileName, "dotnet", StringComparison.OrdinalIgnoreCase); + } + + private static string GetInstallationDirectory(string processPath) + { + if (IsDotnetHost(processPath)) + return Path.GetDirectoryName(typeof(Program).Assembly.Location) ?? Environment.CurrentDirectory; + return Path.GetDirectoryName(processPath) ?? Environment.CurrentDirectory; + } + + private static void ReplaceDeploymentFile(string destination, string staged, string backup) + { + if (!File.Exists(staged)) + throw new InvalidDataException($"Deployment is missing required file '{Path.GetFileName(staged)}'."); + if (File.Exists(destination)) + { + Directory.CreateDirectory(Path.GetDirectoryName(backup)!); + File.Move(destination, backup, overwrite: true); + } + File.Move(staged, destination, overwrite: true); + } + + private static void ExtractDeploymentArchive(string archivePath, string destinationRoot) + { + using var archive = ZipFile.OpenRead(archivePath); + foreach (var entry in archive.Entries) + { + var destination = Path.GetFullPath(Path.Combine(destinationRoot, + entry.FullName.Replace('/', Path.DirectorySeparatorChar))); + var root = Path.GetFullPath(destinationRoot + Path.DirectorySeparatorChar); + if (!destination.StartsWith(root, StringComparison.Ordinal)) + throw new InvalidDataException("Deployment archive contains an unsafe path."); + if (string.IsNullOrEmpty(entry.Name)) + { + Directory.CreateDirectory(destination); + continue; + } + Directory.CreateDirectory(Path.GetDirectoryName(destination)!); + entry.ExtractToFile(destination, overwrite: true); + } + } + + private static void EnsureDeploymentPayload(string deploymentDirectory) + { + var runtimePath = Path.Combine(deploymentDirectory, "XTMF2.dll"); + var modulesPath = Path.Combine(deploymentDirectory, "Modules"); + if (File.Exists(runtimePath) && Directory.Exists(modulesPath)) + return; + + var archivePath = Path.Combine(deploymentDirectory, "deployment.zip"); + if (!File.Exists(archivePath)) + throw new InvalidDataException("Deployment staging data is incomplete."); + ExtractDeploymentArchive(archivePath, deploymentDirectory); + if (!File.Exists(runtimePath) || !Directory.Exists(modulesPath)) + throw new InvalidDataException("Deployment archive is missing XTMF2.dll or Modules."); + } + + private static void RestoreDeploymentFile(string destination, string backup) + { + if (File.Exists(backup)) + File.Move(backup, destination, overwrite: true); + } + + private static void ReplaceDeploymentDirectory(string destination, string staged, string backup) + { + if (!Directory.Exists(staged)) + throw new InvalidDataException("Deployment is missing the required Modules directory."); + if (Directory.Exists(destination)) + Directory.Move(destination, backup); + Directory.Move(staged, destination); + } + + private static void RestoreDeploymentDirectory(string destination, string backup) + { + if (!Directory.Exists(backup)) + return; + if (Directory.Exists(destination)) + Directory.Delete(destination, recursive: true); + Directory.Move(backup, destination); + } + + private static void LogDeploymentStartupSummary() + { + var manifestPath = Environment.GetEnvironmentVariable("XTMF2_DEPLOYMENT_MANIFEST"); + if (string.IsNullOrWhiteSpace(manifestPath) || !File.Exists(manifestPath)) + return; + + try + { + var modules = File.ReadAllLines(manifestPath) + .Where(module => module.Length <= 256 && + !module.Any(char.IsControl) && + (module.Equals("XTMF2.dll", StringComparison.Ordinal) || + module.StartsWith("Modules/", StringComparison.Ordinal))) + .Take(4096) + .ToArray(); + Console.WriteLine($"RunServer restarted with deployment containing {modules.Length} module file(s):"); + foreach (var module in modules) + Console.WriteLine($" {module}"); + Console.Out.Flush(); + Environment.SetEnvironmentVariable("XTMF2_DEPLOYMENT_MANIFEST", null); + } + catch (IOException exception) + { + Console.WriteLine($"RunServer restarted, but deployment overview could not be read: {exception.Message}"); + Console.Out.Flush(); + } } } } \ No newline at end of file diff --git a/src/XTMF2.GUI/App.axaml b/src/XTMF2.GUI/App.axaml index fe3db5b..71d38d7 100644 --- a/src/XTMF2.GUI/App.axaml +++ b/src/XTMF2.GUI/App.axaml @@ -332,6 +332,12 @@ + +