/* Copyright NetFoundry Inc. Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at https://www.apache.org/licenses/LICENSE-2.0 Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License. */ using System; using System.Collections.Generic; using System.IO; using System.IO.Pipes; using System.Security.Principal; using System.Security.AccessControl; using System.Threading; using System.Threading.Tasks; using Newtonsoft.Json; using NLog; using ZitiDesktopEdge.DataStructures; using ZitiDesktopEdge.Server; using System.Diagnostics; using System.Reflection; using ZitiDesktopEdge.Utility; /// /// The implementation will abstract away the setup of the communication to /// the monitor service. This implementation will communicate to the service over a /// a NamedPipe. /// /// All communication is effectively serial - one or more messages sent and /// one or more messages returned. /// /// namespace ZitiDesktopEdge.ServiceClient { public class MonitorClient : AbstractClient { private static readonly Logger _logger = LogManager.GetCurrentClassLogger(); protected override Logger Logger { get { return _logger; } } public const int EXPECTED_API_VERSION = 1; public event EventHandler OnServiceStatusEvent; public event EventHandler OnNotificationEvent; public event EventHandler OnCaptureFeedbackProgressEvent; public DateTime LastFeedbackHeartbeat { get; private set; } = DateTime.MinValue; public bool IsServiceCapturingFeedback => (DateTime.UtcNow - LastFeedbackHeartbeat).TotalSeconds < 10; // Serializes the send/read RPC pairs so concurrent callers don't get each other's // responses off the shared pipe. Without this, e.g. clicking "Capture Feedback" and // then "Check for updates" from the tray will interleave on `ipcReader` and deadlock // both awaits forever. private readonly SemaphoreSlim _rpcLock = new SemaphoreSlim(1, 1); protected virtual void ServiceStatusEvent(MonitorServiceStatusEvent e) { OnServiceStatusEvent?.Invoke(this, e); } protected virtual void InstallationNotificationEvent(InstallationNotificationEvent e) { OnNotificationEvent?.Invoke(this, e); } protected virtual void CaptureFeedbackProgressEvent(MonitorServiceStatusEvent e) { LastFeedbackHeartbeat = DateTime.UtcNow; OnCaptureFeedbackProgressEvent?.Invoke(this, e); } public MonitorClient(string id) : base(id) { } async protected override Task ConnectPipesAsync() { await semaphoreSlim.WaitAsync(); try { pipeClient = new NamedPipeClientStream(localPipeServer, IPCServer.PipeName, PipeDirection.InOut); eventClient = new NamedPipeClientStream(localPipeServer, IPCServer.EventPipeName, PipeDirection.In); await eventClient.ConnectAsync(ServiceConnectTimeout); await pipeClient.ConnectAsync(ServiceConnectTimeout); ClientConnected(null); } catch (Exception ex) { semaphoreSlim.Release(); throw new MonitorServiceException("Could not connect to the monitor service.", ex); } semaphoreSlim.Release(); } protected override void ProcessLine(string line) { var evt = serializer.Deserialize(new JsonTextReader(new StringReader(line))); switch (evt.Type) { case "Notification": var instEvt = serializer.Deserialize(new JsonTextReader(new StringReader(line))); InstallationNotificationEvent(instEvt); break; case "CaptureFeedbackProgress": CaptureFeedbackProgressEvent(evt); break; default: ServiceStatusEvent(evt); break; } } async private Task sendMonitorClientAsync(object objtoSend) { try { await sendAsync("monitor", objtoSend); } catch (Exception ex) { throw new MonitorServiceException("Could not connect to the monitor service.", ex); } } async protected Task readMonitorClientAsync(StreamReader reader) where T : SvcResponse { return await readAsync("monitor", reader, DefaultReadTimeout); } async protected Task readMonitorClientAsync(StreamReader reader, TimeSpan timeout) where T : SvcResponse { return await readAsync("monitor", reader, timeout); } async public Task StopServiceAsync() { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "Stop", Action = "Normal" }; await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } finally { _rpcLock.Release(); } } async public Task StartServiceAsync(TimeSpan timeout) { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "Start", Action = "Normal" }; await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader, timeout); } finally { _rpcLock.Release(); } } async public Task ForceTerminateAsync() { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "Stop", Action = "Force" }; try { await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } catch (Exception ex) { Logger.Error(ex, "Unexpected error"); } return null; } finally { _rpcLock.Release(); } } async public Task StatusAsync() { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "Status", Action = "" }; try { await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } catch (Exception ex) { Logger.Error(ex, "Unexpected error"); } return null; } finally { _rpcLock.Release(); } } async public Task CaptureLogsAsync() { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "CaptureLogs", Action = "Normal" }; await sendMonitorClientAsync(action); LastFeedbackHeartbeat = DateTime.UtcNow; Task readTask = readMonitorClientAsync(ipcReader, TimeSpan.FromMinutes(30)); while (!readTask.IsCompleted) { await Task.WhenAny(readTask, Task.Delay(2000)); if (!IsServiceCapturingFeedback) { throw new MonitorServiceException("Feedback collection stopped responding."); } } return await readTask; } finally { _rpcLock.Release(); } } async public Task SetLogLevelAsync(string level) { if ("verbose".Equals(level?.ToLower())) { //only the data client understands verbose - so use trace... level = "TRACE"; } await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "SetLogLevel", Action = level }; await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } finally { _rpcLock.Release(); } } async public Task DoUpdateCheck() { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "DoUpdateCheck", Action = "" }; await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } finally { _rpcLock.Release(); } } async public Task TriggerUpdate(bool forceDefer = false) { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "TriggerUpdate", Action = forceDefer ? "defer" : "" }; await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } finally { _rpcLock.Release(); } } async public Task SetAutomaticUpgradeDisabledAsync(bool disabled) { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "SetAutomaticUpgradeDisabled", Action = (disabled ? "true" : "false") }; await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } finally { _rpcLock.Release(); } } async public Task SetAutomaticUpgradeURLAsync(string url) { await _rpcLock.WaitAsync(); try { ActionEvent action = new ActionEvent() { Op = "SetAutomaticUpgradeURL", Action = (url) }; await sendMonitorClientAsync(action); return await readMonitorClientAsync(ipcReader); } finally { _rpcLock.Release(); } } async public Task SetMaintenanceWindowAsync(MaintenanceWindowConfigRequest req) { await _rpcLock.WaitAsync(); try { req.Op = "SetMaintenanceWindow"; await sendMonitorClientAsync(req); return await readMonitorClientAsync(ipcReader); } finally { _rpcLock.Release(); } } } }