Moved logic related to command to MtRpcClient

This commit is contained in:
Viacheslav Demydiuk
2024-01-12 23:24:47 +02:00
parent 06082539c3
commit 4f9dc57c25
5 changed files with 165 additions and 143 deletions
+1 -3
View File
@@ -1,6 +1,4 @@
using System; namespace MtApi5
namespace MtApi5
{ {
public class Mt5QuoteEventArgs: EventArgs public class Mt5QuoteEventArgs: EventArgs
{ {
+8
View File
@@ -0,0 +1,8 @@
namespace MtApi5
{
public class Mt5QuotesEventArgs(IEnumerable<Mt5Quote> quotes) : EventArgs
{
public IEnumerable<Mt5Quote> Quotes { get; } = quotes;
}
}
+28 -125
View File
@@ -3,9 +3,7 @@ using MtApi5.Requests;
using Newtonsoft.Json; using Newtonsoft.Json;
using MtApi5.Events; using MtApi5.Events;
using MtClient; using MtClient;
using System.ComponentModel.Design;
using System.Reflection.Metadata; using System.Reflection.Metadata;
using System.Collections.Generic;
namespace MtApi5 namespace MtApi5
{ {
@@ -40,8 +38,6 @@ namespace MtApi5
private HashSet<int> _experts = []; private HashSet<int> _experts = [];
private Dictionary<int, Mt5Quote> _quotes = []; private Dictionary<int, Mt5Quote> _quotes = [];
private readonly EventWaitHandle _quotesWaiter = new AutoResetEvent(false); private readonly EventWaitHandle _quotesWaiter = new AutoResetEvent(false);
private int _commandId = 0;
private readonly Dictionary<int, CommandTask> _tasks = [];
#endregion #endregion
#region Public Methods #region Public Methods
@@ -3319,7 +3315,7 @@ namespace MtApi5
public event EventHandler<Mt5BookEventArgs>? OnBookEvent; public event EventHandler<Mt5BookEventArgs>? OnBookEvent;
public event EventHandler<Mt5TimeBarArgs>? OnLastTimeBar; public event EventHandler<Mt5TimeBarArgs>? OnLastTimeBar;
public event EventHandler<Mt5LockTicksEventArgs>? OnLockTicks; public event EventHandler<Mt5LockTicksEventArgs>? OnLockTicks;
public event EventHandler<IEnumerable<Mt5Quote>>? QuoteList; public event EventHandler<Mt5QuotesEventArgs>? QuoteList;
#endregion #endregion
#region Private Methods #region Private Methods
@@ -3351,7 +3347,10 @@ namespace MtApi5
ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(Mt5ConnectionState.Connecting, message)); ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(Mt5ConnectionState.Connecting, message));
var client = new MtRpcClient(host, port); 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.ConnectionFailed += Client_OnConnectionFailed;
client.Disconnected += Client_Disconnected; client.Disconnected += Client_Disconnected;
@@ -3384,36 +3383,28 @@ namespace MtApi5
ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(state, message)); 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}"); Task.Run(() => _mtEventHandlers[(Mt5EventTypes)e.EventType](e.ExpertHandle, e.Payload));
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;
}
} }
private void ProcessExpertList(MtExpertListMsg? msg) private void Client_ExpertList(object? sender, MtExpertListEventArgs e)
{ {
if (msg == null) Task.Run(()=>ProcessExpertList(e.Experts));
return; }
HashSet<int> 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<int> experts)
{
if (experts == null || experts.Count == 0) if (experts == null || experts.Count == 0)
{ {
Log?.Warn("ProcessExpertList: expert list invalid or empty"); Log?.Warn("ProcessExpertList: expert list invalid or empty");
@@ -3435,16 +3426,11 @@ namespace MtApi5
} }
_quotesWaiter.Set(); _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}"); Log?.Debug($"ProcessExpertAdded: {handle}");
bool added; bool added;
@@ -3473,13 +3459,8 @@ namespace MtApi5
Log?.Warn($"ProcessExpertAdded: expert handle {handle} is already exist"); 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}"); Log?.Debug($"ProcessExpertRemoved: {handle}");
Mt5Quote? quote = null; Mt5Quote? quote = null;
@@ -3494,37 +3475,6 @@ namespace MtApi5
QuoteRemoved?.Invoke(this, new Mt5QuoteEventArgs(quote)); 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) private Mt5Quote? GetQuote(int expertHandle)
{ {
Log?.Debug($"GetQuote: expertHandle = {expertHandle}"); Log?.Debug($"GetQuote: expertHandle = {expertHandle}");
@@ -3625,7 +3575,6 @@ namespace MtApi5
_quotes.Clear(); _quotes.Clear();
_experts.Clear(); _experts.Clear();
_tasks.Clear();
} }
client?.Disconnect(); client?.Disconnect();
@@ -3635,31 +3584,6 @@ namespace MtApi5
ConnectionStateChanged?.Invoke(this, new Mt5ConnectionEventArgs(state, message)); 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<T>(int expertHandle, Mt5CommandType commandType, object? payload = null) private T? SendCommand<T>(int expertHandle, Mt5CommandType commandType, object? payload = null)
{ {
var client = Client; var client = Client;
@@ -3670,22 +3594,9 @@ namespace MtApi5
} }
var payloadJson = JsonConvert.SerializeObject(payload); var payloadJson = JsonConvert.SerializeObject(payload);
Log?.Debug($"SendCommand: sending '{payloadJson}' ...");
MtCommand command = new(expertHandle, (int)commandType, _commandId++, payloadJson); var responseJson = client.SendCommand(expertHandle, (int)commandType, 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);
}
Log?.Debug($"SendCommand: received response JSON [{responseJson}]"); Log?.Debug($"SendCommand: received response JSON [{responseJson}]");
@@ -3780,19 +3691,11 @@ namespace MtApi5
return response.Value; 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() private void OnConnected()
{ {
Log?.Debug("OnConnected: begin"); Log?.Debug("OnConnected: begin");
Client?.Send(new MtNotification(MtNotificationType.ClientReady)); Client?.NotifyClientReady();
_isBacktestingMode = IsTesting(); _isBacktestingMode = IsTesting();
+11 -11
View File
@@ -1,6 +1,6 @@
namespace MtClient namespace MtClient
{ {
public enum MessageType internal enum MessageType
{ {
Command = 0, Command = 0,
Response = 1, Response = 1,
@@ -11,12 +11,12 @@
Notification = 6 Notification = 6
} }
public enum MtNotificationType internal enum MtNotificationType
{ {
ClientReady = 0 ClientReady = 0
} }
public abstract class MtMessage internal abstract class MtMessage
{ {
public abstract MessageType MsgType { get; } public abstract MessageType MsgType { get; }
@@ -28,7 +28,7 @@
protected abstract string GetMessageBody(); 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; 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; public override MessageType MsgType => MessageType.Notification;
@@ -55,7 +55,7 @@
public MtNotificationType NotificationType { private set; get; } = notificationType; 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; 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; 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; public override MessageType MsgType => MessageType.ExpertRemoved;
@@ -119,7 +119,7 @@
} }
} }
public class MtExpertListMsg(HashSet<int> experts) : MtMessage internal class MtExpertListMsg(HashSet<int> experts) : MtMessage
{ {
public override MessageType MsgType => MessageType.ExpertList; 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; public override MessageType MsgType => MessageType.Response;
@@ -177,7 +177,7 @@
} }
} }
public static class MtMessageParser internal static class MtMessageParser
{ {
static MtMessageParser() static MtMessageParser()
{ {
+117 -4
View File
@@ -56,7 +56,35 @@ namespace MtClient
Log($"Disconnect: success"); 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_) lock (pendingMessages_)
{ {
@@ -164,8 +192,45 @@ namespace MtClient
Log("OnReceive: Failed parse message payload"); Log("OnReceive: Failed parse message payload");
return; 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) private void Log(string msg)
@@ -173,9 +238,12 @@ namespace MtClient
Console.WriteLine($"[{Environment.CurrentManagedThreadId}] {msg}"); Console.WriteLine($"[{Environment.CurrentManagedThreadId}] {msg}");
} }
public event EventHandler<MtMessage>? MessageReceived;
public event EventHandler<EventArgs>? ConnectionFailed; public event EventHandler<EventArgs>? ConnectionFailed;
public event EventHandler<EventArgs>? Disconnected; public event EventHandler<EventArgs>? Disconnected;
public event EventHandler<MtExpertListEventArgs>? ExpertList;
public event EventHandler<MtExpertEventArgs>? ExpertAdded;
public event EventHandler<MtExpertEventArgs>? ExpertRemoved;
public event EventHandler<MtEventArgs>? MtEventReceived;
private readonly ClientWebSocket ws_ = new(); private readonly ClientWebSocket ws_ = new();
private readonly string host_; private readonly string host_;
@@ -186,5 +254,50 @@ namespace MtClient
private readonly Thread receiveThread_; private readonly Thread receiveThread_;
private readonly Thread sendThread_; private readonly Thread sendThread_;
private readonly EventWaitHandle sendWaiter_ = new AutoResetEvent(false); private readonly EventWaitHandle sendWaiter_ = new AutoResetEvent(false);
private int nextCommandId = 0;
private readonly Dictionary<int, CommandTask> _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<int> experts) : EventArgs
{
public HashSet<int> Experts { get; } = experts;
}
public class MtExpertEventArgs(int expert) : EventArgs
{
public int Expert { get; }= expert;
} }
} }