From e6731d882834a29dfd1ab48cc5b51aa764d25bc8 Mon Sep 17 00:00:00 2001 From: Viacheslav Demydiuk Date: Fri, 12 Jan 2024 18:22:26 +0200 Subject: [PATCH] Processing incoming messages in MtApi5 --- MtApi5/MtApi5Client.cs | 216 ++++++++++++++++++++++++----------------- 1 file changed, 128 insertions(+), 88 deletions(-) diff --git a/MtApi5/MtApi5Client.cs b/MtApi5/MtApi5Client.cs index c656b90b..8cb1a899 100755 --- a/MtApi5/MtApi5Client.cs +++ b/MtApi5/MtApi5Client.cs @@ -3,6 +3,9 @@ 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 { @@ -3344,13 +3347,13 @@ namespace MtApi5 _connectionState = Mt5ConnectionState.Connecting; } - string message = $"Connecting to {host}:{port}"; + string message = $"Connect: connecting to {host}:{port}"; ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(Mt5ConnectionState.Connecting, message)); var client = new MtRpcClient(host, port); - client.MessageReceived += _client_OnMessageReceived; - client.ConnectionFailed += _client_OnConnectionFailed; - client.Disconnected += _client_Disconnected; + client.MessageReceived += Client_OnMessageReceived; + client.ConnectionFailed += Client_OnConnectionFailed; + client.Disconnected += Client_Disconnected; var state = Mt5ConnectionState.Failed; try @@ -3360,15 +3363,16 @@ namespace MtApi5 } catch (Exception e) { - Log?.Warn($"Failed connection to {host}:{port}. {e.Message}"); + Log?.Warn($"Connect: Failed connection to {host}:{port}. {e.Message}"); } + Log?.Info($"Connect: connection to {host}:{port} is {state}"); + lock (_locker) { if (state == Mt5ConnectionState.Connected) { _client = client; - Log?.Info($"Connected to {host}:{port}"); } _connectionState = state; @@ -3378,35 +3382,30 @@ namespace MtApi5 OnConnected(); ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(state, message)); - - Log?.Info($"Connected finished"); } - private void _client_OnMessageReceived(object? o, MtMessage msg) + private void Client_OnMessageReceived(object? o, MtMessage msg) { - Log?.Debug($"Message received: {msg}"); + Log?.Debug($"Message received: type = {msg}"); - Task.Run(() => + switch (msg.MsgType) { - switch (msg.MsgType) - { - case MessageType.ExpertList: - ProcessExpertList(msg as MtExpertListMsg); - break; - case MessageType.ExpertAdded: - //ProcessExpertAdded(msg as MtExpertAddedMsg); - 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; - } - }); + 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; + } } private void ProcessExpertList(MtExpertListMsg? msg) @@ -3417,19 +3416,13 @@ namespace MtApi5 HashSet experts = msg.Experts; if (experts == null || experts.Count == 0) { - Console.WriteLine($"ProcessExpertList: expert list invalid"); + Log?.Warn("ProcessExpertList: expert list invalid or empty"); return; } - - lock(_locker) - { - _experts = experts; - } - + Dictionary quotes = []; foreach (var handle in experts) { - Console.WriteLine($"ProcessExpertList: {handle}"); var quote = GetQuote(handle); if (quote != null) quotes[handle] = quote; @@ -3437,11 +3430,82 @@ namespace MtApi5 lock (_locker) { - _quotes= quotes; + _experts = experts; + _quotes = quotes; } _quotesWaiter.Set(); - QuoteList?.Invoke(this, _quotes.Values.ToList()); + QuoteList?.Invoke(this, quotes.Values.ToList()); + } + + private void ProcessExpertAdded(MtExpertAddedMsg? msg) + { + if (msg == null) + return; + + int handle = msg.ExpertHandle; + + Log?.Debug($"ProcessExpertAdded: {handle}"); + + bool added; + lock (_locker) + { + added = _experts.Add(handle); + } + + if (added) + { + var quote = GetQuote(handle); + if (quote != null) + { + lock (_locker) + { + _quotes[handle] = quote; + } + + QuoteAdded?.Invoke(this, new Mt5QuoteEventArgs(quote)); + } + else + Log?.Warn($"ProcessExpertAdded: failed to get quote for expert {handle}"); + + } + else + Log?.Warn($"ProcessExpertAdded: expert handle {handle} is already exist"); + } + + private void ProcessExpertRemoved(MtExpertRemovedMsg? msg) + { + if (msg == null) + return; + + int handle = msg.ExpertHandle; + + Log?.Debug($"ProcessExpertRemoved: {handle}"); + + Mt5Quote? quote = null; + lock (_locker) + { + _experts.Remove(handle); + if (_quotes.TryGetValue(handle, out quote)) + _quotes.Remove(handle); + } + + if (quote != null) + 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) @@ -3461,29 +3525,27 @@ namespace MtApi5 value.SetResponse(payload); } } - private Mt5Quote? GetQuote(int expertHandle) { - var response = SendCommand(expertHandle, Mt5CommandType.GetQuote); - return response; + Log?.Debug($"GetQuote: expertHandle = {expertHandle}"); + var quote = SendCommand(expertHandle, Mt5CommandType.GetQuote); + if (quote != null) + quote.ExpertHandle = expertHandle; + return quote; } - private void _client_OnConnectionFailed(object? sender, EventArgs e) + private void Client_OnConnectionFailed(object? sender, EventArgs e) { + Log?.Info("Received connection failed"); Disconnect(true); } - private void _client_Disconnected(object? sender, EventArgs e) + private void Client_Disconnected(object? sender, EventArgs e) { + Log?.Info("Received normal disconnection"); Disconnect(false); } - //private void _client_MtEventReceived(MtEvent e) - //{ - // var eventType = (Mt5EventTypes)e.EventType; - // _mtEventHandlers[eventType](e.ExpertHandle, e.Payload); - //} - private void ReceivedOnTradeTransactionEvent(int expertHandler, string payload) { var e = JsonConvert.DeserializeObject(payload); @@ -3515,6 +3577,9 @@ namespace MtApi5 var e = JsonConvert.DeserializeObject(payload); if (e == null || string.IsNullOrEmpty(e.Instrument) || e.Tick == null) return; + + QuoteUpdated?.Invoke(this, e.Instrument, e.Tick.bid, e.Tick.ask); + var quote = new Mt5Quote(e.Instrument, e.Tick.bid, e.Tick.ask) { ExpertHandle = expertHandler, @@ -3522,8 +3587,6 @@ namespace MtApi5 Time = e.Tick.time, Last = e.Tick.last }; - - QuoteUpdated?.Invoke(this, quote.Instrument, quote.Bid, quote.Ask); QuoteUpdate?.Invoke(this, new Mt5QuoteEventArgs(quote)); } @@ -3609,16 +3672,19 @@ namespace MtApi5 var payloadJson = JsonConvert.SerializeObject(payload); MtCommand command = new(expertHandle, (int)commandType, _commandId++, payloadJson); - _tasks[command.CommandId] = new CommandTask(); - - Log?.Debug($"SendCommand: {command.CommandId}"); - client.Send(command); - - var responseJson = _tasks[command.CommandId].WaitResponse(10000); // 10 sec - if (_tasks.Remove(command.CommandId) == false) + CommandTask commandTask = new(); + lock (_locker) { - Log?.Warn($"SendCommand: task {command.CommandId} is not found in collection"); - return default; + _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); } Log?.Debug($"SendCommand: received response JSON [{responseJson}]"); @@ -3722,32 +3788,6 @@ namespace MtApi5 // QuoteUpdated?.Invoke(this, quote.Instrument, quote.Bid, quote.Ask); //} - //private void _client_ServerDisconnected(object sender, EventArgs e) - //{ - // Disconnect(false); - //} - - //private void _client_ServerFailed(object sender, EventArgs e) - //{ - // Disconnect(true); - //} - - //private void _client_QuoteRemoved(MtQuote quote) - //{ - // if (quote != null) - // { - // QuoteRemoved?.Invoke(this, new Mt5QuoteEventArgs(new Mt5Quote(quote))); - // } - //} - - //private void _client_QuoteAdded(MtQuote quote) - //{ - // if (quote != null) - // { - // QuoteAdded?.Invoke(this, new Mt5QuoteEventArgs(new Mt5Quote(quote))); - // } - //} - private void OnConnected() { Log?.Debug("OnConnected: begin");