1c3d849ad5
* Reconcile duplicated code * Add License header * CS0219 Fixes: Value assigned, but never used * CA1507: Use nameof in place of string literals * CS0108 : Hides Inherited Member; Use new keyword to overwrite formally * CS0114: Hides inherited member; use override keyword * CS0168: Variable is declared but never used * Tests CS1062; using obsolete implicit Symbol -> String; fix via .ToString() * CS0472: Non Nullable Obj getting Null Checked * CS0067 Member not used; ignore all cases for future use * CS00162 : Unreachable code; either removed or ignored for debugging and test cases * CS0169 Remove non-used fields; ignore those that may be used in future * CS0414; Field is assigned but never used. * CS0618; Obsolete properties and members; Only fixes simple ones, rest will have to broken up * CS0649; Field never assigned too * CS0659 & CS0661 ; Overwrite operators and equals but not hashcode; I don't really override it but just call base * Small comment fix * Cleanup pragma statement
278 lines
8.3 KiB
C#
278 lines
8.3 KiB
C#
/*
|
|
* QUANTCONNECT.COM - Democratizing Finance, Empowering Individuals.
|
|
* Lean Algorithmic Trading Engine v2.0. Copyright 2014 QuantConnect Corporation.
|
|
*
|
|
* 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 http://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.Collections.Generic;
|
|
using QuantConnect.Interfaces;
|
|
using QuantConnect.Logging;
|
|
using QuantConnect.Notifications;
|
|
using QuantConnect.Packets;
|
|
|
|
namespace QuantConnect.Messaging
|
|
{
|
|
/// <summary>
|
|
/// Desktop implementation of messaging system for Lean Engine
|
|
/// </summary>
|
|
public class EventMessagingHandler : IMessagingHandler
|
|
{
|
|
private AlgorithmNodePacket _job;
|
|
private volatile bool _loaded;
|
|
private Queue<Packet> _queue;
|
|
|
|
/// <summary>
|
|
/// Gets or sets whether this messaging handler has any current subscribers.
|
|
/// When set to false, messages won't be sent.
|
|
/// </summary>
|
|
public bool HasSubscribers
|
|
{
|
|
get;
|
|
set;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Initialize the Messaging System Plugin.
|
|
/// </summary>
|
|
public void Initialize()
|
|
{
|
|
_queue = new Queue<Packet>();
|
|
|
|
ConsumerReadyEvent += () => { _loaded = true; };
|
|
}
|
|
|
|
/// <summary>
|
|
/// Set Loaded to true
|
|
/// </summary>
|
|
public void LoadingComplete()
|
|
{
|
|
_loaded = true;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Set the user communication channel
|
|
/// </summary>
|
|
/// <param name="job"></param>
|
|
public void SetAuthentication(AlgorithmNodePacket job)
|
|
{
|
|
_job = job;
|
|
}
|
|
|
|
#pragma warning disable 1591
|
|
public delegate void DebugEventRaised(DebugPacket packet);
|
|
public event DebugEventRaised DebugEvent;
|
|
|
|
public delegate void SystemDebugEventRaised(SystemDebugPacket packet);
|
|
#pragma warning disable 0067 // SystemDebugEvent is not used currently; ignore the warning
|
|
public event SystemDebugEventRaised SystemDebugEvent;
|
|
#pragma warning restore 0067
|
|
|
|
public delegate void LogEventRaised(LogPacket packet);
|
|
public event LogEventRaised LogEvent;
|
|
|
|
public delegate void RuntimeErrorEventRaised(RuntimeErrorPacket packet);
|
|
public event RuntimeErrorEventRaised RuntimeErrorEvent;
|
|
|
|
public delegate void HandledErrorEventRaised(HandledErrorPacket packet);
|
|
public event HandledErrorEventRaised HandledErrorEvent;
|
|
|
|
public delegate void BacktestResultEventRaised(BacktestResultPacket packet);
|
|
public event BacktestResultEventRaised BacktestResultEvent;
|
|
|
|
public delegate void ConsumerReadyEventRaised();
|
|
public event ConsumerReadyEventRaised ConsumerReadyEvent;
|
|
#pragma warning restore 1591
|
|
|
|
/// <summary>
|
|
/// Send any message with a base type of Packet.
|
|
/// </summary>
|
|
public void Send(Packet packet)
|
|
{
|
|
//Until we're loaded queue it up
|
|
if (!_loaded)
|
|
{
|
|
_queue.Enqueue(packet);
|
|
return;
|
|
}
|
|
|
|
//Catch up if this is the first time
|
|
while (_queue.Count > 0)
|
|
{
|
|
ProcessPacket(_queue.Dequeue());
|
|
}
|
|
|
|
//Finally process this new packet
|
|
ProcessPacket(packet);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Send any notification with a base type of Notification.
|
|
/// </summary>
|
|
/// <param name="notification">The notification to be sent.</param>
|
|
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();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Send any message with a base type of Packet that has been enqueued.
|
|
/// </summary>
|
|
public void SendEnqueuedPackets()
|
|
{
|
|
while (_queue.Count > 0 && _loaded)
|
|
{
|
|
ProcessPacket(_queue.Dequeue());
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Packet processing implementation
|
|
/// </summary>
|
|
private void ProcessPacket(Packet packet)
|
|
{
|
|
//Packets we handled in the UX.
|
|
switch (packet.Type)
|
|
{
|
|
case PacketType.Debug:
|
|
var debug = (DebugPacket)packet;
|
|
OnDebugEvent(debug);
|
|
break;
|
|
|
|
case PacketType.SystemDebug:
|
|
var systemDebug = (SystemDebugPacket)packet;
|
|
OnSystemDebugEvent(systemDebug);
|
|
break;
|
|
|
|
case PacketType.Log:
|
|
var log = (LogPacket)packet;
|
|
OnLogEvent(log);
|
|
break;
|
|
|
|
case PacketType.RuntimeError:
|
|
var runtime = (RuntimeErrorPacket)packet;
|
|
OnRuntimeErrorEvent(runtime);
|
|
break;
|
|
|
|
case PacketType.HandledError:
|
|
var handled = (HandledErrorPacket)packet;
|
|
OnHandledErrorEvent(handled);
|
|
break;
|
|
|
|
case PacketType.BacktestResult:
|
|
var result = (BacktestResultPacket)packet;
|
|
OnBacktestResultEvent(result);
|
|
break;
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Raise a debug event safely
|
|
/// </summary>
|
|
protected virtual void OnDebugEvent(DebugPacket packet)
|
|
{
|
|
var handler = DebugEvent;
|
|
|
|
if (handler != null)
|
|
{
|
|
handler(packet);
|
|
}
|
|
}
|
|
|
|
|
|
/// <summary>
|
|
/// Raise a system debug event safely
|
|
/// </summary>
|
|
protected virtual void OnSystemDebugEvent(SystemDebugPacket packet)
|
|
{
|
|
var handler = DebugEvent;
|
|
|
|
if (handler != null)
|
|
{
|
|
handler(packet);
|
|
}
|
|
}
|
|
|
|
|
|
/// <summary>
|
|
/// Handler for consumer ready code.
|
|
/// </summary>
|
|
public virtual void OnConsumerReadyEvent()
|
|
{
|
|
var handler = ConsumerReadyEvent;
|
|
if (handler != null)
|
|
{
|
|
handler();
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Raise a log event safely
|
|
/// </summary>
|
|
protected virtual void OnLogEvent(LogPacket packet)
|
|
{
|
|
var handler = LogEvent;
|
|
if (handler != null)
|
|
{
|
|
handler(packet);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Raise a handled error event safely
|
|
/// </summary>
|
|
protected virtual void OnHandledErrorEvent(HandledErrorPacket packet)
|
|
{
|
|
var handler = HandledErrorEvent;
|
|
if (handler != null)
|
|
{
|
|
handler(packet);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Raise runtime error safely
|
|
/// </summary>
|
|
protected virtual void OnRuntimeErrorEvent(RuntimeErrorPacket packet)
|
|
{
|
|
var handler = RuntimeErrorEvent;
|
|
if (handler != null)
|
|
{
|
|
handler(packet);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Raise a backtest result event safely.
|
|
/// </summary>
|
|
protected virtual void OnBacktestResultEvent(BacktestResultPacket packet)
|
|
{
|
|
var handler = BacktestResultEvent;
|
|
if (handler != null)
|
|
{
|
|
handler(packet);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Dispose of any resources
|
|
/// </summary>
|
|
public virtual void Dispose()
|
|
{
|
|
}
|
|
}
|
|
} |