From 4f9dc57c251cc4b334cb245c842f0914810cc63e Mon Sep 17 00:00:00 2001 From: Viacheslav Demydiuk Date: Fri, 12 Jan 2024 23:24:47 +0200 Subject: [PATCH] Moved logic related to command to MtRpcClient --- MtApi5/Mt5QuoteEventArgs.cs | 4 +- MtApi5/Mt5QuotesEventArgs.cs | 8 ++ MtApi5/MtApi5Client.cs | 153 +++++++---------------------------- MtClient/MtMessage.cs | 22 ++--- MtClient/MtRpcClient.cs | 121 ++++++++++++++++++++++++++- 5 files changed, 165 insertions(+), 143 deletions(-) create mode 100755 MtApi5/Mt5QuotesEventArgs.cs diff --git a/MtApi5/Mt5QuoteEventArgs.cs b/MtApi5/Mt5QuoteEventArgs.cs index 068b100b..ccd6528f 100755 --- a/MtApi5/Mt5QuoteEventArgs.cs +++ b/MtApi5/Mt5QuoteEventArgs.cs @@ -1,6 +1,4 @@ -using System; - -namespace MtApi5 +namespace MtApi5 { public class Mt5QuoteEventArgs: EventArgs { diff --git a/MtApi5/Mt5QuotesEventArgs.cs b/MtApi5/Mt5QuotesEventArgs.cs new file mode 100755 index 00000000..11a17a59 --- /dev/null +++ b/MtApi5/Mt5QuotesEventArgs.cs @@ -0,0 +1,8 @@ +namespace MtApi5 +{ + public class Mt5QuotesEventArgs(IEnumerable quotes) : EventArgs + { + public IEnumerable Quotes { get; } = quotes; + + } +} diff --git a/MtApi5/MtApi5Client.cs b/MtApi5/MtApi5Client.cs index 57e57963..363a6527 100755 --- a/MtApi5/MtApi5Client.cs +++ b/MtApi5/MtApi5Client.cs @@ -3,9 +3,7 @@ using MtApi5.Requests; using Newtonsoft.Json; using MtApi5.Events; using MtClient; -using System.ComponentModel.Design; using System.Reflection.Metadata; -using System.Collections.Generic; namespace MtApi5 { @@ -40,8 +38,6 @@ namespace MtApi5 private HashSet _experts = []; private Dictionary _quotes = []; private readonly EventWaitHandle _quotesWaiter = new AutoResetEvent(false); - private int _commandId = 0; - private readonly Dictionary _tasks = []; #endregion #region Public Methods @@ -3319,7 +3315,7 @@ namespace MtApi5 public event EventHandler? OnBookEvent; public event EventHandler? OnLastTimeBar; public event EventHandler? OnLockTicks; - public event EventHandler>? QuoteList; + public event EventHandler? QuoteList; #endregion #region Private Methods @@ -3351,7 +3347,10 @@ namespace MtApi5 ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(Mt5ConnectionState.Connecting, message)); var client = new MtRpcClient(host, port); - client.MessageReceived += Client_OnMessageReceived; + client.ExpertList += Client_ExpertList; + client.ExpertAdded += Client_ExpertAdded; + client.ExpertRemoved += Client_ExpertRemoved; + client.MtEventReceived += Client_MtEventReceived; client.ConnectionFailed += Client_OnConnectionFailed; client.Disconnected += Client_Disconnected; @@ -3384,36 +3383,28 @@ namespace MtApi5 ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(state, message)); } - private void Client_OnMessageReceived(object? o, MtMessage msg) + private void Client_MtEventReceived(object? sender, MtEventArgs e) { - Log?.Debug($"Message received: type = {msg}"); - - switch (msg.MsgType) - { - case MessageType.ExpertList: - Task.Run(() => ProcessExpertList(msg as MtExpertListMsg)); //must be runned on another thread to avoid blocking - break; - case MessageType.ExpertAdded: - Task.Run(() => ProcessExpertAdded(msg as MtExpertAddedMsg)); //must be runned on another thread to avoid blocking - break; - case MessageType.ExpertRemoved: - ProcessExpertRemoved(msg as MtExpertRemovedMsg); - break; - case MessageType.Event: - ProcessEvent(msg as MtEvent); - break; - case MessageType.Response: - ProcessResponse(msg as MtResponse); - break; - } + Task.Run(() => _mtEventHandlers[(Mt5EventTypes)e.EventType](e.ExpertHandle, e.Payload)); } - private void ProcessExpertList(MtExpertListMsg? msg) + private void Client_ExpertList(object? sender, MtExpertListEventArgs e) { - if (msg == null) - return; + Task.Run(()=>ProcessExpertList(e.Experts)); + } - HashSet experts = msg.Experts; + private void Client_ExpertAdded(object? sender, MtExpertEventArgs e) + { + Task.Run(() => ProcessExpertAdded(e.Expert)); + } + + private void Client_ExpertRemoved(object? sender, MtExpertEventArgs e) + { + Task.Run(() => ProcessExpertRemoved(e.Expert)); + } + + private void ProcessExpertList(HashSet experts) + { if (experts == null || experts.Count == 0) { Log?.Warn("ProcessExpertList: expert list invalid or empty"); @@ -3435,16 +3426,11 @@ namespace MtApi5 } _quotesWaiter.Set(); - QuoteList?.Invoke(this, quotes.Values.ToList()); + QuoteList?.Invoke(this, new(quotes.Values.ToList())); } - private void ProcessExpertAdded(MtExpertAddedMsg? msg) + private void ProcessExpertAdded(int handle) { - if (msg == null) - return; - - int handle = msg.ExpertHandle; - Log?.Debug($"ProcessExpertAdded: {handle}"); bool added; @@ -3473,13 +3459,8 @@ namespace MtApi5 Log?.Warn($"ProcessExpertAdded: expert handle {handle} is already exist"); } - private void ProcessExpertRemoved(MtExpertRemovedMsg? msg) + private void ProcessExpertRemoved(int handle) { - if (msg == null) - return; - - int handle = msg.ExpertHandle; - Log?.Debug($"ProcessExpertRemoved: {handle}"); Mt5Quote? quote = null; @@ -3494,37 +3475,6 @@ namespace MtApi5 QuoteRemoved?.Invoke(this, new Mt5QuoteEventArgs(quote)); } - private void ProcessEvent(MtEvent? msg) - { - if (msg == null) - return; - - var handle = msg.ExpertHandle; - var eventType = (Mt5EventTypes)msg.EventType; - var payload = msg.Payload; - - Log?.Debug($"ProcessEvent: {handle}, {eventType}, {payload}"); - - _mtEventHandlers[eventType](handle, payload); - } - - private void ProcessResponse(MtResponse? msg) - { - if (msg == null) - return; - - var handle = msg.ExpertHandle; - var commandId = msg.CommandId; - var payload = msg.Payload; - - Log?.Debug($"ProcessResponse: {handle}, {commandId}, [{payload}]"); - - lock (_locker) - { - if (_tasks.TryGetValue(commandId, out CommandTask? value)) - value.SetResponse(payload); - } - } private Mt5Quote? GetQuote(int expertHandle) { Log?.Debug($"GetQuote: expertHandle = {expertHandle}"); @@ -3625,7 +3575,6 @@ namespace MtApi5 _quotes.Clear(); _experts.Clear(); - _tasks.Clear(); } client?.Disconnect(); @@ -3635,31 +3584,6 @@ namespace MtApi5 ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(state, message)); } - internal class CommandTask - { - private readonly EventWaitHandle responseWaiter_ = new AutoResetEvent(false); - private string? response_; - private readonly object locker_ = new(); - - public string? WaitResponse(int time) - { - responseWaiter_.WaitOne(time); - lock (locker_) - { - return response_; - } - } - - public void SetResponse(string result) - { - lock (locker_) - { - response_ = result; - } - responseWaiter_.Set(); - } - } - private T? SendCommand(int expertHandle, Mt5CommandType commandType, object? payload = null) { var client = Client; @@ -3670,22 +3594,9 @@ namespace MtApi5 } var payloadJson = JsonConvert.SerializeObject(payload); + Log?.Debug($"SendCommand: sending '{payloadJson}' ..."); - MtCommand command = new(expertHandle, (int)commandType, _commandId++, payloadJson); - CommandTask commandTask = new(); - lock (_locker) - { - _tasks[command.CommandId] = commandTask; - } - - Log?.Debug($"SendCommand: sending {command.CommandId} ..."); - client.Send(command); - var responseJson = commandTask.WaitResponse(10000); // 10 sec - - lock(_locker) - { - _tasks.Remove(command.CommandId); - } + var responseJson = client.SendCommand(expertHandle, (int)commandType, payloadJson); Log?.Debug($"SendCommand: received response JSON [{responseJson}]"); @@ -3780,19 +3691,11 @@ namespace MtApi5 return response.Value; } - - //private void _client_QuoteUpdated(MtQuote quote) - //{ - // if (quote == null) return; - // QuoteUpdate?.Invoke(this, new Mt5QuoteEventArgs(new Mt5Quote(quote))); - // QuoteUpdated?.Invoke(this, quote.Instrument, quote.Bid, quote.Ask); - //} - private void OnConnected() { Log?.Debug("OnConnected: begin"); - Client?.Send(new MtNotification(MtNotificationType.ClientReady)); + Client?.NotifyClientReady(); _isBacktestingMode = IsTesting(); diff --git a/MtClient/MtMessage.cs b/MtClient/MtMessage.cs index fe25c93a..4618edfa 100755 --- a/MtClient/MtMessage.cs +++ b/MtClient/MtMessage.cs @@ -1,6 +1,6 @@ namespace MtClient { - public enum MessageType + internal enum MessageType { Command = 0, Response = 1, @@ -11,12 +11,12 @@ Notification = 6 } - public enum MtNotificationType + internal enum MtNotificationType { ClientReady = 0 } - public abstract class MtMessage + internal abstract class MtMessage { public abstract MessageType MsgType { get; } @@ -28,7 +28,7 @@ protected abstract string GetMessageBody(); } - public class MtCommand(int expertHandle, int commandType, int commandId, string payload) : MtMessage + internal class MtCommand(int expertHandle, int commandType, int commandId, string payload) : MtMessage { public override MessageType MsgType => MessageType.Command; @@ -43,7 +43,7 @@ } } - public class MtNotification(MtNotificationType notificationType) : MtMessage + internal class MtNotification(MtNotificationType notificationType) : MtMessage { public override MessageType MsgType => MessageType.Notification; @@ -55,7 +55,7 @@ public MtNotificationType NotificationType { private set; get; } = notificationType; } - public class MtEvent(int expertHandle, int eventType, string payload) : MtMessage + internal class MtEvent(int expertHandle, int eventType, string payload) : MtMessage { public override MessageType MsgType => MessageType.Event; @@ -79,7 +79,7 @@ } } - public class MtExpertAddedMsg(int expertHandle) : MtMessage + internal class MtExpertAddedMsg(int expertHandle) : MtMessage { public override MessageType MsgType => MessageType.ExpertAdded; @@ -99,7 +99,7 @@ } } - public class MtExpertRemovedMsg(int expertHandle) : MtMessage + internal class MtExpertRemovedMsg(int expertHandle) : MtMessage { public override MessageType MsgType => MessageType.ExpertRemoved; @@ -119,7 +119,7 @@ } } - public class MtExpertListMsg(HashSet experts) : MtMessage + internal class MtExpertListMsg(HashSet experts) : MtMessage { public override MessageType MsgType => MessageType.ExpertList; @@ -145,7 +145,7 @@ } } - public class MtResponse(int expertHandle, int commandId, string payload) : MtMessage + internal class MtResponse(int expertHandle, int commandId, string payload) : MtMessage { public override MessageType MsgType => MessageType.Response; @@ -177,7 +177,7 @@ } } - public static class MtMessageParser + internal static class MtMessageParser { static MtMessageParser() { diff --git a/MtClient/MtRpcClient.cs b/MtClient/MtRpcClient.cs index d3a4cff1..c5f08dca 100755 --- a/MtClient/MtRpcClient.cs +++ b/MtClient/MtRpcClient.cs @@ -56,7 +56,35 @@ namespace MtClient Log($"Disconnect: success"); } - public void Send(MtMessage message) + public string? SendCommand(int expertHandle, int commandType, string payload) + { + CommandTask commandTask = new(); + int commandId; + lock (_tasks) + { + commandId = nextCommandId++; + _tasks[commandId] = commandTask; + } + + MtCommand command = new(expertHandle, commandType, commandId, payload); + Send(command); + + var response = commandTask.WaitResponse(10000); // 10 sec + lock (_tasks) + { + _tasks.Remove(commandId); + } + + return response; + } + + public void NotifyClientReady() + { + MtNotification notification = new(MtNotificationType.ClientReady); + Send(notification); + } + + private void Send(MtMessage message) { lock (pendingMessages_) { @@ -164,8 +192,45 @@ namespace MtClient Log("OnReceive: Failed parse message payload"); return; } - - MessageReceived?.Invoke(this, message); + + switch (message.MsgType) + { + case MessageType.Event: + if (message is MtEvent e) + MtEventReceived?.Invoke(this, new(e.ExpertHandle, e.EventType, e.Payload)); + break; + case MessageType.Response: + if (message is MtResponse response) + ProcessResponse(response); + break; + case MessageType.ExpertList: + if (message is MtExpertListMsg expertListMsg) + ExpertList?.Invoke(this, new(expertListMsg.Experts)); + break; + case MessageType.ExpertAdded: + if (message is MtExpertAddedMsg expertAddedMsg) + ExpertAdded?.Invoke(this, new(expertAddedMsg.ExpertHandle)); + break; + case MessageType.ExpertRemoved: + if (message is MtExpertRemovedMsg expertRemovedMsg) + ExpertRemoved?.Invoke(this, new(expertRemovedMsg.ExpertHandle)); + break; + } + } + + private void ProcessResponse(MtResponse msg) + { + var handle = msg.ExpertHandle; + var commandId = msg.CommandId; + var payload = msg.Payload; + + Log($"ProcessResponse: {handle}, {commandId}, [{payload}]"); + + lock (_tasks) + { + if (_tasks.TryGetValue(commandId, out CommandTask? value)) + value.SetResponse(payload); + } } private void Log(string msg) @@ -173,9 +238,12 @@ namespace MtClient Console.WriteLine($"[{Environment.CurrentManagedThreadId}] {msg}"); } - public event EventHandler? MessageReceived; public event EventHandler? ConnectionFailed; public event EventHandler? Disconnected; + public event EventHandler? ExpertList; + public event EventHandler? ExpertAdded; + public event EventHandler? ExpertRemoved; + public event EventHandler? MtEventReceived; private readonly ClientWebSocket ws_ = new(); private readonly string host_; @@ -186,5 +254,50 @@ namespace MtClient private readonly Thread receiveThread_; private readonly Thread sendThread_; private readonly EventWaitHandle sendWaiter_ = new AutoResetEvent(false); + + private int nextCommandId = 0; + private readonly Dictionary _tasks = []; + } + + internal class CommandTask + { + private readonly EventWaitHandle responseWaiter_ = new AutoResetEvent(false); + private string? response_; + private readonly object locker_ = new(); + + public string? WaitResponse(int time) + { + responseWaiter_.WaitOne(time); + lock (locker_) + { + return response_; + } + } + + public void SetResponse(string result) + { + lock (locker_) + { + response_ = result; + } + responseWaiter_.Set(); + } + } + + public class MtEventArgs(int expertHandle, int eventType, string payload) : EventArgs + { + public int ExpertHandle { get; } = expertHandle; + public int EventType { get; } = eventType; + public string Payload { get; } = payload; + } + + public class MtExpertListEventArgs(HashSet experts) : EventArgs + { + public HashSet Experts { get; } = experts; + } + + public class MtExpertEventArgs(int expert) : EventArgs + { + public int Expert { get; }= expert; } } \ No newline at end of file