using System; using System.Net; using System.Net.Sockets; using Newtonsoft.Json; using QuantConnect.Configuration; using QuantConnect.Interfaces; using QuantConnect.Logging; using QuantConnect.Notifications; using QuantConnect.Packets; using NetMQ; using NetMQ.Sockets; namespace QuantConnect.Messaging { /// /// Message handler that sends messages over tcp using NetMQ. /// public class StreamingMessageHandler : IMessagingHandler { private string _port; private PushSocket _server; private AlgorithmNodePacket _job; /// /// Gets or sets whether this messaging handler has any current subscribers. /// This is not used in this message handler. Messages are sent via tcp as they arrive /// public bool HasSubscribers { get; set; } /// /// Initialize the messaging system /// public void Initialize() { _port = Config.Get("desktop-http-port"); CheckPort(); _server = new PushSocket("@tcp://*:" + _port); } /// /// Set the user communication channel /// /// public void SetAuthentication(AlgorithmNodePacket job) { _job = job; Transmit(_job); } /// /// Send any notification with a base type of Notification. /// /// The notification to be sent. public void SendNotification(Notification notification) { var type = notification.GetType(); if (type == typeof(NotificationEmail) || type == typeof(NotificationWeb) || type == typeof(NotificationSms)) { Log.Error("Messaging.SendNotification(): Send not implemented for notification of type: " + type.Name); return; } notification.Send(); } /// /// Send all types of packets /// public void Send(Packet packet) { Transmit(packet); if (StreamingApi.IsEnabled) { StreamingApi.Transmit(_job.UserId, _job.Channel, packet); } } /// /// Send a message to the _server using ZeroMQ /// /// Packet to transmit public void Transmit(Packet packet) { var payload = JsonConvert.SerializeObject(packet); var message = new NetMQMessage(); message.Append(payload); _server.SendMultipartMessage(message); } /// /// Check if port to be used by the desktop application is available. /// private void CheckPort() { try { TcpListener tcpListener = new TcpListener(IPAddress.Any, _port.ToInt32()); tcpListener.Start(); tcpListener.Stop(); } catch { throw new Exception("The port configured in config.json is either being used or blocked by a firewall." + "Please choose a new port or open the port in the firewall."); } } } }