Advertisement IC Markets

How to ask support for historical spread data and what they send

Discuss 1-minute to 15-minute price action setups, fading intraday momentum, key support/resistance zones, and proven short-term trading methodologies.
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

I guess that you allready know that FILE_SHARE_READ buffering—is an architectural anti-pattern. MetaTrader runs OnTick() synchronously. If the disk controller stalls, or if you are running 20 charts across 5 brokers on your HP Z1 workstation, your I/O thread will eventually bottleneck and miss ticks during high-volatility news events.

The true enterprise architecture decouples the MetaTrader client from the storage layer. We will turn MT4/MT5 into a pure, stateless telemetry sensor.

This version queues ticks in a memory buffer (an array of structs) and flushes them as a batched JSON payload via HTTP POST to a local ASP.NET Core microservice. Your C# backend can then handle the database ingestion (e.g., to MS SQL Server via Entity Framework Core) asynchronously.
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Recommended broker for automated trading & scalping IC Markets
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

1. The Universal MT4/MT5 Telemetry Sensor

This EA builds a raw JSON array in memory to avoid external MQL JSON library dependencies, then uses MetaTrader's native WebRequest to fire the batch to your C# API.

Note: You must add http://localhost:5000 to the allowed WebRequest URLs in MetaTrader (Tools -> Options -> Expert Advisors).

Code: Select all

//+------------------------------------------------------------------+
//|                                     EnterpriseSpreadSensor.mq4/5 |
//|                     Stateless JSON Telemetry over HTTP POST      |
//+------------------------------------------------------------------+
#property copyright "Copyright 2026"
#property strict

//--- Inputs
input string   InpApiEndpoint    = "http://localhost:5000/api/ticks"; // C# Microservice URL
input int      InpBatchSize      = 50;                                // Ticks per JSON payload
input string   InpSessionStart   = "07:00:00";                        // UTC Start
input string   InpSessionEnd     = "10:00:00";                        // UTC End

//--- Tick Buffer Struct
struct TickRecord {
   ulong    time_msc;
   double   bid;
   double   ask;
   double   spread_pips;
};

TickRecord     g_tickQueue[];
int            g_queueCount = 0;
double         g_lastBid = 0.0;
double         g_lastAsk = 0.0;

//+------------------------------------------------------------------+
int OnInit()
{
   ArrayResize(g_tickQueue, InpBatchSize);
   Print("[Telemetry Sensor] Initialized. Target: ", InpApiEndpoint);
   return(INIT_SUCCEEDED);
}

//+------------------------------------------------------------------+
void OnTick()
{
   MqlTick tick;
   if(!SymbolInfoTick(_Symbol, tick)) return;

   // Filter redundant quotes
   if(tick.bid == g_lastBid && tick.ask == g_lastAsk) return;

   datetime gmtNow = TimeGMT();
   if(!IsWithinSession(gmtNow, InpSessionStart, InpSessionEnd)) return;

   double pipSize = GetPipUnit();
   double spreadPips = (pipSize > 0) ? (tick.ask - tick.bid) / pipSize : 0.0;

   // Enqueue Tick
   g_tickQueue[g_queueCount].time_msc = tick.time_msc;
   g_tickQueue[g_queueCount].bid = tick.bid;
   g_tickQueue[g_queueCount].ask = tick.ask;
   g_tickQueue[g_queueCount].spread_pips = spreadPips;
   g_queueCount++;

   g_lastBid = tick.bid;
   g_lastAsk = tick.ask;

   // Flush Batch over HTTP
   if(g_queueCount >= InpBatchSize)
   {
      FlushTelemetry();
      g_queueCount = 0;
   }
}

//+------------------------------------------------------------------+
//| HTTP POST & JSON Serialization                                   |
//+------------------------------------------------------------------+
void FlushTelemetry()
{
   // Manually build JSON array to avoid external DLLs/libraries in MT
   string json = "[";
   for(int i = 0; i < g_queueCount; i++)
   {
      string record = StringFormat("{\"symbol\":\"%s\",\"time_msc\":%I64u,\"bid\":%.5f,\"ask\":%.5f,\"spread\":%.2f}",
                                   _Symbol, g_tickQueue[i].time_msc, 
                                   g_tickQueue[i].bid, g_tickQueue[i].ask, g_tickQueue[i].spread_pips);
      json += record;
      if(i < g_queueCount - 1) json += ",";
   }
   json += "]";

   char postData[];
   StringToCharArray(json, postData, 0, WHOLE_ARRAY, CP_UTF8);
   
   // Strip trailing null terminator injected by StringToCharArray
   int dataSize = ArraySize(postData) - 1; 

   char result[];
   string resultHeaders;
   string headers = "Content-Type: application/json\r\n";

   int res = WebRequest("POST", InpApiEndpoint, headers, 1000, postData, dataSize, result, resultHeaders);
   
   if(res != 200)
   {
      Print(StringFormat("[Telemetry Error] HTTP %d. Is the C# backend running?", res));
   }
}

//+------------------------------------------------------------------+
//| Utilities                                                        |
//+------------------------------------------------------------------+
double GetPipUnit() { return (_Digits == 3 || _Digits == 5) ? _Point * 10.0 : _Point; }

bool IsWithinSession(datetime time, string startStr, string endStr)
{
   datetime sod = StringToTime(TimeToString(time, TIME_DATE));
   int currentSec = (int)(time - sod);
   
   string sParts[], eParts[];
   StringSplit(startStr, ':', sParts);
   StringSplit(endStr, ':', eParts);
   
   int startSec = (int)StringToInteger(sParts[0])*3600 + (int)StringToInteger(sParts[1])*60 + (int)StringToInteger(sParts[2]);
   int endSec = (int)StringToInteger(eParts[0])*3600 + (int)StringToInteger(eParts[1])*60 + (int)StringToInteger(eParts[2]);

   return (startSec <= endSec) ? (currentSec >= startSec && currentSec <= endSec) 
                               : (currentSec >= startSec || currentSec <= endSec);
}
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

2. The ASP.NET Core "Catcher" (High-Performance Ingestion)

To handle the incoming HTTP POSTs without blocking, we use an ASP.NET Core Minimal API coupled with System.Threading.Channels. This creates an in-memory queue on the C# side, allowing the API to instantly return HTTP 200 to MetaTrader while a background service bulk-inserts the data into MS SQL Server or a master CSV.

Create a new .NET 8 Web API project and replace Program.cs:

Code: Select all

using System.Threading.Channels;

var builder = WebApplication.CreateBuilder(args);

// Create an unbounded channel for high-throughput, non-blocking ingestion
var tickChannel = Channel.CreateUnbounded<TickRecord>();
builder.Services.AddSingleton(tickChannel.Writer);
builder.Services.AddSingleton(tickChannel.Reader);
builder.Services.AddHostedService<TickIngestionWorker>();

var app = builder.Build();

// Minimal API Endpoint - O(1) execution time, instantly frees the MT4/MT5 thread
app.MapPost("/api/ticks", async (TickRecord[] payload, ChannelWriter<TickRecord> writer) =>
{
    foreach (var tick in payload)
    {
        await writer.WriteAsync(tick);
    }
    return Results.Ok();
});

app.Run("http://localhost:5000");

// --- Models ---
public record TickRecord(string Symbol, ulong Time_msc, double Bid, double Ask, double Spread);

// --- Background Worker (Decoupled Database Writer) ---
public class TickIngestionWorker : BackgroundService
{
    private readonly ChannelReader<TickRecord> _reader;
    private readonly ILogger<TickIngestionWorker> _logger;

    public TickIngestionWorker(ChannelReader<TickRecord> reader, ILogger<TickIngestionWorker> logger)
    {
        _reader = reader;
        _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("Tick Ingestion Worker started.");

        // Read dynamically as ticks arrive from the channel
        await foreach (var tick in _reader.ReadAllAsync(stoppingToken))
        {
            // TODO: Map to Entity Framework Core -> MS SQL Server
            // or batch write to a daily partitioned CSV file.
            
            _logger.LogInformation("[{Symbol}] Spread: {Spread} | Bid: {Bid}", tick.Symbol, tick.Spread, tick.Bid);
        }
    }
}
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

Here is the pure MQL4 implementation. Because MT4 Expert Advisors are strictly single-threaded, WebRequest blocks the entire EA execution until the HTTP response is received. This makes the decoupled architecture even more critical for MT4 than MT5—your C# backend must ingest the payload and return an HTTP status immediately.

1. The MT4 Sensor (SpreadTelemetry.mq4)
This script is compiled specifically for the MetaEditor 4 compiler. It builds the JSON payload natively to avoid relying on external MT4 DLLs that often break across terminal updates.

Code: Select all

//+------------------------------------------------------------------+
//|                                              SpreadTelemetry.mq4 |
//|                     Stateless JSON Telemetry over HTTP POST      |
//+------------------------------------------------------------------+
#property copyright "Enterprise Telemetry"
#property strict

//--- Inputs
input string   InpApiEndpoint    = "http://localhost:5000/api/ticks"; 
input int      InpBatchSize      = 50;                                
input string   InpSessionStart   = "07:00:00";                        
input string   InpSessionEnd     = "10:00:00";                        

//--- Structs & Globals
struct TickRecord {
   long     time_msc;
   double   bid;
   double   ask;
   double   spread_pips;
};

TickRecord     g_tickQueue[];
int            g_queueCount = 0;
double         g_lastBid = 0.0;
double         g_lastAsk = 0.0;

//+------------------------------------------------------------------+
int OnInit()
{
   ArrayResize(g_tickQueue, InpBatchSize);
   Print("[MT4 Telemetry] Initialized. Target: ", InpApiEndpoint);
   return(INIT_SUCCEEDED);
}

//+------------------------------------------------------------------+
void OnTick()
{
   MqlTick tick;
   if(!SymbolInfoTick(_Symbol, tick)) return;

   if(tick.bid == g_lastBid && tick.ask == g_lastAsk) return;

   datetime gmtNow = TimeGMT();
   if(!IsWithinSession(gmtNow, InpSessionStart, InpSessionEnd)) return;

   double pipSize = (_Digits == 3 || _Digits == 5) ? _Point * 10.0 : _Point;
   double spreadPips = (pipSize > 0) ? (tick.ask - tick.bid) / pipSize : 0.0;

   // Enqueue
   g_tickQueue[g_queueCount].time_msc = tick.time_msc;
   g_tickQueue[g_queueCount].bid = tick.bid;
   g_tickQueue[g_queueCount].ask = tick.ask;
   g_tickQueue[g_queueCount].spread_pips = spreadPips;
   g_queueCount++;

   g_lastBid = tick.bid;
   g_lastAsk = tick.ask;

   if(g_queueCount >= InpBatchSize)
   {
      FlushTelemetry();
      g_queueCount = 0;
   }
}

//+------------------------------------------------------------------+
void FlushTelemetry()
{
   string json = "[";
   for(int i = 0; i < g_queueCount; i++)
   {
      string record = StringFormat("{\"symbol\":\"%s\",\"time_msc\":%I64d,\"bid\":%.5f,\"ask\":%.5f,\"spread\":%.2f}",
                                   _Symbol, g_tickQueue[i].time_msc, 
                                   g_tickQueue[i].bid, g_tickQueue[i].ask, g_tickQueue[i].spread_pips);
      json += record;
      if(i < g_queueCount - 1) json += ",";
   }
   json += "]";

   char postData[];
   StringToCharArray(json, postData, 0, WHOLE_ARRAY, CP_UTF8);
   
   // MT4 StringToCharArray appends a null terminator, strip it for clean JSON
   int dataSize = ArraySize(postData) - 1; 

   char result[];
   string resultHeaders;
   string headers = "Content-Type: application/json\r\n";

   // Timeout set to 500ms to prevent MT4 thread locking
   int res = WebRequest("POST", InpApiEndpoint, headers, 500, postData, dataSize, result, resultHeaders);
   
   if(res != 200)
   {
      if(res == 4014) 
         Print("[MT4 Telemetry] ERR 4014: Add ", InpApiEndpoint, " to allowed WebRequest URLs in Options -> Expert Advisors.");
      else 
         Print("[MT4 Telemetry] HTTP Error: ", res);
   }
}

//+------------------------------------------------------------------+
bool IsWithinSession(datetime time, string startStr, string endStr)
{
   datetime sod = StringToTime(TimeToString(time, TIME_DATE));
   int currentSec = (int)(time - sod);
   
   string sParts[], eParts[];
   StringSplit(startStr, ':', sParts);
   StringSplit(endStr, ':', eParts);
   
   int startSec = (int)StringToInteger(sParts[0])*3600 + (int)StringToInteger(sParts[1])*60 + (int)StringToInteger(sParts[2]);
   int endSec = (int)StringToInteger(eParts[0])*3600 + (int)StringToInteger(eParts[1])*60 + (int)StringToInteger(eParts[2]);

   return (startSec <= endSec) ? (currentSec >= startSec && currentSec <= endSec) 
                               : (currentSec >= startSec || currentSec <= endSec);
}
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

2. The ASP.NET Core Ingestion Service (EF Core & MS SQL)

Because MT4 will freeze if the HTTP response is delayed, the C# API acts solely as a high-speed router. It accepts the JSON, drops it into an unbounded memory channel, and instantly returns 200 OK. A background IHostedService asynchronously flushes the channel into an MS SQL Server database via Entity Framework Core.

Code: Select all

using System.Threading.Channels;
using Microsoft.EntityFrameworkCore;

var builder = WebApplication.CreateBuilder(args);

// Configure Entity Framework Core for MS SQL Server
builder.Services.AddDbContext<TelemetryDbContext>(options =>
    options.UseSqlServer("Server=localhost;Database=MarketTelemetry;Trusted_Connection=True;TrustServerCertificate=True;"));

// Create high-throughput channel
var tickChannel = Channel.CreateUnbounded<TickRecord>();
builder.Services.AddSingleton(tickChannel.Writer);
builder.Services.AddSingleton(tickChannel.Reader);

// Register Background Worker
builder.Services.AddHostedService<SqlIngestionWorker>();

var app = builder.Build();

// O(1) Minimal API - Unblocks the MT4 thread instantly
app.MapPost("/api/ticks", async (TickRecord[] payload, ChannelWriter<TickRecord> writer) =>
{
    foreach (var tick in payload)
    {
        await writer.WriteAsync(tick);
    }
    return Results.Ok();
});

app.Run("http://localhost:5000");

// --- Models ---
public record TickRecord(string Symbol, long Time_msc, double Bid, double Ask, double Spread);

public class TelemetryEntity
{
    public int Id { get; set; }
    public string Symbol { get; set; } = string.Empty;
    public DateTime TimestampUtc { get; set; }
    public double Bid { get; set; }
    public double Ask { get; set; }
    public double SpreadPips { get; set; }
}

public class TelemetryDbContext : DbContext
{
    public TelemetryDbContext(DbContextOptions<TelemetryDbContext> options) : base(options) { }
    public DbSet<TelemetryEntity> Ticks { get; set; }
}

// --- Background Worker ---
public class SqlIngestionWorker : BackgroundService
{
    private readonly ChannelReader<TickRecord> _reader;
    private readonly IServiceProvider _serviceProvider;
    private readonly ILogger<SqlIngestionWorker> _logger;

    public SqlIngestionWorker(ChannelReader<TickRecord> reader, IServiceProvider serviceProvider, ILogger<SqlIngestionWorker> logger)
    {
        _reader = reader;
        _serviceProvider = serviceProvider;
        _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("SQL Ingestion Worker Active.");

        while (!stoppingToken.IsCancellationRequested)
        {
            var batch = new List<TelemetryEntity>();

            // Read up to 500 records at a time from the channel
            while (batch.Count < 500 && _reader.TryRead(out var tick))
            {
                batch.Add(new TelemetryEntity
                {
                    Symbol = tick.Symbol,
                    TimestampUtc = DateTimeOffset.FromUnixTimeMilliseconds(tick.Time_msc).UtcDateTime,
                    Bid = tick.Bid,
                    Ask = tick.Ask,
                    SpreadPips = tick.Spread
                });
            }

            if (batch.Any())
            {
                using var scope = _serviceProvider.CreateScope();
                var dbContext = scope.ServiceProvider.GetRequiredService<TelemetryDbContext>();
                
                await dbContext.Ticks.AddRangeAsync(batch, stoppingToken);
                await dbContext.SaveChangesAsync(stoppingToken);
                
                _logger.LogInformation("Inserted {Count} ticks into MS SQL Database.", batch.Count);
            }
            else
            {
                // Wait for data if channel is empty
                await _reader.WaitToReadAsync(stoppingToken); 
            }
        }
    }
}
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

Here is the pure MT5 implementation. MetaTrader 5 offers a much more robust execution environment than MT4, but EAs still run synchronously on a single thread per chart. To prevent the EA from missing ticks during NFP or high-liquidity sweeps, the WebRequest timeout is kept ultra-short, and the batching architecture remains strictly decoupled.

I have updated the headers and optimized the MQL5 syntax specifically for MT5’s memory management.
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

1. The MT5 Sensor (SpreadTelemetry.mq5)

Because this avoids external DLLs or third-party JSON libraries, you can compile this directly in MetaEditor 5 and drop it onto any chart.

Code: Select all

//+------------------------------------------------------------------+
//|                                              SpreadTelemetry.mq5 |
//|                               Copyright 2026, AI Profi Solutions |
//|                                     https://aiprofisolutions.com |
//+------------------------------------------------------------------+
#property strict

//--- Inputs
input string   InpApiEndpoint    = "http://localhost:5000/api/ticks"; 
input int      InpBatchSize      = 50;                                
input string   InpSessionStart   = "07:00:00";                        
input string   InpSessionEnd     = "10:00:00";                        

//--- Structs & Globals
struct TickRecord {
   ulong    time_msc;
   double   bid;
   double   ask;
   double   spread_pips;
};

TickRecord     g_tickQueue[];
int            g_queueCount = 0;
double         g_lastBid = 0.0;
double         g_lastAsk = 0.0;

//+------------------------------------------------------------------+
int OnInit()
{
   ArrayResize(g_tickQueue, InpBatchSize);
   Print("[MT5 Telemetry] Initialized. Target: ", InpApiEndpoint);
   return(INIT_SUCCEEDED);
}

//+------------------------------------------------------------------+
void OnTick()
{
   MqlTick tick;
   if(!SymbolInfoTick(_Symbol, tick)) return;

   // Ignore ticks where bid/ask haven't changed (e.g., volume-only updates)
   if(tick.bid == g_lastBid && tick.ask == g_lastAsk) return;

   datetime gmtNow = TimeTradeServer();
   if(!IsWithinSession(gmtNow, InpSessionStart, InpSessionEnd)) return;

   double pipSize = (_Digits == 3 || _Digits == 5) ? _Point * 10.0 : _Point;
   double spreadPips = (pipSize > 0) ? (tick.ask - tick.bid) / pipSize : 0.0;

   // Enqueue
   g_tickQueue[g_queueCount].time_msc = tick.time_msc;
   g_tickQueue[g_queueCount].bid = tick.bid;
   g_tickQueue[g_queueCount].ask = tick.ask;
   g_tickQueue[g_queueCount].spread_pips = spreadPips;
   g_queueCount++;

   g_lastBid = tick.bid;
   g_lastAsk = tick.ask;

   if(g_queueCount >= InpBatchSize)
   {
      FlushTelemetry();
      g_queueCount = 0;
   }
}

//+------------------------------------------------------------------+
void FlushTelemetry()
{
   // Manual JSON array construction for zero-dependency compilation
   string json = "[";
   for(int i = 0; i < g_queueCount; i++)
   {
      string record = StringFormat("{\"symbol\":\"%s\",\"time_msc\":%I64u,\"bid\":%.5f,\"ask\":%.5f,\"spread\":%.2f}",
                                   _Symbol, g_tickQueue[i].time_msc, 
                                   g_tickQueue[i].bid, g_tickQueue[i].ask, g_tickQueue[i].spread_pips);
      json += record;
      if(i < g_queueCount - 1) json += ",";
   }
   json += "]";

   char postData[];
   StringToCharArray(json, postData, 0, WHOLE_ARRAY, CP_UTF8);
   
   // StringToCharArray appends a null terminator, strip it for clean JSON payload
   int dataSize = ArraySize(postData) - 1; 

   char result[];
   string resultHeaders;
   string headers = "Content-Type: application/json\r\n";

   // Timeout set to 500ms. MT5 will not block longer than this.
   int res = WebRequest("POST", InpApiEndpoint, headers, 500, postData, dataSize, result, resultHeaders);
   
   if(res != 200)
   {
      if(res == 4014) 
         Print("[MT5 Telemetry] ERR 4014: Add ", InpApiEndpoint, " to allowed WebRequest URLs in Tools -> Options -> Expert Advisors.");
      else 
         Print("[MT5 Telemetry] HTTP Error: ", res);
   }
}

//+------------------------------------------------------------------+
bool IsWithinSession(datetime time, string startStr, string endStr)
{
   datetime sod = time - (time % 86400); // Fast MT5 Start of Day calculation
   int currentSec = (int)(time - sod);
   
   string sParts[], eParts[];
   StringSplit(startStr, ':', sParts);
   StringSplit(endStr, ':', eParts);
   
   int startSec = (int)StringToInteger(sParts[0])*3600 + (int)StringToInteger(sParts[1])*60 + (int)StringToInteger(sParts[2]);
   int endSec = (int)StringToInteger(eParts[0])*3600 + (int)StringToInteger(eParts[1])*60 + (int)StringToInteger(eParts[2]);

   return (startSec <= endSec) ? (currentSec >= startSec && currentSec <= endSec) 
                               : (currentSec >= startSec || currentSec <= endSec);
}
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
PTScalper
Site Admin
Posts: 3027
Joined: Mon Jul 20, 2026 1:28 pm

Re: How to ask support for historical spread data and what they send

Post by PTScalper »

2. The ASP.NET Core Ingestion Service

For a multi-core workstation like the HP Z1, this C# backend fully utilizes hardware concurrency. The Minimal API instantly accepts the JSON payload and pushes it into an unbounded memory channel (returning 200 OK to MT5 in less than 2 milliseconds).

A background worker service then pulls from that channel and performs bulk inserts into your local MS SQL Server instance via Entity Framework Core.

Code: Select all

using System.Threading.Channels;
using Microsoft.EntityFrameworkCore;

var builder = WebApplication.CreateBuilder(args);

// Configure Entity Framework Core for MS SQL Server
builder.Services.AddDbContext<TelemetryDbContext>(options =>
    options.UseSqlServer("Server=localhost;Database=MarketTelemetry;Trusted_Connection=True;TrustServerCertificate=True;"));

// Create high-throughput channel for thread decoupling
var tickChannel = Channel.CreateUnbounded<TickRecord>();
builder.Services.AddSingleton(tickChannel.Writer);
builder.Services.AddSingleton(tickChannel.Reader);

// Register Background Worker
builder.Services.AddHostedService<SqlIngestionWorker>();

var app = builder.Build();

// O(1) Minimal API - Unblocks the MT5 thread instantly
app.MapPost("/api/ticks", async (TickRecord[] payload, ChannelWriter<TickRecord> writer) =>
{
    foreach (var tick in payload)
    {
        await writer.WriteAsync(tick);
    }
    return Results.Ok();
});

app.Run("http://localhost:5000");

// --- Models ---
public record TickRecord(string Symbol, ulong Time_msc, double Bid, double Ask, double Spread);

public class TelemetryEntity
{
    public int Id { get; set; }
    public string Symbol { get; set; } = string.Empty;
    public DateTime TimestampUtc { get; set; }
    public double Bid { get; set; }
    public double Ask { get; set; }
    public double SpreadPips { get; set; }
}

public class TelemetryDbContext : DbContext
{
    public TelemetryDbContext(DbContextOptions<TelemetryDbContext> options) : base(options) { }
    public DbSet<TelemetryEntity> Ticks { get; set; }
}

// --- Background Worker ---
public class SqlIngestionWorker : BackgroundService
{
    private readonly ChannelReader<TickRecord> _reader;
    private readonly IServiceProvider _serviceProvider;
    private readonly ILogger<SqlIngestionWorker> _logger;

    public SqlIngestionWorker(ChannelReader<TickRecord> reader, IServiceProvider serviceProvider, ILogger<SqlIngestionWorker> logger)
    {
        _reader = reader;
        _serviceProvider = serviceProvider;
        _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("SQL Ingestion Worker Active.");

        while (!stoppingToken.IsCancellationRequested)
        {
            var batch = new List<TelemetryEntity>();

            // Read up to 1000 records at a time from the channel
            while (batch.Count < 1000 && _reader.TryRead(out var tick))
            {
                batch.Add(new TelemetryEntity
                {
                    Symbol = tick.Symbol,
                    TimestampUtc = DateTimeOffset.FromUnixTimeMilliseconds((long)tick.Time_msc).UtcDateTime,
                    Bid = tick.Bid,
                    Ask = tick.Ask,
                    SpreadPips = tick.Spread
                });
            }

            if (batch.Any())
            {
                using var scope = _serviceProvider.CreateScope();
                var dbContext = scope.ServiceProvider.GetRequiredService<TelemetryDbContext>();
                
                // Bulk insert into MS SQL Server
                await dbContext.Ticks.AddRangeAsync(batch, stoppingToken);
                await dbContext.SaveChangesAsync(stoppingToken);
                
                _logger.LogInformation("Inserted {Count} ticks into MS SQL Database.", batch.Count);
            }
            else
            {
                // Yield thread if channel is empty
                await _reader.WaitToReadAsync(stoppingToken); 
            }
        }
    }
}
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Post Reply