using MtApi5; using Microsoft.AspNetCore.Mvc; using System.Reflection; using System.Globalization; using System.Collections.Concurrent; using System.Net.WebSockets; using System.Text; using Newtonsoft.Json; var builder = WebApplication.CreateBuilder(args); builder.WebHost.UseUrls("http://0.0.0.0:8080"); builder.Services.AddSingleton(sp => { var client = new MtApi5Client(); try { var prop = typeof(MtApi5Client).GetProperty("CommandTimeout"); if (prop != null && prop.CanWrite) { prop.SetValue(client, 120000); } else { var field = typeof(MtApi5Client).GetField("_command_timeout", BindingFlags.NonPublic | BindingFlags.Instance); field?.SetValue(client, 120000); } Console.WriteLine("CommandTimeout set to 120000ms"); } catch (Exception ex) { Console.WriteLine($"Could not set CommandTimeout: {ex.Message}"); } return client; }); builder.Services.AddSingleton(new SemaphoreSlim(1, 1)); builder.Services.AddHostedService(); builder.Services.AddCors(options => { options.AddDefaultPolicy(policy => { policy.AllowAnyOrigin() .AllowAnyMethod() .AllowAnyHeader(); }); }); var app = builder.Build(); app.UseCors(); app.UseWebSockets(); app.UseExceptionHandler(errorApp => { errorApp.Run(async context => { var logger = context.RequestServices.GetRequiredService().CreateLogger("GlobalExceptionHandler"); var exception = context.Features.Get()?.Error; if (exception != null) { logger.LogError(exception, "Unhandled exception while processing {Method} {Path}", context.Request.Method, context.Request.Path); } context.Response.StatusCode = StatusCodes.Status500InternalServerError; await context.Response.WriteAsJsonAsync(new { detail = "The bridge could not process the request.", status = StatusCodes.Status500InternalServerError }); }); }); // API Key 认证 const string API_KEY = "UiHMqtaYLZzwBdcuS4RFmEGhgDO8N2eI"; app.Use(async (context, next) => { var key = context.Request.Headers["X-API-Key"].FirstOrDefault() ?? context.Request.Query["key"].FirstOrDefault(); if (key != API_KEY) { context.Response.StatusCode = 401; await context.Response.WriteAsync("Unauthorized"); return; } await next(); }); app.MapGet("/health", (MtApi5Client mt) => { var connected = mt.ConnectionState == Mt5ConnectionState.Connected; return connected ? Results.Ok(new { status = "healthy", mt5_connected = true, mt5_version = "unknown", api_version = "1.0.0" }) : Results.Json(new { status = "disconnected", mt5_connected = false, mt5_version = "unknown", api_version = "1.0.0" }, statusCode: StatusCodes.Status503ServiceUnavailable); }); app.MapGet("/account", (MtApi5Client mt) => { var login = (ulong)mt.AccountInfoInteger(ENUM_ACCOUNT_INFO_INTEGER.ACCOUNT_LOGIN); var leverage = (uint)mt.AccountInfoInteger(ENUM_ACCOUNT_INFO_INTEGER.ACCOUNT_LEVERAGE); var tradeAllowed = mt.AccountInfoInteger(ENUM_ACCOUNT_INFO_INTEGER.ACCOUNT_TRADE_ALLOWED) != 0; var tradeExpert = mt.AccountInfoInteger(ENUM_ACCOUNT_INFO_INTEGER.ACCOUNT_TRADE_EXPERT) != 0; var currency = mt.AccountInfoString(ENUM_ACCOUNT_INFO_STRING.ACCOUNT_CURRENCY) ?? "USD"; var server = mt.AccountInfoString(ENUM_ACCOUNT_INFO_STRING.ACCOUNT_SERVER) ?? ""; var name = mt.AccountInfoString(ENUM_ACCOUNT_INFO_STRING.ACCOUNT_NAME) ?? ""; var company = mt.AccountInfoString(ENUM_ACCOUNT_INFO_STRING.ACCOUNT_COMPANY) ?? ""; var balance = mt.AccountInfoDouble(ENUM_ACCOUNT_INFO_DOUBLE.ACCOUNT_BALANCE); var equity = mt.AccountInfoDouble(ENUM_ACCOUNT_INFO_DOUBLE.ACCOUNT_EQUITY); var profit = mt.AccountInfoDouble(ENUM_ACCOUNT_INFO_DOUBLE.ACCOUNT_PROFIT); var margin = mt.AccountInfoDouble(ENUM_ACCOUNT_INFO_DOUBLE.ACCOUNT_MARGIN); var marginFree = mt.AccountInfoDouble(ENUM_ACCOUNT_INFO_DOUBLE.ACCOUNT_MARGIN_FREE); var marginLevel = mt.AccountInfoDouble(ENUM_ACCOUNT_INFO_DOUBLE.ACCOUNT_MARGIN_LEVEL); return Results.Ok(new { data = new[] { new { login, leverage, trade_allowed = tradeAllowed, trade_expert = tradeExpert, currency, currency_digits = 2, server, name, company, balance, equity, profit, margin, margin_free = marginFree, margin_level = marginLevel } }, count = 1, format = "json" }); }); app.MapGet("/symbols/{symbol}", (string symbol, MtApi5Client mt, ILoggerFactory loggerFactory) => { if (string.IsNullOrWhiteSpace(symbol)) return Results.BadRequest(new { detail = "Symbol is required." }); try { // ponytail: SymbolSelect(true) subscribes to MT5 quotes for this symbol — without it // SymbolInfoDouble returns 0 for bid/ask until a tick arrives naturally. mt.SymbolSelect(symbol, true); var digits = (int)mt.SymbolInfoInteger(symbol, ENUM_SYMBOL_INFO_INTEGER.SYMBOL_DIGITS); var spreadFloat = mt.SymbolInfoInteger(symbol, ENUM_SYMBOL_INFO_INTEGER.SYMBOL_SPREAD_FLOAT) != 0; var spread = (int)mt.SymbolInfoInteger(symbol, ENUM_SYMBOL_INFO_INTEGER.SYMBOL_SPREAD); var point = mt.SymbolInfoDouble(symbol, ENUM_SYMBOL_INFO_DOUBLE.SYMBOL_POINT); double bid = 0, ask = 0; var tickInfo = WaitForTick(mt, symbol); if (tickInfo != null) { bid = tickInfo.bid; ask = tickInfo.ask; } if (bid == 0) bid = mt.SymbolInfoDouble(symbol, ENUM_SYMBOL_INFO_DOUBLE.SYMBOL_BID); if (ask == 0) ask = mt.SymbolInfoDouble(symbol, ENUM_SYMBOL_INFO_DOUBLE.SYMBOL_ASK); var volumeMin = mt.SymbolInfoDouble(symbol, ENUM_SYMBOL_INFO_DOUBLE.SYMBOL_VOLUME_MIN); var volumeMax = mt.SymbolInfoDouble(symbol, ENUM_SYMBOL_INFO_DOUBLE.SYMBOL_VOLUME_MAX); var volumeStep = mt.SymbolInfoDouble(symbol, ENUM_SYMBOL_INFO_DOUBLE.SYMBOL_VOLUME_STEP); var contractSize = mt.SymbolInfoDouble(symbol, ENUM_SYMBOL_INFO_DOUBLE.SYMBOL_TRADE_CONTRACT_SIZE); string description = "", currencyBase = "", currencyProfit = "", path = ""; try { description = mt.SymbolInfoString(symbol, ENUM_SYMBOL_INFO_STRING.SYMBOL_DESCRIPTION) ?? ""; } catch { } try { currencyBase = mt.SymbolInfoString(symbol, ENUM_SYMBOL_INFO_STRING.SYMBOL_CURRENCY_BASE) ?? ""; } catch { } try { currencyProfit = mt.SymbolInfoString(symbol, ENUM_SYMBOL_INFO_STRING.SYMBOL_CURRENCY_PROFIT) ?? ""; } catch { } try { path = mt.SymbolInfoString(symbol, ENUM_SYMBOL_INFO_STRING.SYMBOL_PATH) ?? ""; } catch { } return Results.Ok(new { data = new[] { new { name = symbol, description, digits, point, bid, ask, spread, spread_float = spreadFloat, volume_min = volumeMin, volume_max = volumeMax, volume_step = volumeStep, trade_contract_size = contractSize, currency_base = currencyBase, currency_profit = currencyProfit, category = path } }, count = 1, format = "json" }); } catch (Exception ex) { loggerFactory.CreateLogger("SymbolsEndpoint").LogError(ex, "Failed to read symbol info for {Symbol}", symbol); return Results.Problem("Failed to read symbol info.", statusCode: StatusCodes.Status502BadGateway); } }); app.MapGet("/symbols/{symbol}/tick", (string symbol, MtApi5Client mt) => { if (string.IsNullOrWhiteSpace(symbol)) return Results.BadRequest(new { detail = "Symbol is required." }); mt.SymbolSelect(symbol, true); var tick = WaitForTick(mt, symbol); if (tick == null) return Results.NotFound(new { detail = $"Tick for {symbol} not found" }); return Results.Json(new { data = new[] { new { time = tick.time.ToString("yyyy-MM-ddTHH:mm:ss"), bid = tick.bid, ask = tick.ask, last = tick.last, volume = tick.volume, time_msc = tick.time.ToString("yyyy-MM-ddTHH:mm:ss.fff000"), flags = 6, volume_real = tick.volume_real } }, count = 1, format = "json" }); }); // Tick stream via WebSocket. Connect: ws://host:8080/stream/ticks/{symbol} // MT5 must be in Market Watch for the symbol; this endpoint auto-adds it on connect // and removes it when the last subscriber disconnects. app.Map("/stream/ticks/{symbol}", async (string symbol, HttpContext ctx, MtApi5Client mt, SemaphoreSlim gate, ILoggerFactory loggerFactory) => { var logger = loggerFactory.CreateLogger("TickStream"); if (string.IsNullOrWhiteSpace(symbol)) { ctx.Response.StatusCode = 400; await ctx.Response.WriteAsJsonAsync(new { detail = "Symbol is required." }); return; } if (!ctx.WebSockets.IsWebSocketRequest) { ctx.Response.StatusCode = 400; await ctx.Response.WriteAsJsonAsync(new { detail = "WebSocket required. Use ws://.../stream/ticks/{symbol}." }); return; } if (mt.ConnectionState != Mt5ConnectionState.Connected) { ctx.Response.StatusCode = 503; await ctx.Response.WriteAsJsonAsync(new { detail = "MT5 is not connected." }); return; } await gate.WaitAsync(ctx.RequestAborted); try { if (!mt.SymbolSelect(symbol, true)) { ctx.Response.StatusCode = 404; await ctx.Response.WriteAsJsonAsync(new { detail = $"Symbol '{symbol}' not found." }); return; } } finally { gate.Release(); } var subs = MtConnectionService.TickSubscribers.GetOrAdd(symbol, _ => new ConcurrentDictionary()); if (subs.Count >= MtConnectionService.MaxSubscribersPerSymbol) { ctx.Response.StatusCode = 429; await ctx.Response.WriteAsJsonAsync(new { detail = $"Too many subscribers for {symbol} (limit={MtConnectionService.MaxSubscribersPerSymbol})." }); return; } using var ws = await ctx.WebSockets.AcceptWebSocketAsync(); var subscriber = new WsSubscriber(ws); subs[subscriber] = 1; logger.LogInformation("Tick stream subscribed for {Symbol} via WS (subscribers={Count})", symbol, subs.Count); var buf = new byte[1024]; try { while (ws.State == WebSocketState.Open && !ctx.RequestAborted.IsCancellationRequested) { var result = await ws.ReceiveAsync(buf, ctx.RequestAborted); if (result.MessageType == WebSocketMessageType.Close) { await ws.CloseAsync(WebSocketCloseStatus.NormalClosure, "bye", CancellationToken.None); break; } } } catch (OperationCanceledException) { } catch (WebSocketException) { } finally { subs.TryRemove(subscriber, out _); MtConnectionService.TickSubscribers.TryRemove(symbol, out _); if (subs.IsEmpty) { try { if (mt.ConnectionState == Mt5ConnectionState.Connected) { await gate.WaitAsync(); try { mt.SymbolSelect(symbol, false); } finally { gate.Release(); } } } catch (Exception ex) { logger.LogWarning(ex, "Failed to remove {Symbol} from Market Watch", symbol); } } logger.LogInformation("Tick stream disconnected for {Symbol} via WS (remaining={Count})", symbol, subs.Count); } }); // Tick stream via Server-Sent Events. Same semantics as /stream/ticks/{symbol}, // one-way push over HTTP for clients that can't open a WebSocket. // ponytail: SSE = no receive loop, just keep the response open and let the broadcast // loop (OnQuoteUpdate + heartbeat timer) push data frames. app.MapGet("/stream/ticks-sse/{symbol}", async (string symbol, HttpContext ctx, MtApi5Client mt, SemaphoreSlim gate, ILoggerFactory loggerFactory) => { var logger = loggerFactory.CreateLogger("TickStream"); if (string.IsNullOrWhiteSpace(symbol)) { ctx.Response.StatusCode = 400; await ctx.Response.WriteAsJsonAsync(new { detail = "Symbol is required." }); return; } if (mt.ConnectionState != Mt5ConnectionState.Connected) { ctx.Response.StatusCode = 503; await ctx.Response.WriteAsJsonAsync(new { detail = "MT5 is not connected." }); return; } await gate.WaitAsync(ctx.RequestAborted); try { if (!mt.SymbolSelect(symbol, true)) { ctx.Response.StatusCode = 404; await ctx.Response.WriteAsJsonAsync(new { detail = $"Symbol '{symbol}' not found." }); return; } } finally { gate.Release(); } var subs = MtConnectionService.TickSubscribers.GetOrAdd(symbol, _ => new ConcurrentDictionary()); if (subs.Count >= MtConnectionService.MaxSubscribersPerSymbol) { ctx.Response.StatusCode = 429; await ctx.Response.WriteAsJsonAsync(new { detail = $"Too many subscribers for {symbol} (limit={MtConnectionService.MaxSubscribersPerSymbol})." }); return; } ctx.Response.StatusCode = 200; ctx.Response.ContentType = "text/event-stream"; ctx.Response.Headers.CacheControl = "no-cache"; ctx.Response.Headers.Connection = "keep-alive"; ctx.Response.Headers["X-Accel-Buffering"] = "no"; // disable nginx buffering var subscriber = new SseSubscriber(ctx.Response.Body); subs[subscriber] = 1; logger.LogInformation("Tick stream subscribed for {Symbol} via SSE (subscribers={Count})", symbol, subs.Count); try { // Initial comment to flush headers and confirm connection await ctx.Response.WriteAsync(": connected\n\n", ctx.RequestAborted); await ctx.Response.Body.FlushAsync(ctx.RequestAborted); // Hold the request open until client disconnects; ticks/heartbeats arrive via OnQuoteUpdate + timer. await Task.Delay(Timeout.Infinite, ctx.RequestAborted); } catch (OperationCanceledException) { } catch (Exception ex) { logger.LogDebug(ex, "SSE stream ended for {Symbol}", symbol); } finally { subs.TryRemove(subscriber, out _); MtConnectionService.TickSubscribers.TryRemove(symbol, out _); if (subs.IsEmpty) { try { if (mt.ConnectionState == Mt5ConnectionState.Connected) { await gate.WaitAsync(); try { mt.SymbolSelect(symbol, false); } finally { gate.Release(); } } } catch (Exception ex) { logger.LogWarning(ex, "Failed to remove {Symbol} from Market Watch", symbol); } } logger.LogInformation("Tick stream disconnected for {Symbol} via SSE (remaining={Count})", symbol, subs.Count); } }); app.MapGet("/rates/from-pos", (string symbol, string timeframe, int start_pos, int count, MtApi5Client mt, ILoggerFactory loggerFactory) => { if (string.IsNullOrWhiteSpace(symbol)) return Results.BadRequest(new { detail = "Symbol is required." }); if (!TryParseTimeframe(timeframe, out var tf)) return Results.BadRequest(new { detail = $"Unsupported timeframe '{timeframe}'." }); if (start_pos < 0) return Results.BadRequest(new { detail = "start_pos must be greater than or equal to 0." }); if (count <= 0 || count > 10000) return Results.BadRequest(new { detail = "count must be between 1 and 10000." }); try { const int chunkSize = 1000; var allRates = new List(); int remaining = count; int currentPos = start_pos; while (remaining > 0) { int chunk = Math.Min(remaining, chunkSize); var result = mt.CopyRates(symbol, tf, currentPos, chunk, out MqlRates[]? rates); if (rates != null && rates.Length > 0) { allRates.AddRange(rates); currentPos += rates.Length; remaining -= rates.Length; if (rates.Length < chunk) break; } else { break; } } var data = allRates.Select(r => new { time = r.time.ToString("yyyy-MM-ddTHH:mm:ss"), r.open, r.high, r.low, r.close, tick_volume = r.tick_volume, r.spread, real_volume = r.real_volume }).ToArray(); return Results.Ok(new { data, count = data.Length, format = "json" }); } catch (Exception ex) { loggerFactory.CreateLogger("RatesEndpoint").LogError(ex, "Failed to read rates for {Symbol} {Timeframe}", symbol, timeframe); return Results.Problem("Failed to read rates.", statusCode: StatusCodes.Status502BadGateway); } }); app.MapGet("/rates/from-date", (string symbol, string timeframe, string date_from, string date_to, MtApi5Client mt, ILoggerFactory loggerFactory) => { if (string.IsNullOrWhiteSpace(symbol)) return Results.BadRequest(new { detail = "Symbol is required." }); if (!TryParseTimeframe(timeframe, out var tf)) return Results.BadRequest(new { detail = $"Unsupported timeframe '{timeframe}'." }); if (!TryParseClientDate(date_from, out var fromDate)) return Results.BadRequest(new { detail = "date_from must be a valid ISO-8601 or yyyy-MM-dd value." }); if (!TryParseClientDate(date_to, out var toDate)) return Results.BadRequest(new { detail = "date_to must be a valid ISO-8601 or yyyy-MM-dd value." }); if (fromDate > toDate) return Results.BadRequest(new { detail = "date_from must be earlier than or equal to date_to." }); try { // ponytail: Bars(from, to) is half-open [from, to). Same-day query would yield 0 // without this nudge — bump toDate by one day to include the entire end day. var effectiveTo = (toDate.Date == fromDate.Date) ? toDate.Date.AddDays(1) : toDate; var total = mt.Bars(symbol, tf, fromDate, effectiveTo); if (total <= 0) return Results.Ok(new { data = Array.Empty(), count = 0, format = "json" }); if (total > 10000) return Results.BadRequest(new { detail = $"Date range covers {total} bars (max 10000). Narrow the range." }); var result = mt.CopyRates(symbol, tf, total - 1, total, out MqlRates[]? rates); if (rates == null || rates.Length == 0) return Results.Ok(new { data = Array.Empty(), count = 0, format = "json" }); var data = rates.Select(r => new { time = r.time.ToString("yyyy-MM-ddTHH:mm:ss"), r.open, r.high, r.low, r.close, tick_volume = r.tick_volume, r.spread, real_volume = r.real_volume }).ToArray(); return Results.Ok(new { data, count = data.Length, format = "json" }); } catch (Exception ex) { loggerFactory.CreateLogger("RatesFromDateEndpoint").LogError(ex, "Failed to read rates for {Symbol} {Timeframe} {DateFrom}~{DateTo}", symbol, timeframe, date_from, date_to); return Results.Problem("Failed to read rates.", statusCode: StatusCodes.Status502BadGateway); } }); app.MapGet("/positions", (string? symbol, MtApi5Client mt) => { var total = mt.PositionsTotal(); var positions = new List(); for (int i = 0; i < total; i++) { var ticket = mt.PositionGetTicket(i); if (ticket == 0) continue; var posSymbol = mt.PositionGetString(ENUM_POSITION_PROPERTY_STRING.POSITION_SYMBOL); if (!string.IsNullOrWhiteSpace(symbol) && posSymbol != symbol) continue; positions.Add(new { ticket, symbol = posSymbol, type = (int)mt.PositionGetInteger(ENUM_POSITION_PROPERTY_INTEGER.POSITION_TYPE), volume = mt.PositionGetDouble(ENUM_POSITION_PROPERTY_DOUBLE.POSITION_VOLUME), price_open = mt.PositionGetDouble(ENUM_POSITION_PROPERTY_DOUBLE.POSITION_PRICE_OPEN), sl = mt.PositionGetDouble(ENUM_POSITION_PROPERTY_DOUBLE.POSITION_SL), tp = mt.PositionGetDouble(ENUM_POSITION_PROPERTY_DOUBLE.POSITION_TP), price_current = mt.PositionGetDouble(ENUM_POSITION_PROPERTY_DOUBLE.POSITION_PRICE_CURRENT), swap = mt.PositionGetDouble(ENUM_POSITION_PROPERTY_DOUBLE.POSITION_SWAP), profit = mt.PositionGetDouble(ENUM_POSITION_PROPERTY_DOUBLE.POSITION_PROFIT), comment = mt.PositionGetString(ENUM_POSITION_PROPERTY_STRING.POSITION_COMMENT), magic = (ulong)mt.PositionGetInteger(ENUM_POSITION_PROPERTY_INTEGER.POSITION_MAGIC) }); } return Results.Ok(new { data = positions, count = positions.Count, format = "json" }); }); app.MapGet("/orders", (string? symbol, MtApi5Client mt) => { var total = mt.OrdersTotal(); var orders = new List(); for (int i = 0; i < total; i++) { var ticket = mt.OrderGetTicket(i); if (ticket == 0) continue; var orderSymbol = mt.OrderGetString(ENUM_ORDER_PROPERTY_STRING.ORDER_SYMBOL); if (!string.IsNullOrWhiteSpace(symbol) && orderSymbol != symbol) continue; orders.Add(new { ticket, symbol = orderSymbol, type = (int)mt.OrderGetInteger(ENUM_ORDER_PROPERTY_INTEGER.ORDER_TYPE), volume_initial = mt.OrderGetDouble(ENUM_ORDER_PROPERTY_DOUBLE.ORDER_VOLUME_INITIAL), price_open = mt.OrderGetDouble(ENUM_ORDER_PROPERTY_DOUBLE.ORDER_PRICE_OPEN), sl = mt.OrderGetDouble(ENUM_ORDER_PROPERTY_DOUBLE.ORDER_SL), tp = mt.OrderGetDouble(ENUM_ORDER_PROPERTY_DOUBLE.ORDER_TP), magic = (ulong)mt.OrderGetInteger(ENUM_ORDER_PROPERTY_INTEGER.ORDER_MAGIC), comment = mt.OrderGetString(ENUM_ORDER_PROPERTY_STRING.ORDER_COMMENT) }); } return Results.Ok(new { data = orders, count = orders.Count, format = "json" }); }); app.MapPost("/order/check", async (TradeRequestDto body, MtApi5Client mt, SemaphoreSlim gate) => { var validationError = ValidateTradeRequest(body); if (validationError != null) return Results.BadRequest(new { detail = validationError }); body.type_filling = (uint)ResolveFillingMode(mt, body.symbol, body.type_filling ?? 0); var request = ToMqlTradeRequest(body); await gate.WaitAsync(); try { mt.OrderCheck(request, out MqlTradeCheckResult? checkResult); if (checkResult == null) return Results.Ok(new { data = new { retcode = 0u, balance = 0.0, equity = 0.0, profit = 0.0, margin = 0.0, margin_free = 0.0, margin_level = 0.0, comment = "Check failed" }, count = 1, format = "json" }); return Results.Ok(new { data = new { retcode = checkResult.Retcode, balance = checkResult.Balance, equity = checkResult.Equity, profit = checkResult.Profit, margin = checkResult.Margin, margin_free = checkResult.Margin_free, margin_level = checkResult.Margin_level, comment = checkResult.Comment ?? "" }, count = 1, format = "json" }); } finally { gate.Release(); } }); app.MapPost("/order/send", async (TradeRequestBody body, MtApi5Client mt, SemaphoreSlim gate) => { var validationError = ValidateTradeRequest(body.request); if (validationError != null) return Results.BadRequest(new { detail = validationError }); body.request.type_filling = (uint)ResolveFillingMode(mt, body.request.symbol, body.request.type_filling ?? 0); var request = ToMqlTradeRequest(body.request); await gate.WaitAsync(); try { mt.OrderSend(request, out MqlTradeResult? result); if (result == null) return Results.Ok(new { data = new { retcode = 0u, order = 0ul, comment = "Send failed" }, count = 1, format = "json" }); return Results.Ok(new { data = new { retcode = result.Retcode, order = result.Order, comment = result.Comment ?? "" }, count = 1, format = "json" }); } finally { gate.Release(); } }); app.MapPost("/position/close", async (ClosePositionDto body, MtApi5Client mt, SemaphoreSlim gate) => { if (body == null || body.ticket == 0) return Results.BadRequest(new { detail = "ticket is required." }); await gate.WaitAsync(); try { bool ok; if (body.volume is > 0) { // ponytail: PositionClosePartial has no MqlTradeResult overload — return bool only. ok = mt.PositionClosePartial(body.ticket, body.volume.Value, body.deviation ?? 10); } else { mt.PositionClose(body.ticket, body.deviation ?? 10, out MqlTradeResult? result); if (result == null) return Results.Ok(new { data = new { retcode = 0u, order = 0ul, comment = "Close failed" }, count = 1, format = "json" }); return Results.Ok(new { data = new { retcode = result.Retcode, order = result.Order, comment = result.Comment ?? "" }, count = 1, format = "json" }); } return Results.Ok(new { data = new { retcode = ok ? 10009u : 10004u, order = body.ticket, comment = ok ? (body.volume is > 0 ? "Partial closed" : "Closed") : "Close failed" }, count = 1, format = "json" }); } finally { gate.Release(); } }); app.MapPost("/position/modify", async (ModifyPositionDto body, MtApi5Client mt, SemaphoreSlim gate) => { if (body == null || body.ticket == 0) return Results.BadRequest(new { detail = "ticket is required." }); if (body.sl < 0 || body.tp < 0) return Results.BadRequest(new { detail = "sl/tp must be non-negative." }); await gate.WaitAsync(); try { // ponytail: MtApi5 PositionModify returns bool, no MqlTradeResult — pack into our uniform envelope. var ok = mt.PositionModify(body.ticket, body.sl, body.tp); return Results.Ok(new { data = new { retcode = ok ? 10009u : 10004u, order = body.ticket, comment = ok ? "Modified" : "Modify failed" }, count = 1, format = "json" }); } finally { gate.Release(); } }); app.MapPost("/order/cancel", async (OrderTicketDto body, MtApi5Client mt, SemaphoreSlim gate) => { if (body == null || body.ticket == 0) return Results.BadRequest(new { detail = "ticket is required." }); await gate.WaitAsync(); try { // ponytail: MtApi5 has no OrderDelete — reuse OrderSend with TRADE_ACTION_REMOVE. var req = new MqlTradeRequest { Action = ENUM_TRADE_REQUEST_ACTIONS.TRADE_ACTION_REMOVE, Order = body.ticket, }; mt.OrderSend(req, out MqlTradeResult? result); if (result == null) return Results.Ok(new { data = new { retcode = 0u, order = 0ul, comment = "Cancel failed" }, count = 1, format = "json" }); return Results.Ok(new { data = new { retcode = result.Retcode, order = result.Order, comment = result.Comment ?? "" }, count = 1, format = "json" }); } finally { gate.Release(); } }); app.MapPost("/position/close-by", async (CloseByDto body, MtApi5Client mt, SemaphoreSlim gate) => { if (body == null || body.position == 0 || body.position_by == 0) return Results.BadRequest(new { detail = "position and position_by are required." }); await gate.WaitAsync(); try { // ponytail: TRADE_ACTION_CLOSE_BY lets MT5 net two opposite positions instead of // two separate market closes — saves one side of the spread. var req = new MqlTradeRequest { Action = ENUM_TRADE_REQUEST_ACTIONS.TRADE_ACTION_CLOSE_BY, Position = body.position, PositionBy = body.position_by, }; mt.OrderSend(req, out MqlTradeResult? result); if (result == null) return Results.Ok(new { data = new { retcode = 0u, order = 0ul, comment = "Close-by failed" }, count = 1, format = "json" }); return Results.Ok(new { data = new { retcode = result.Retcode, order = result.Order, comment = result.Comment ?? "" }, count = 1, format = "json" }); } finally { gate.Release(); } }); app.MapPost("/positions/close-batch", async (BatchCloseDto body, MtApi5Client mt, SemaphoreSlim gate, ILoggerFactory loggerFactory) => { if (body == null || (string.IsNullOrWhiteSpace(body.symbol) && body.magic == null)) return Results.BadRequest(new { detail = "At least one of `symbol` or `magic` is required." }); var logger = loggerFactory.CreateLogger("BatchClose"); var results = new List(); int closed = 0, failed = 0; // ponytail: hold the gate for the whole batch — one logical op, brief duration. // Clients wanting concurrent writes should not use batch close. await gate.WaitAsync(); try { var total = mt.PositionsTotal(); for (int i = 0; i < total; i++) { var ticket = mt.PositionGetTicket(i); if (ticket == 0) continue; var posSymbol = mt.PositionGetString(ENUM_POSITION_PROPERTY_STRING.POSITION_SYMBOL); if (!string.IsNullOrWhiteSpace(body.symbol) && posSymbol != body.symbol) continue; if (body.magic is ulong m && mt.PositionGetInteger(ENUM_POSITION_PROPERTY_INTEGER.POSITION_MAGIC) != (long)m) continue; mt.PositionClose(ticket, body.deviation ?? 10, out MqlTradeResult? result); if (result != null && result.Retcode == 10009) { closed++; results.Add(new { ticket, symbol = posSymbol, retcode = result.Retcode, comment = "closed" }); } else { failed++; var comment = result?.Comment ?? "unknown"; logger.LogWarning("Batch close failed for ticket {Ticket} ({Symbol}): {Comment}", ticket, posSymbol, comment); results.Add(new { ticket, symbol = posSymbol, retcode = result?.Retcode ?? 0u, comment }); } } return Results.Ok(new { data = results, count = results.Count, closed, failed, format = "json" }); } finally { gate.Release(); } }); app.MapPost("/order/modify", async (ModifyOrderDto body, MtApi5Client mt, SemaphoreSlim gate) => { if (body == null || body.ticket == 0) return Results.BadRequest(new { detail = "ticket is required." }); if (body.price <= 0) return Results.BadRequest(new { detail = "price must be positive." }); if (body.sl < 0 || body.tp < 0) return Results.BadRequest(new { detail = "sl/tp must be non-negative." }); await gate.WaitAsync(); try { var req = new MqlTradeRequest { Action = ENUM_TRADE_REQUEST_ACTIONS.TRADE_ACTION_MODIFY, Order = body.ticket, Price = body.price, Sl = body.sl, Tp = body.tp, }; mt.OrderSend(req, out MqlTradeResult? result); if (result == null) return Results.Ok(new { data = new { retcode = 0u, order = 0ul, comment = "Modify failed" }, count = 1, format = "json" }); return Results.Ok(new { data = new { retcode = result.Retcode, order = result.Order, comment = result.Comment ?? "" }, count = 1, format = "json" }); } finally { gate.Release(); } }); app.MapGet("/history/deals", (string date_from, string date_to, string? symbol, MtApi5Client mt) => { if (!TryParseClientDate(date_from, out var fromDt) || !TryParseClientDate(date_to, out var toDt)) return Results.BadRequest(new { detail = "date_from and date_to must be valid ISO-8601 or yyyy-MM-dd values." }); if (fromDt > toDt) return Results.BadRequest(new { detail = "date_from must be earlier than or equal to date_to." }); mt.HistorySelect(fromDt, toDt); var total = mt.HistoryDealsTotal(); var dealList = new List(); for (int i = 0; i < total; i++) { var ticket = mt.HistoryDealGetTicket(i); if (ticket == 0) continue; var dealSymbol = mt.HistoryDealGetString(ticket, ENUM_DEAL_PROPERTY_STRING.DEAL_SYMBOL); if (!string.IsNullOrWhiteSpace(symbol) && dealSymbol != symbol) continue; var dealTime = MtTimeToDateTime(mt.HistoryDealGetInteger(ticket, ENUM_DEAL_PROPERTY_INTEGER.DEAL_TIME)); dealList.Add(new { ticket, time = dealTime.ToString("yyyy-MM-ddTHH:mm:ss"), entry = (int)mt.HistoryDealGetInteger(ticket, ENUM_DEAL_PROPERTY_INTEGER.DEAL_ENTRY), magic = (ulong)mt.HistoryDealGetInteger(ticket, ENUM_DEAL_PROPERTY_INTEGER.DEAL_MAGIC), volume = mt.HistoryDealGetDouble(ticket, ENUM_DEAL_PROPERTY_DOUBLE.DEAL_VOLUME), price = mt.HistoryDealGetDouble(ticket, ENUM_DEAL_PROPERTY_DOUBLE.DEAL_PRICE), profit = mt.HistoryDealGetDouble(ticket, ENUM_DEAL_PROPERTY_DOUBLE.DEAL_PROFIT), commission = mt.HistoryDealGetDouble(ticket, ENUM_DEAL_PROPERTY_DOUBLE.DEAL_COMMISSION), swap = mt.HistoryDealGetDouble(ticket, ENUM_DEAL_PROPERTY_DOUBLE.DEAL_SWAP), symbol = dealSymbol, comment = mt.HistoryDealGetString(ticket, ENUM_DEAL_PROPERTY_STRING.DEAL_COMMENT) ?? "" }); } return Results.Ok(new { data = dealList, count = dealList.Count, format = "json" }); }); app.MapGet("/gvar", (HttpContext context) => { var mt = context.RequestServices.GetRequiredService(); var total = mt.GlobalVariablesTotal(); var vars = new List(); for (int i = 0; i < total; i++) { var name = mt.GlobalVariableName(i) ?? ""; var value = mt.GlobalVariableGet(name); vars.Add(new { name, value }); } return Results.Ok(new { data = vars, count = vars.Count, format = "json" }); }); app.MapGet("/gvar/{name}", (HttpContext context, string name) => { if (string.IsNullOrWhiteSpace(name)) return Results.BadRequest(new { detail = "Global variable name is required." }); var mt = context.RequestServices.GetRequiredService(); if (!mt.GlobalVariableCheck(name)) return Results.NotFound(new { detail = $"GlobalVariable '{name}' not found" }); var value = mt.GlobalVariableGet(name); return Results.Ok(new { name, value }); }); app.MapPost("/gvar/{name}", async (HttpContext context, string name, SemaphoreSlim gate) => { if (string.IsNullOrWhiteSpace(name)) return Results.BadRequest(new { detail = "Global variable name is required." }); var mt = context.RequestServices.GetRequiredService(); var body = await context.Request.ReadFromJsonAsync(); if (body == null || double.IsNaN(body.value) || double.IsInfinity(body.value)) return Results.BadRequest(new { detail = "Invalid body." }); await gate.WaitAsync(); try { mt.GlobalVariableSet(name, body.value); return Results.Ok(new { name, value = body.value }); } finally { gate.Release(); } }); app.MapDelete("/gvar/{name}", async (HttpContext context, string name, SemaphoreSlim gate) => { if (string.IsNullOrWhiteSpace(name)) return Results.BadRequest(new { detail = "Global variable name is required." }); var mt = context.RequestServices.GetRequiredService(); await gate.WaitAsync(); try { if (!mt.GlobalVariableCheck(name)) return Results.NotFound(new { detail = $"GlobalVariable '{name}' not found" }); mt.GlobalVariableDel(name); return Results.Ok(new { deleted = name }); } finally { gate.Release(); } }); app.Run(); DateTime MtTimeToDateTime(long mtTime) { return new DateTime(1970, 1, 1).AddSeconds(mtTime); } bool TryParseTimeframe(string timeframe, out ENUM_TIMEFRAMES parsed) { parsed = timeframe switch { "TIMEFRAME_M1" => ENUM_TIMEFRAMES.PERIOD_M1, "TIMEFRAME_M5" => ENUM_TIMEFRAMES.PERIOD_M5, "TIMEFRAME_M15" => ENUM_TIMEFRAMES.PERIOD_M15, "TIMEFRAME_M30" => ENUM_TIMEFRAMES.PERIOD_M30, "TIMEFRAME_H1" => ENUM_TIMEFRAMES.PERIOD_H1, "TIMEFRAME_H4" => ENUM_TIMEFRAMES.PERIOD_H4, "TIMEFRAME_D1" => ENUM_TIMEFRAMES.PERIOD_D1, _ => default }; return timeframe is "TIMEFRAME_M1" or "TIMEFRAME_M5" or "TIMEFRAME_M15" or "TIMEFRAME_M30" or "TIMEFRAME_H1" or "TIMEFRAME_H4" or "TIMEFRAME_D1"; } bool TryParseClientDate(string value, out DateTime parsed) { return DateTime.TryParse(value, CultureInfo.InvariantCulture, DateTimeStyles.AssumeLocal, out parsed) || DateTime.TryParseExact(value, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.AssumeLocal, out parsed); } MqlTick? WaitForTick(MtApi5Client mt, string symbol, int timeoutMs = 3000) { for (int i = 0; i < timeoutMs / 100; i++) { var tick = mt.SymbolInfoTick(symbol); if (tick != null && tick.bid > 0) return tick; Thread.Sleep(100); } return mt.SymbolInfoTick(symbol); } uint ResolveFillingMode(MtApi5Client mt, string symbol, uint requestedFilling) { try { var modeFlags = (long)mt.SymbolInfoInteger(symbol, ENUM_SYMBOL_INFO_INTEGER.SYMBOL_FILLING_MODE); bool fokSupported = (modeFlags & 1) != 0; bool iocSupported = (modeFlags & 2) != 0; bool requestedSupported = requestedFilling switch { 0 => fokSupported, 1 => iocSupported, 2 => true, _ => false }; if (requestedSupported) return requestedFilling; if (iocSupported) return 1; if (fokSupported) return 0; return 2; } catch { return requestedFilling; } } string? ValidateTradeRequest(TradeRequestDto request) { if (!Enum.IsDefined(typeof(ENUM_TRADE_REQUEST_ACTIONS), (int)request.action)) return "action is invalid."; if (string.IsNullOrWhiteSpace(request.symbol)) return "symbol is required."; if (request.volume is not null && (!double.IsFinite(request.volume.Value) || request.volume.Value <= 0)) return "volume must be a positive number."; if (request.price is not null && (!double.IsFinite(request.price.Value) || request.price.Value < 0)) return "price must be a non-negative number."; if (request.sl is not null && (!double.IsFinite(request.sl.Value) || request.sl.Value < 0)) return "sl must be a non-negative number."; if (request.tp is not null && (!double.IsFinite(request.tp.Value) || request.tp.Value < 0)) return "tp must be a non-negative number."; if (request.order_type is not null && !Enum.IsDefined(typeof(ENUM_ORDER_TYPE), (int)request.order_type.Value)) return "order_type is invalid."; if (request.deviation is not null && request.deviation.Value > 100000) return "deviation is unreasonably large."; return null; } MqlTradeRequest ToMqlTradeRequest(TradeRequestDto r) { return new MqlTradeRequest { Action = (ENUM_TRADE_REQUEST_ACTIONS)r.action, Symbol = r.symbol, Volume = r.volume ?? 0, Type = (ENUM_ORDER_TYPE)(r.order_type ?? 0), Price = r.price ?? 0, Sl = r.sl ?? 0, Tp = r.tp ?? 0, Magic = r.magic ?? 0, Comment = r.comment ?? "", Order = r.order ?? 0, Position = r.position ?? 0, Deviation = r.deviation ?? 0, Type_filling = (ENUM_ORDER_TYPE_FILLING)(r.type_filling ?? 0), Type_time = (ENUM_ORDER_TYPE_TIME)(r.type_time ?? 0), }; } class TradeRequestDto { public uint action { get; set; } public string symbol { get; set; } = ""; public double? volume { get; set; } public uint? order_type { get; set; } public double? price { get; set; } public double? sl { get; set; } public double? tp { get; set; } public ulong? magic { get; set; } public string? comment { get; set; } public ulong? order { get; set; } public ulong? position { get; set; } public ulong? deviation { get; set; } public uint? type_filling { get; set; } // 0=FOK, 1=IOC, 2=RETURN — 经纪商常要求显式指定 public uint? type_time { get; set; } // 0=GTC, 1=DAY, 2=SPECIFIED, 3=SPECIFIED_DAY } class TradeRequestBody { public TradeRequestDto request { get; set; } = new(); } class GvarSetDto { public double value { get; set; } } class ClosePositionDto { public ulong ticket { get; set; } public double? volume { get; set; } // null/0 = 全平;>0 = 部分平仓 public ulong? deviation { get; set; } } class ModifyPositionDto { public ulong ticket { get; set; } public double sl { get; set; } public double tp { get; set; } } class OrderTicketDto { public ulong ticket { get; set; } } class ModifyOrderDto { public ulong ticket { get; set; } public double price { get; set; } public double sl { get; set; } public double tp { get; set; } } class CloseByDto { public ulong position { get; set; } public ulong position_by { get; set; } } class BatchCloseDto { public string? symbol { get; set; } public ulong? magic { get; set; } public ulong? deviation { get; set; } } abstract class TickSubscriber { public abstract ValueTask SendAsync(ReadOnlyMemory payload); public abstract bool IsAlive { get; } } sealed class WsSubscriber : TickSubscriber { readonly WebSocket _ws; public WsSubscriber(WebSocket ws) => _ws = ws; public override bool IsAlive => _ws.State == WebSocketState.Open; public override ValueTask SendAsync(ReadOnlyMemory payload) => _ws.SendAsync(payload, WebSocketMessageType.Text, true, CancellationToken.None); } sealed class SseSubscriber : TickSubscriber { readonly Stream _stream; static readonly byte[] Prefix = Encoding.UTF8.GetBytes("data: "); static readonly byte[] Suffix = Encoding.UTF8.GetBytes("\n\n"); public SseSubscriber(Stream stream) => _stream = stream; public override bool IsAlive => _stream.CanWrite; public override async ValueTask SendAsync(ReadOnlyMemory payload) { await _stream.WriteAsync(Prefix); await _stream.WriteAsync(payload); await _stream.WriteAsync(Suffix); await _stream.FlushAsync(); } } class MtConnectionService : BackgroundService { // Tick subscribers, keyed by symbol name. Each value is the set of connected // client sinks (WebSocket or SSE) listening to that symbol's tick stream. public static readonly ConcurrentDictionary> TickSubscribers = new(); public const int MaxSubscribersPerSymbol = 50; private readonly MtApi5Client _client; private readonly ILogger _logger; private readonly TimeSpan _connectTimeout; private const int Port = 8228; private static readonly TimeSpan HeartbeatInterval = TimeSpan.FromSeconds(30); private static readonly JsonSerializerSettings JsonSettings = new() { DateFormatString = "yyyy-MM-ddTHH:mm:ss", Formatting = Formatting.None }; public MtConnectionService(MtApi5Client client, ILogger logger) { _client = client; _logger = logger; _connectTimeout = TimeSpan.FromSeconds(15); _client.QuoteUpdate += OnQuoteUpdate; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { if (_client.ConnectionState != Mt5ConnectionState.Connected) { var firstConnect = _client.ConnectionState == Mt5ConnectionState.Disconnected; _logger.LogWarning(firstConnect ? "Connecting to MT5 via MtApi5 on port {Port}..." : "MT5 disconnected, attempting reconnect on port {Port}...", Port); var connected = await TryConnectAsync(stoppingToken); if (!connected) { _logger.LogError("MT5 is not connected. Ensure MtApi5 EA is loaded and listening on port {Port}.", Port); await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken); continue; } _logger.LogInformation(firstConnect ? "MT5 Bridge ready on http://localhost:8080" : "MT5 reconnected."); } await Task.Delay(TimeSpan.FromSeconds(30), stoppingToken); } } private async Task TryConnectAsync(CancellationToken stoppingToken) { if (_client.ConnectionState == Mt5ConnectionState.Connected) { return true; } var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); EventHandler? handler = null; handler = (_, e) => { _logger.LogInformation("MT5 Connection: {Status}", e.Status); if (e.Status == Mt5ConnectionState.Connected || e.Status == Mt5ConnectionState.Failed || e.Status == Mt5ConnectionState.Disconnected) { tcs.TrySetResult(e.Status); } }; _client.ConnectionStateChanged += handler; try { _client.BeginConnect(Port); using var timeoutCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); timeoutCts.CancelAfter(_connectTimeout); using var registration = timeoutCts.Token.Register(() => tcs.TrySetCanceled(timeoutCts.Token)); var status = await tcs.Task; return status == Mt5ConnectionState.Connected; } catch (OperationCanceledException) when (!stoppingToken.IsCancellationRequested) { _logger.LogWarning("Timed out waiting {TimeoutSeconds}s for MT5 connection state change.", _connectTimeout.TotalSeconds); return false; } finally { _client.ConnectionStateChanged -= handler; } } private void OnQuoteUpdate(object? sender, Mt5QuoteEventArgs e) { if (!MtConnectionService.TickSubscribers.TryGetValue(e.Quote.Instrument, out var subs)) return; var payload = JsonConvert.SerializeObject(new { type = "tick", symbol = e.Quote.Instrument, time = e.Quote.Time, bid = e.Quote.Bid, ask = e.Quote.Ask, last = e.Quote.Last, volume = e.Quote.Volume, time_msc = e.Quote.Time.ToString("yyyy-MM-ddTHH:mm:ss.fff000"), flags = 6 }, JsonSettings); var bytes = Encoding.UTF8.GetBytes(payload); // ponytail: fire-and-forget. Slow clients get dropped by SendAsync fault; cleanup happens on disconnect. foreach (var sub in subs.Keys) { if (sub.IsAlive) _ = sub.SendAsync(bytes); } } // ponytail: periodic heartbeat so clients behind NAT/firewall can detect liveness // even on low-liquidity symbols with no ticks for minutes. private static readonly byte[] HeartbeatPayload = Encoding.UTF8.GetBytes( JsonConvert.SerializeObject(new { type = "heartbeat", time = DateTime.UtcNow })); private static void BroadcastHeartbeat(object? _) { foreach (var subs in TickSubscribers.Values) { foreach (var sub in subs.Keys) { if (sub.IsAlive) _ = sub.SendAsync(HeartbeatPayload); } } } private readonly Timer _heartbeat = new(BroadcastHeartbeat, null, HeartbeatInterval, HeartbeatInterval); public override async Task StopAsync(CancellationToken cancellationToken) { await _heartbeat.DisposeAsync(); await base.StopAsync(cancellationToken); } }