2879 строки
Без EOL
137 КиБ
MQL5
2879 строки
Без EOL
137 КиБ
MQL5
//+------------------------------------------------------------------+
|
|
//| AppController.mqh |
|
|
//| Copyright 2026, MetaQuotes Ltd. |
|
|
//| www.mql5.com |
|
|
//+------------------------------------------------------------------+
|
|
#ifndef REQUEST_LATENCY_LAB_APP_CONTROLLER_MQH
|
|
#define REQUEST_LATENCY_LAB_APP_CONTROLLER_MQH
|
|
|
|
#include "..\..\Include\RequestLatencyLab\Models.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\Configuration.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\TimeSource.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\EventJournal.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\RequestTracker.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\ExperimentRunner.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\Legacy\TradeRequestTools.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\CompletionPolicy.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\StateReader.mqh"
|
|
#include "..\..\Include\RequestLatencyLab\MarketRegime.mqh"
|
|
|
|
enum ENUM_LAB_STATE
|
|
{
|
|
LAB_STATE_IDLE = 0,
|
|
LAB_STATE_PREFLIGHT,
|
|
LAB_STATE_READY,
|
|
LAB_STATE_WARMUP,
|
|
LAB_STATE_MAIN,
|
|
LAB_STATE_DRAINING,
|
|
LAB_STATE_EXPORTING,
|
|
LAB_STATE_FINISHED,
|
|
LAB_STATE_PAUSED,
|
|
LAB_STATE_BLOCKED
|
|
};
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
class CAppController
|
|
{
|
|
private:
|
|
LabSettings m_settings;
|
|
SymbolRules m_rules;
|
|
CRequestTracker *m_tracker;
|
|
CEventJournal *m_journal;
|
|
CExperimentRunner*m_runner;
|
|
CTimeSource *m_clock;
|
|
ENUM_LAB_STATE m_state;
|
|
string m_session_id;
|
|
ulong m_active_sequence; // текущая незавершённая операция
|
|
ulong m_last_pause_us; // момент завершения паузы
|
|
ulong m_pause_until; // следующая разрешённая отправка
|
|
bool m_operator_started;
|
|
string m_block_reason;
|
|
ulong m_n_cleanup_attempts;
|
|
ulong m_last_cleanup_sequence;
|
|
//--- B1 (audit-3): cleanup проходит отдельный lifecycle как запись CLEANUP
|
|
bool m_cleanup_in_flight; // отправка cleanup выполняется
|
|
ulong m_cleanup_sequence; // sequence записи CLEANUP
|
|
ulong m_cleanup_parent; // parent main sequence
|
|
int m_cleanup_attempt;
|
|
//--- E3: калибровка рынка и последний классифицированный режим
|
|
Calibration m_calibration;
|
|
bool m_calibration_loaded;
|
|
ENUM_LAB_MARKET_REGIME m_last_regime;
|
|
string m_last_window_id;
|
|
ulong m_last_regime_us;
|
|
ulong m_e3_start_us; // E3: начало ожидания режима
|
|
ulong m_e3_window_seq; // R4-B5: счётчик окон E3
|
|
string m_termination_reason; // формальный статус завершения
|
|
//--- REQUEST-ID служебных cleanup-отправок (не unresolved для основной записи)
|
|
ulong m_cleanup_request_ids[256];
|
|
ulong m_cleanup_request_us[256];
|
|
int m_cleanup_req_count;
|
|
MarketWindow m_windows[];
|
|
int m_window_count;
|
|
//--- R6-B6: сохранность контрольных точек и устойчивого журнала намерений
|
|
bool m_checkpoint_ok;
|
|
bool m_durable_ok;
|
|
ulong m_durable_uncertain;
|
|
//--- R6-S2: монотонный возраст последней наблюдённой котировки
|
|
ulong m_last_quote_observed_us;
|
|
//--- R7-B5: идентичность последнего НОВОГО обновления рабочего символа
|
|
//--- (серверная метка тика): повторное чтение прежней котировки не
|
|
//--- обновляет m_last_quote_observed_us.
|
|
long m_last_quote_msc;
|
|
//--- R8-S3: буфер SEND_RETURN durable — дописывание только в safe-фазе
|
|
//--- (после закрытия окна), синхронной записи между T1 и callbacks нет.
|
|
ulong m_pending_sr_seq[];
|
|
ulong m_pending_sr_us[];
|
|
int m_pending_sr_count;
|
|
//--- R8-S2: номер последнего классифицированного окна E3 (для samples)
|
|
ulong m_last_window_seq;
|
|
//--- R9-S2: якорь (time_msc) последнего классифицированного окна E3
|
|
long m_last_anchor_msc;
|
|
//--- R9-S3: сырые тики последнего окна E3 (в памяти до safe-фазы)
|
|
MqlTick m_window_ticks[];
|
|
ulong m_window_ticks_seq;
|
|
bool m_window_ticks_pending;
|
|
//--- R13-B1: инъекция отказа durable-записи SyntheticDispatch ТОЛЬКО
|
|
//--- для регрессий (счётчик: до какого AppendDurableIntent отказать;
|
|
//--- 0 = выключено; 1 = первый INTENT, 2 = первый SEND_RETURN, ...).
|
|
int m_test_fail_intent_counter;
|
|
//--- R14-B1: монотонный счётчик выданных session_id процесса EA
|
|
//--- (не сбрасывается при ResetForRestart/StartRun).
|
|
int m_session_counter;
|
|
|
|
public:
|
|
CAppController(void);
|
|
void Configure(CRequestTracker *tracker, CEventJournal *journal,
|
|
CExperimentRunner *runner, CTimeSource *clock);
|
|
//--- отправка конкретного плана (единая точка)
|
|
bool DispatchPlan(const RequestPlan &plan, const bool do_send, ulong &sent_sequence,
|
|
LabError &error)
|
|
{
|
|
return(PrepareAndDispatch(plan, do_send, sent_sequence, error));
|
|
}
|
|
static bool AppendDurableLine(const string path, const string line, LabError &err);
|
|
void QueueDurableSendReturn(const ulong seq, const ulong t_us);
|
|
int PendingDurableSendReturnCount() const { return(m_pending_sr_count); }
|
|
bool FlushPendingSendReturns(LabError &err);
|
|
static bool MeasurementPointsComplete(const RequestMetadata &md);
|
|
static bool CleanupCollectionReady(const RequestMetadata &md);
|
|
bool ConfirmDurableSendReturns(void);
|
|
static string MarketTickHeader(void);
|
|
static string FormatMarketTickLine(const string session_id, const string wid,
|
|
const ulong window_seq, const MqlTick &t);
|
|
static bool TickFileEmpty(const string path);
|
|
bool FlushWindowTicks(LabError &err);
|
|
bool EnsureMarketTicksFile(void);
|
|
void ReturnToSchedule(void);
|
|
bool E3ContextStale(const ulong now_us);
|
|
bool DriveNext(LabError &error);
|
|
void RequirePause(void);
|
|
bool PauseElapsed() const;
|
|
bool InitRun(LabError &error);
|
|
bool EndMainCycle(LabError &error);
|
|
bool DriveCleanup(LabError &error);
|
|
static bool ZeroResidualExposure(const ENUM_LAB_OPERATION op, const string symbol,
|
|
const ulong magic);
|
|
CRequestTracker *Tracker();
|
|
CExperimentRunner *Runner();
|
|
SymbolRules Rules();
|
|
bool CheckIntegrity(LabError &error);
|
|
ENUM_LAB_STATE State() const;
|
|
string SessionId() const;
|
|
string BlockReason() const;
|
|
ulong ActiveSequence() const;
|
|
int MarketWindowCount() const;
|
|
bool GetMarketWindow(const int index, MarketWindow &out) const;
|
|
string StateName() const;
|
|
bool Preflight(LabError &error);
|
|
bool StartRun(LabError &error);
|
|
void TerminateRun(const string reason);
|
|
string TerminationReason() const;
|
|
bool CleanupInFlight() const;
|
|
string ConfigurationId() const;
|
|
void MarkFinished();
|
|
void ResetSessionWindows(void);
|
|
bool ResetForRestart(void);
|
|
bool InMeasurementWindow() const;
|
|
void TestSetFailIntentAfter(const int n);
|
|
void TestSetCleanupWindow(const bool open);
|
|
void TestAddMarketWindow(const MarketWindow &w);
|
|
void StopRun(void);
|
|
bool IsRunning() const;
|
|
bool CanCheckpoint() const;
|
|
bool CheckpointOk() const;
|
|
void SetCheckpointState(const bool ok);
|
|
void BlockRun(const string reason);
|
|
void NoteQuoteObserved(const long tick_msc);
|
|
ulong LastQuoteObservedUs() const;
|
|
ulong CountUncertainDurableIntents(const string session_id);
|
|
bool LoadDurableContent(const string path, string &content);
|
|
void OnTransaction(TransactionObservation &obs, LabError &error);
|
|
bool OnTimer(const ulong timer_last_call_us, LabError &error);
|
|
bool CompletedCheck(const ulong record_sequence, ConfirmationSnapshot &snap, LabError &error);
|
|
bool SetCalibration(const Calibration &cal);
|
|
bool CalibrationLoaded() const;
|
|
string CalibrationId() const;
|
|
void UpdateMarketRegime(ENUM_LAB_MARKET_REGIME ®ime);
|
|
void FinalizeTransaction(TransactionObservation &obs, LabError &error);
|
|
bool RunSynthetic(LabError &error);
|
|
bool SyntheticDispatch(const RequestPlan &plan);
|
|
bool SessionDirUsed(const string dir);
|
|
string NewSessionId(void);
|
|
void SetSettings(const LabSettings &s);
|
|
LabSettings Settings();
|
|
|
|
private:
|
|
bool AppendDurableIntent(const ulong seq, const string kind, const ulong t_us);
|
|
ulong SendCleanup(const MqlTradeRequest &request, const ulong parent_sequence,
|
|
const ENUM_LAB_OPERATION op, LabError &error);
|
|
void RememberCleanupRequest(const uint rid);
|
|
bool IsCleanupRequest(const uint rid) const;
|
|
void CheckTriggered(const ulong seq, LabError &error);
|
|
void SyncRecoveredDeals(const ulong sequence, const ConfirmationSnapshot &snap);
|
|
bool PrepareAndDispatch(const RequestPlan &plan, const bool do_send,
|
|
ulong &sent_sequence, LabError &error);
|
|
};
|
|
//--- R6-B6: устойчивый журнал намерений (append-only, каталог сессии).
|
|
//--- INTENT пишется ДО T0; SEND_RETURN - после T1; пара различает
|
|
//--- «отправлено» и «только намерение» при восстановлении после сбоя.
|
|
bool CAppController::AppendDurableIntent(const ulong seq, const string kind, const ulong t_us)
|
|
{
|
|
if(m_session_id == NULL || m_session_id == "")
|
|
return(false);
|
|
const string dir = "RequestLatencyLab\\" + m_settings.campaign_id + "\\" + m_session_id;
|
|
if(!FolderCreate(dir))
|
|
return(false); // каталог сессии должен существовать
|
|
const string line = StringFormat("%s,%I64u,%s,%I64u\r\n",
|
|
m_session_id, seq, kind, t_us);
|
|
LabError derr;
|
|
return(AppendDurableLine(dir + "\\intents.csv", line, derr));
|
|
}
|
|
//--- подготовка плана к отправке (PreCheck, intent, T0/T1)
|
|
//--- B1 (audit-3): cleanup — полноценная запись tracker с role=CLEANUP.
|
|
//--- Отдельный lifecycle: precheck -> reserve -> T0/T1 -> correlation ->
|
|
//--- T6 -> подтверждение нулевой экспозиции -> пауза -> следующий MAIN.
|
|
//--- Не считает dispatch дважды и не продолжает серию до подтверждения.
|
|
ulong CAppController::SendCleanup(const MqlTradeRequest &request, const ulong parent_sequence,
|
|
const ENUM_LAB_OPERATION op, LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
//--- R9-B3: cleanup - тоже торговая отправка: при нарушении
|
|
//--- сохранности (durable/checkpoint) или терминальном состоянии
|
|
//--- новые торговые вызовы запрещены, даже служебные.
|
|
if(m_state == LAB_STATE_BLOCKED || m_state == LAB_STATE_FINISHED ||
|
|
!m_durable_ok || !m_checkpoint_ok)
|
|
{
|
|
error.code = 7;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "cleanup dispatch blocked (durability/safety gate)";
|
|
return(0);
|
|
}
|
|
//--- build & precheck по тем же торговым правилам (S6)
|
|
RequestPlan p;
|
|
p.Zero();
|
|
p.role = LAB_ROLE_CLEANUP;
|
|
p.operation = op;
|
|
p.parent_sequence = parent_sequence;
|
|
p.symbol = request.symbol;
|
|
p.magic = (ulong)request.magic;
|
|
p.side = (request.type == ORDER_TYPE_BUY ? LAB_SIDE_BUY :
|
|
(request.type == ORDER_TYPE_SELL ? LAB_SIDE_SELL : LAB_SIDE_NONE));
|
|
p.mode = LAB_MODE_ASYNC;
|
|
p.logging_mode = LAB_LOG_MINIMAL;
|
|
//--- R4-B1: в план cleanup сохраняется полный фактический контракт
|
|
//--- отправленного запроса (объём позиции/ордера, filling, deviation,
|
|
//--- цена/тикет цели, comment) — иначе T6 для POSITION_CLOSE не сможет
|
|
//--- сверить requested/executed объём в единых units.
|
|
p.volume = request.volume;
|
|
p.deviation_points = (int)request.deviation;
|
|
p.comment = request.comment;
|
|
p.actual_type = request.type;
|
|
p.actual_filling = request.type_filling;
|
|
p.actual_known = true; // R5-S2: 0 — допустимое значение enum
|
|
if(request.action == TRADE_ACTION_REMOVE && request.order != 0)
|
|
p.request_price = (double)request.order; // тикет отложенного ордера
|
|
else
|
|
p.request_price = request.price;
|
|
MqlTradeCheckResult chk;
|
|
if(!LabPreCheck(request, chk, error))
|
|
{
|
|
//--- отказ/неопределённость cleanup -> BLOCKED (безопасность)
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cleanup precheck rejected";
|
|
return(0);
|
|
}
|
|
//--- резервируем отдельную sequence CLEANUP
|
|
ulong seq = m_runner.NextSequence();
|
|
p.sequence = seq;
|
|
if(!m_tracker.Reserve(p, seq))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cannot reserve cleanup";
|
|
return(0);
|
|
}
|
|
LabEvent ev;
|
|
ev.Zero();
|
|
ev.kind = LAB_KIND_DISPATCH_INTENT;
|
|
ev.local_time_us = m_clock.NowUs();
|
|
ev.origin_sequence = seq;
|
|
ev.payload_json = StringFormat("{\"role\":\"CLEANUP\",\"parent\":%I64u}",
|
|
parent_sequence);
|
|
ulong es = 0;
|
|
if(!m_journal.Add(ev, es))
|
|
{
|
|
//--- R5-B2: intent не сохранён — отправка cleanup запрещена
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "journal overflow at cleanup intent";
|
|
return(0);
|
|
}
|
|
//--- R7-S3: cleanup-отправки — тоже фактические торговые вызовы: они
|
|
//--- фиксируются в устойчивом журнале намерений (INTENT до T0), иначе
|
|
//--- durable-журнал покрывал бы не все отправки.
|
|
if(!AppendDurableIntent(seq, "INTENT", m_clock.NowUs()))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "durable intent write failed at cleanup";
|
|
return(0);
|
|
}
|
|
//--- T0/T1
|
|
ResetLastError();
|
|
const ulong t0 = m_clock.NowUs();
|
|
m_tracker.MarkSendStart(seq, t0);
|
|
MqlTradeResult res;
|
|
const bool ok = OrderSendAsync(request, res);
|
|
const ulong t1 = m_clock.NowUs();
|
|
SendObservation send;
|
|
send.Zero();
|
|
send.t0_us = t0;
|
|
send.t1_us = t1;
|
|
send.send_ok = ok;
|
|
send.result = res;
|
|
send.last_error = GetLastError();
|
|
send.result_known = true;
|
|
m_tracker.MarkSendReturn(seq, send);
|
|
//--- R8-B3: SEND_RETURN cleanup — буфер; синхронная запись между T1
|
|
//--- и callbacks не выполняется (дописывание в safe-фазе после
|
|
//--- закрытия окна cleanup).
|
|
QueueDurableSendReturn(seq, t1);
|
|
LabEvent ev2;
|
|
ev2.Zero();
|
|
ev2.kind = LAB_KIND_SEND_RETURN;
|
|
ev2.local_time_us = t1;
|
|
ev2.origin_sequence = seq;
|
|
ev2.payload_json = StringFormat("{\"send_ok\":%s,\"retcode\":%u}",
|
|
(ok ? "true" : "false"), res.retcode);
|
|
if(!m_journal.Add(ev2, es))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "journal overflow at cleanup send_return";
|
|
return(0);
|
|
}
|
|
if(res.request_id != 0)
|
|
RememberCleanupRequest(res.request_id);
|
|
m_tracker.ComputeDeadlines(seq, (ulong)m_settings.outcome_timeout_ms * 1000,
|
|
(ulong)m_settings.late_grace_ms * 1000);
|
|
m_cleanup_in_flight = true;
|
|
m_cleanup_sequence = seq;
|
|
m_cleanup_parent = parent_sequence;
|
|
m_cleanup_attempt++;
|
|
m_n_cleanup_attempts++;
|
|
m_last_cleanup_sequence = seq;
|
|
m_runner.MarkCleanupDispatched(); // однократный учёт при отправке
|
|
return(seq);
|
|
}
|
|
//--- служебный реестр cleanup request_id (без ложных unresolved)
|
|
void CAppController::RememberCleanupRequest(const uint rid)
|
|
{
|
|
if(rid == 0 || m_cleanup_req_count >= 256)
|
|
return;
|
|
m_cleanup_request_ids[m_cleanup_req_count] = rid;
|
|
m_cleanup_request_us[m_cleanup_req_count] = m_clock.NowUs();
|
|
m_cleanup_req_count++;
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::IsCleanupRequest(const uint rid) const
|
|
{
|
|
if(rid == 0)
|
|
return(false);
|
|
for(int i = 0; i < m_cleanup_req_count; i++)
|
|
if((uint)(m_cleanup_request_ids[i] & 0xFFFFFFFF) == rid)
|
|
return(true);
|
|
return(false);
|
|
}
|
|
//--- адресная сверка конкретной записи по сигналу события (аудит-2 B4)
|
|
void CAppController::CheckTriggered(const ulong seq, LabError &error)
|
|
{
|
|
if(seq == 0)
|
|
return;
|
|
RequestMetadata r;
|
|
if(!m_tracker.GetMetadata(seq, r))
|
|
return;
|
|
const bool have_t6 = (r.present_mask & (uint)LAB_MASK_T6) != 0;
|
|
//--- R6-B1: наличие T6 НЕ прекращает перенос результатов сверки
|
|
//--- (callback-покрытие для T5); исход и T6 не переписываются.
|
|
//--- закрытая коллекция без T6 не переоткрывается; с T6 — только
|
|
//--- перенос покрытия сверки ниже.
|
|
if(r.collection_closed && !have_t6)
|
|
return;
|
|
ConfirmationSnapshot snap;
|
|
//--- B5 (audit-3): событийная сверка адресная и короткая — полный
|
|
//--- проход истории вынесен в таймер (light=true, volume из callback).
|
|
//--- R4-B2: callback_volume накапливается в лотах; на границе
|
|
//--- tracker/snapshot он переводится в units ровно один раз.
|
|
double cb_units = 0.0;
|
|
if(r.callback_volume > 0.0)
|
|
CStateReader::LotsToUnits(r.callback_volume, m_rules.volume_step,
|
|
m_settings.volume_units_tolerance, cb_units);
|
|
CStateReader::BuildSnapshot(seq, r.plan.operation, r.order_ticket,
|
|
r.plan.volume / m_rules.volume_step,
|
|
m_rules.volume_step, m_settings.volume_units_tolerance,
|
|
m_clock, r.position_ticket,
|
|
m_settings.history_scan_budget, snap,
|
|
true, cb_units, r.deal_count);
|
|
snap.requested_volume_units = r.plan.volume / m_rules.volume_step;
|
|
snap.order_state_after = snap.order_state_before;
|
|
if(have_t6)
|
|
{
|
|
m_tracker.SyncReconcile(seq, snap, m_rules.volume_step,
|
|
m_settings.volume_units_tolerance);
|
|
return;
|
|
}
|
|
snap.trigger = LAB_FINAL_CALLBACK_CHECK;
|
|
CompletedCheck(seq, snap, error);
|
|
}
|
|
//--- R5-B4: регистрация восстановленных из истории сделок в реестре
|
|
//--- (deals.csv) с происхождением HISTORY_RECOVERED.
|
|
void CAppController::SyncRecoveredDeals(const ulong sequence, const ConfirmationSnapshot &snap)
|
|
{
|
|
if(m_tracker == NULL || snap.deal_count <= 0)
|
|
return;
|
|
for(int k = 0; k < snap.deal_count; k++)
|
|
{
|
|
const ulong dt = snap.deal_tickets[k];
|
|
if(dt == 0)
|
|
continue;
|
|
const double vol = HistoryDealGetDouble(dt, DEAL_VOLUME);
|
|
const long tms = (long)HistoryDealGetInteger(dt, DEAL_TIME_MSC);
|
|
m_tracker.RegisterHistoryDeal(dt, snap.order_ticket, vol, tms, sequence);
|
|
}
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::PrepareAndDispatch(const RequestPlan &plan, const bool do_send,
|
|
ulong &sent_sequence, LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
//--- R8-B1: последняя защита перед транспортом — терминальное/
|
|
//--- блокирующее состояние (E3_WINDOW_CAPACITY_EXCEEDED,
|
|
//--- E3_WINDOW_STORAGE_FAILED и т.п.) запрещает НОВУЮ отправку, даже
|
|
//--- если DriveNext уже вошёл в ветвь MAIN до перехода состояния.
|
|
if(m_state == LAB_STATE_FINISHED || m_state == LAB_STATE_BLOCKED)
|
|
{
|
|
error.code = 7;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "controller not dispatchable in state " + StateName();
|
|
return(false);
|
|
}
|
|
//--- контекст рынка и запрос
|
|
MqlTick quote;
|
|
if(!SymbolInfoTick(m_settings.symbol, quote))
|
|
{
|
|
//--- временное отсутствие котировки => WAIT_DATA (B9: не precheck reject)
|
|
error.code = 5;
|
|
error.message = "WAIT_DATA";
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkWaitSlot(plan.sequence);
|
|
return(false);
|
|
}
|
|
RequestPlan p = plan;
|
|
p.volume = (p.volume > 0.0 ? p.volume : CConfiguration::MinVolume(m_rules));
|
|
//--- S1: фактическое отклонение и comment из конфигурации
|
|
p.deviation_points = m_settings.deviation_points;
|
|
if(StringLen(p.comment) == 0)
|
|
p.comment = StringFormat("RLL_%s_%I64u", plan.condition_id,
|
|
(plan.sequence != 0 ? plan.sequence : (ulong)m_clock.NowUs()));
|
|
MqlTradeRequest request;
|
|
if(!LabBuildRequest(p, m_rules, quote, request, error))
|
|
{
|
|
//--- R4-B7: сбой построения запроса — слот возвращается в очередь
|
|
//--- (ограниченное число повторов), квоту не расходует.
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkRejectedSlot(plan.sequence);
|
|
return(false);
|
|
}
|
|
//--- R4-B8: фактический контракт запроса экспортируется (samples.csv)
|
|
p.actual_type = request.type;
|
|
p.actual_filling = request.type_filling;
|
|
p.request_price = request.price;
|
|
p.actual_known = true; // R5-S2: 0 — допустимое значение enum
|
|
MqlTradeCheckResult check;
|
|
if(!LabPreCheck(request, check, error))
|
|
{
|
|
//--- B9: однократная фиксация precheck-rejected как отдельного статуса
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkRejectedSlot(plan.sequence);
|
|
m_runner.MarkPrecheckRejected();
|
|
return(false);
|
|
}
|
|
//--- intent сохраняется до отправки
|
|
ulong seq = plan.sequence;
|
|
if(seq == 0)
|
|
seq = m_runner.NextSequence();
|
|
if(!m_tracker.Reserve(p, seq))
|
|
{
|
|
//--- R4-B7: переполнение ёмкости — WAIT (повторная попытка позже)
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkWaitSlot(plan.sequence);
|
|
error.code = 6;
|
|
error.message = "cannot reserve sequence";
|
|
return(false);
|
|
}
|
|
p.sequence = seq;
|
|
//--- S2: фиксация рыночного контекста (bid/ask/spread/tick_time)
|
|
m_tracker.SetMarketContext(seq, quote.bid, quote.ask, quote.time_msc);
|
|
//--- R8-S2: E3 — перепроверка возраста рыночного контекста НЕПОСРЕД-
|
|
//--- СТВЕННО перед T0 (после классификации в DriveNext могла пройти
|
|
//--- пауза/события). Застаревший контекст не отправляется; слот
|
|
//--- возвращается в очередь (WAIT), durable-INTENT не пишется.
|
|
if(m_settings.experiment_id == LAB_EXP_E3 &&
|
|
m_settings.feature_max_age_ms > 0 && E3ContextStale(m_clock.NowUs()))
|
|
{
|
|
if(seq != 0)
|
|
m_tracker.ReleaseReservation(seq); // R9-B1: откат резерва
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkWaitSlot(plan.sequence);
|
|
error.code = 8;
|
|
error.message = "E3 quote context stale before T0";
|
|
return(false);
|
|
}
|
|
LabEvent intent;
|
|
intent.Zero();
|
|
intent.kind = LAB_KIND_DISPATCH_INTENT;
|
|
intent.local_time_us = m_clock.NowUs();
|
|
intent.origin_sequence = seq;
|
|
intent.payload_json = StringFormat("{\"role\":\"%s\",\"op\":\"%s\",\"slot\":%I64u}",
|
|
LabRoleName(p.role), LabOperationName(p.operation), p.slot_id);
|
|
ulong ev = 0;
|
|
if(!m_journal.Add(intent, ev))
|
|
{
|
|
//--- R5-B2: intent не сохранён — отправка запрещена (нельзя торговать
|
|
//--- без протокольной записи намерения).
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkWaitSlot(plan.sequence);
|
|
if(seq != 0)
|
|
m_tracker.ReleaseReservation(seq); // R9-B1: откат резерва
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "journal overflow at intent";
|
|
error.code = 61;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "intent not persisted";
|
|
return(false);
|
|
}
|
|
//--- R6-B6: устойчивый журнал намерений ДО T0 (crash-safe intent);
|
|
//--- локальный прогон (do_send=false) в durable-журнал не пишется.
|
|
if(do_send && !AppendDurableIntent(seq, "INTENT", m_clock.NowUs()))
|
|
{
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkWaitSlot(plan.sequence);
|
|
if(seq != 0)
|
|
m_tracker.ReleaseReservation(seq); // R9-B1: откат резерва
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "durable intent write failed";
|
|
error.code = 63;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "durable intent not persisted";
|
|
return(false);
|
|
}
|
|
if(!do_send)
|
|
{
|
|
//--- LOcal-only: транспорт-счётчик вызывается в тестах/драйвере
|
|
return(true);
|
|
}
|
|
//--- R9-S2: повторная проверка свежести ВЫБРАННОГО контекста
|
|
//--- НЕПОСРЕДСТВЕННО перед транспортным вызовом (между предыдущей
|
|
//--- проверкой и T0 выполнялись синхронные записи durable-INTENT,
|
|
//--- которые могли затянуться). При застое durable-INTENT явно
|
|
//--- отменяется (INTENT_CANCELLED), резерв откатывается (R9-B1),
|
|
//--- слот повторно ставится в очередь — без фантомного send.
|
|
if(m_settings.experiment_id == LAB_EXP_E3 &&
|
|
m_settings.feature_max_age_ms > 0 && E3ContextStale(m_clock.NowUs()))
|
|
{
|
|
if(!AppendDurableIntent(seq, "INTENT_CANCELLED", m_clock.NowUs()))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "durable intent cancel write failed";
|
|
error.code = 63;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "durable intent cancel not persisted";
|
|
if(seq != 0)
|
|
m_tracker.ReleaseReservation(seq);
|
|
return(false);
|
|
}
|
|
if(seq != 0)
|
|
m_tracker.ReleaseReservation(seq);
|
|
if(plan.sequence != 0)
|
|
m_runner.MarkWaitSlot(plan.sequence);
|
|
error.code = 8;
|
|
error.message = "E3 quote context stale at T0";
|
|
return(false);
|
|
}
|
|
//--- единственная точка торговой отправки
|
|
SendObservation send;
|
|
send.Zero();
|
|
send.result_known = false;
|
|
ResetLastError();
|
|
const ulong t0 = m_clock.NowUs();
|
|
m_tracker.MarkSendStart(seq, t0);
|
|
send.t0_us = t0;
|
|
if(p.mode == LAB_MODE_SYNC)
|
|
{
|
|
const bool ok = OrderSend(request, send.result);
|
|
send.send_ok = ok;
|
|
send.last_error = GetLastError();
|
|
//--- режим фиксируется в плане (SYNC/ASYNC)
|
|
}
|
|
else
|
|
{
|
|
const bool ok = OrderSendAsync(request, send.result);
|
|
send.send_ok = ok;
|
|
send.last_error = GetLastError();
|
|
//--- async: send_ok=false тоже регистрируется как вызов
|
|
}
|
|
const ulong t1 = m_clock.NowUs();
|
|
send.t1_us = t1;
|
|
send.result_known = true;
|
|
m_tracker.MarkSendReturn(seq, send);
|
|
//--- R5-S3: первичный журнал фиксирует фактические данные отправки
|
|
//--- (для MAIN так же, как для cleanup).
|
|
LabEvent esx;
|
|
esx.Zero();
|
|
esx.kind = LAB_KIND_SEND_START;
|
|
esx.local_time_us = t0;
|
|
esx.origin_sequence = seq;
|
|
esx.payload_json = StringFormat(
|
|
"{\"role\":\"%s\",\"action\":%d,\"symbol\":\"%s\",\"volume\":%.10g,\"type\":%d}",
|
|
LabRoleName(p.role), (int)request.action, request.symbol, request.volume,
|
|
(int)request.type);
|
|
ulong esr_id = 0;
|
|
m_journal.Add(esx, esr_id);
|
|
LabEvent esr;
|
|
esr.Zero();
|
|
esr.kind = LAB_KIND_SEND_RETURN;
|
|
esr.local_time_us = t1;
|
|
esr.origin_sequence = seq;
|
|
esr.request_id = send.result.request_id;
|
|
esr.order_ticket = send.result.order;
|
|
esr.deal_ticket = send.result.deal;
|
|
esr.payload_json = StringFormat(
|
|
"{\"send_ok\":%s,\"retcode\":%u,\"request_id\":%u,\"order\":%I64u,\"deal\":%I64u}",
|
|
(send.send_ok ? "true" : "false"), send.result.retcode, send.result.request_id,
|
|
(long)send.result.order, (long)send.result.deal);
|
|
m_journal.Add(esr, esr_id);
|
|
//--- B9: DISPATCHED только после фактического входа в OrderSend*;
|
|
//--- счётчик фактических вызовов растёт здесь (не в NextMain).
|
|
if(p.role == LAB_ROLE_MAIN && seq != 0)
|
|
m_runner.MarkDispatchedSlot(seq);
|
|
else
|
|
if(p.role == LAB_ROLE_WARMUP && seq != 0)
|
|
m_runner.MarkWarmupDispatched(seq); // R6-B4: прогрев по факту отправки
|
|
m_runner.CountDispatched(p.role);
|
|
//--- E5: VERBOSE - один фиксированный PrintFormat после регистрации возврата отправки
|
|
if(p.logging_mode == LAB_LOG_VERBOSE)
|
|
PrintFormat("[LAB][%s] send_return seq=%I64u retcode=%u request_id=%u dur=%I64u us",
|
|
m_session_id, seq, send.result.retcode, send.result.request_id, (ulong)(t1 - t0));
|
|
//--- возврат события достаточно: async обычно возвращает retcode==DONE и request_id
|
|
m_tracker.ComputeDeadlines(seq, (ulong)m_settings.outcome_timeout_ms * 1000,
|
|
(ulong)m_settings.late_grace_ms * 1000);
|
|
//--- R5-B6: рыночный режим/окно/калибровка закрепляются за записью
|
|
//--- ДО отправки; поздний callback не перезаписывает контекст следующей
|
|
//--- операции (FinalizeTransaction режим не трогает).
|
|
if(m_settings.experiment_id == LAB_EXP_E3 &&
|
|
m_last_regime != LAB_REGIME_UNKNOWN)
|
|
m_tracker.SetMarketRegime(seq, m_last_regime, m_last_window_id,
|
|
m_calibration.calibration_id,
|
|
m_last_window_seq); // R8-S2
|
|
//--- R8-B3: SEND_RETURN durable НЕ пишется синхронно между T1 и
|
|
//--- обработкой callbacks (один поток исполнения советника [D2]).
|
|
//--- Факт возврата буферизуется в памяти и дописывается в safe-фазе
|
|
//--- после закрытия окна (FlushPendingSendReturns). При аварии до
|
|
//--- этого INTENT без SEND_RETURN честно остаётся DISPATCH_UNCERTAIN.
|
|
if(do_send)
|
|
QueueDurableSendReturn(seq, t1);
|
|
m_active_sequence = seq;
|
|
return(true);
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| публичные методы (реализации) |
|
|
//+------------------------------------------------------------------+
|
|
CAppController::CAppController()
|
|
{
|
|
m_tracker = NULL;
|
|
m_journal = NULL;
|
|
m_runner = NULL;
|
|
m_clock = NULL;
|
|
m_state = LAB_STATE_IDLE;
|
|
m_active_sequence = 0;
|
|
m_operator_started = false;
|
|
m_pause_until = 0;
|
|
m_block_reason = "";
|
|
m_n_cleanup_attempts = 0;
|
|
m_last_cleanup_sequence = 0;
|
|
m_cleanup_in_flight = false;
|
|
m_cleanup_sequence = 0;
|
|
m_cleanup_parent = 0;
|
|
m_cleanup_attempt = 0;
|
|
m_last_pause_us = 0;
|
|
m_calibration.Zero();
|
|
m_calibration_loaded = false;
|
|
m_last_regime = LAB_REGIME_UNKNOWN;
|
|
m_last_window_id = "";
|
|
m_last_regime_us = 0;
|
|
m_e3_start_us = 0;
|
|
m_e3_window_seq = 0;
|
|
m_termination_reason = "";
|
|
m_cleanup_req_count = 0;
|
|
m_window_count = 0;
|
|
m_checkpoint_ok = true;
|
|
m_durable_ok = true;
|
|
m_durable_uncertain = 0;
|
|
m_last_quote_observed_us = 0;
|
|
m_last_quote_msc = 0; // R7-B5
|
|
m_pending_sr_count = 0; // R8-B3
|
|
m_last_window_seq = 0; // R8-S2
|
|
m_last_anchor_msc = 0; // R9-S2
|
|
m_window_ticks_seq = 0; // R9-S3
|
|
m_window_ticks_pending = false; // R9-S3
|
|
m_test_fail_intent_counter = 0; // R13-B1: тестовая инъекция (0=выкл)
|
|
m_session_counter = 0; // R14-B1
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
void CAppController::Configure(CRequestTracker *tracker, CEventJournal *journal,
|
|
CExperimentRunner *runner, CTimeSource *clock)
|
|
{
|
|
m_tracker = tracker;
|
|
m_journal = journal;
|
|
m_runner = runner;
|
|
m_clock = clock;
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
//--- R7-B1: запись одной строки durable-журнала. StringToCharArray с
|
|
//--- WHOLE_ARRAY копирует также завершающий ноль, поэтому массив содержит
|
|
//--- n элементов (n-1 полезных). FileWriteArray(handle,array,start,count)
|
|
//--- принимает НАЧАЛЬНЫЙ индекс и ЧИСЛО элементов [API1],[API2], поэтому
|
|
//--- пишем ровно n-1 байт с индекса 0 (без NUL) и сверяем число записей.
|
|
//--- Публично для файлового round-trip теста (RequestLatencyTests DURA-01).
|
|
static bool CAppController::AppendDurableLine(const string path, const string line, LabError &err)
|
|
{
|
|
err.Reset();
|
|
err.component = LAB_COMP_JOURNAL;
|
|
uchar b[];
|
|
const int n = StringToCharArray(line, b, 0, WHOLE_ARRAY, CP_UTF8);
|
|
if(n <= 1)
|
|
{
|
|
err.code = 65;
|
|
err.severity = LAB_SEV_BLOCKER;
|
|
err.message = "durable line encode failed";
|
|
return(false);
|
|
}
|
|
int h = FileOpen(path, FILE_READ | FILE_WRITE | FILE_BIN);
|
|
if(h == INVALID_HANDLE)
|
|
h = FileOpen(path, FILE_WRITE | FILE_BIN); // создать файл
|
|
if(h == INVALID_HANDLE)
|
|
{
|
|
err.code = 66;
|
|
err.severity = LAB_SEV_BLOCKER;
|
|
err.message = "durable line open failed";
|
|
return(false);
|
|
}
|
|
FileSeek(h, 0, SEEK_END);
|
|
const uint wrote = FileWriteArray(h, b, 0, n - 1);
|
|
FileClose(h);
|
|
if(wrote != (uint)(n - 1))
|
|
{
|
|
err.code = 67;
|
|
err.severity = LAB_SEV_BLOCKER;
|
|
err.message = "durable line short write";
|
|
return(false);
|
|
}
|
|
return(true);
|
|
}
|
|
//--- R8-B3: SEND_RETURN durable НЕ пишется синхронно между T1 и обработкой
|
|
//--- callbacks (один поток исполнения советника [D2]): намерение INTENT
|
|
//--- устойчиво до T0, а факт возврата отправки буферизуется в памяти и
|
|
//--- дописывается в safe-фазе после закрытия соответствующего окна
|
|
//--- (FlushPendingSendReturns). При аварии до этого INTENT без SEND_RETURN
|
|
//--- честно остаётся DISPATCH_UNCERTAIN. Публично для регрессий R8.
|
|
void CAppController::QueueDurableSendReturn(const ulong seq, const ulong t_us)
|
|
{
|
|
if(m_pending_sr_count >= ArraySize(m_pending_sr_seq) ||
|
|
m_pending_sr_count >= ArraySize(m_pending_sr_us))
|
|
{
|
|
const int newcap = m_pending_sr_count + 16;
|
|
if(ArrayResize(m_pending_sr_seq, newcap) != newcap ||
|
|
ArrayResize(m_pending_sr_us, newcap) != newcap)
|
|
{
|
|
//--- R9-B3: нехватка буфера SEND_RETURN — fail-closed
|
|
m_durable_ok = false;
|
|
if(m_state != LAB_STATE_BLOCKED)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "durable SEND_RETURN buffer allocation failed";
|
|
}
|
|
return;
|
|
}
|
|
}
|
|
m_pending_sr_seq[m_pending_sr_count] = seq;
|
|
m_pending_sr_us[m_pending_sr_count] = t_us;
|
|
m_pending_sr_count++;
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
|
|
//--- дописать накопленные SEND_RETURN в durable-журнал (safe-фаза: окно
|
|
//--- закрыто / экспорт). Сбой — false: продолжать торговлю нельзя.
|
|
bool CAppController::FlushPendingSendReturns(LabError &err)
|
|
{
|
|
if(m_pending_sr_count == 0)
|
|
return(true);
|
|
if(m_session_id == NULL || m_session_id == "")
|
|
return(false); // без сессии записать некуда — неопределённость
|
|
const string dir = "RequestLatencyLab\\" + m_settings.campaign_id + "\\" + m_session_id;
|
|
if(!FolderCreate(dir))
|
|
return(false);
|
|
for(int i = 0; i < m_pending_sr_count; i++)
|
|
{
|
|
const string line = StringFormat("%s,%I64u,SEND_RETURN,%I64u\r\n",
|
|
m_session_id, m_pending_sr_seq[i],
|
|
m_pending_sr_us[i]);
|
|
if(!AppendDurableLine(dir + "\\intents.csv", line, err))
|
|
return(false);
|
|
}
|
|
m_pending_sr_count = 0;
|
|
return(true);
|
|
}
|
|
//--- R8-B2: ЕДИНАЯ политика полноты ИЗМЕРИТЕЛЬНЫХ точек перед закрытием
|
|
//--- collection: REQUEST (T2) — точка всех операций с торговым вызовом;
|
|
//--- ORDER_ADD (T3) — операций, где ордер является ожидаемым событием
|
|
//--- (кроме однозначного отказа без ордера). Используется CompletedCheck,
|
|
//--- OnTimer и (через них) cleanup; при неполноте закрывает grace.
|
|
static bool CAppController::MeasurementPointsComplete(const RequestMetadata &md)
|
|
{
|
|
const uint pm = md.present_mask;
|
|
const ENUM_LAB_OPERATION op = md.plan.operation;
|
|
const bool request_expected = (op == LAB_OP_MARKET_OPEN ||
|
|
op == LAB_OP_POSITION_CLOSE ||
|
|
op == LAB_OP_PENDING_CREATE ||
|
|
op == LAB_OP_PENDING_DELETE);
|
|
if(request_expected && (pm & (uint)LAB_MASK_T2) == 0)
|
|
return(false);
|
|
bool order_point_expected = (op == LAB_OP_MARKET_OPEN ||
|
|
op == LAB_OP_POSITION_CLOSE ||
|
|
op == LAB_OP_PENDING_CREATE);
|
|
//--- однозначный отказ без ордера: ORDER_ADD не ожидается
|
|
if(order_point_expected && md.outcome == LAB_OUT_REJECTED &&
|
|
md.order_ticket == 0)
|
|
order_point_expected = false;
|
|
if(order_point_expected && (pm & (uint)LAB_MASK_T3) == 0)
|
|
return(false);
|
|
return(true);
|
|
}
|
|
//--- R10-B3: критерий готовности закрытия коллекции cleanup - единая
|
|
//--- политика с CompletedCheck: измерительные точки (T2/T3 по операции)
|
|
//--- полны И полнота DEAL-callback (T4/T5/callback_coverage_complete)
|
|
//--- подтверждена для операций с ожидаемой сделкой; REQUEST-only
|
|
//--- (PENDING_DELETE без сделки) DEAL не ждёт; неполноту закрывает grace.
|
|
static bool CAppController::CleanupCollectionReady(const RequestMetadata &md)
|
|
{
|
|
if(!MeasurementPointsComplete(md))
|
|
return(false);
|
|
const bool deal_callback_expected = (md.plan.operation == LAB_OP_MARKET_OPEN ||
|
|
md.plan.operation == LAB_OP_POSITION_CLOSE ||
|
|
md.plan.operation == LAB_OP_PENDING_CREATE);
|
|
return(!deal_callback_expected || md.callback_coverage_complete);
|
|
}
|
|
//--- R8-B3/R9-B3: подтверждение durable-фазы после закрытия окна —
|
|
//--- дописывание накопленных SEND_RETURN и (R9-S3) сырых тиков окна.
|
|
//--- Сбой записи — fail-closed: возвращает false, новые отправки
|
|
//--- блокируются, вызывающий код ОБЯЗАН прервать текущую ветвь и НЕ
|
|
//--- перезаписывать терминальное состояние BLOCKED.
|
|
bool CAppController::ConfirmDurableSendReturns(void)
|
|
{
|
|
LabError derr;
|
|
if(!FlushPendingSendReturns(derr))
|
|
{
|
|
m_durable_ok = false;
|
|
if(m_state != LAB_STATE_BLOCKED)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "durable SEND_RETURN flush failed";
|
|
}
|
|
return(false);
|
|
}
|
|
LabError terr;
|
|
if(!FlushWindowTicks(terr))
|
|
{
|
|
m_durable_ok = false;
|
|
if(m_state != LAB_STATE_BLOCKED)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "window ticks dump failed";
|
|
}
|
|
return(false);
|
|
}
|
|
return(true);
|
|
}
|
|
//--- R9-S3: выгрузка сырых тиков последнего окна E3 в safe-фазе
|
|
//--- (после закрытия соответствующего sample; внутри измерительной
|
|
//--- фазы синхронная запись запрещена R8-B3). Append-only market_ticks.csv:
|
|
//--- session_id,window_id,window_seq,time_msc,bid,ask,flags,volume.
|
|
//--- Сбой записи — false; вызывающий код блокирует торговлю.
|
|
//--- R10-B4: заголовок/строка market_ticks.csv; статические, чтобы формат
|
|
//--- проверялся отдельным round-trip тестом (WINDOW-ROUNDTRIP).
|
|
static string CAppController::MarketTickHeader()
|
|
{
|
|
return("session_id,window_id,window_seq,time_msc,bid,ask,flags,volume\r\n");
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
static string CAppController::FormatMarketTickLine(const string session_id, const string wid,
|
|
const ulong window_seq, const MqlTick &t)
|
|
{
|
|
return(StringFormat("%s,%s,%I64u,%I64u,%.17g,%.17g,%d,%I64u\r\n",
|
|
session_id, wid, window_seq, (ulong)t.time_msc,
|
|
t.bid, t.ask, (int)t.flags, (long)t.volume));
|
|
}
|
|
//--- R10-B4: файл отсутствует/пуст - нужен заголовок (FileSize требует
|
|
//--- handle, а не путь).
|
|
static bool CAppController::TickFileEmpty(const string path)
|
|
{
|
|
const int h = FileOpen(path, FILE_READ | FILE_BIN);
|
|
if(h == INVALID_HANDLE)
|
|
return(true);
|
|
const bool empty = (FileSize(h) <= 0);
|
|
FileClose(h);
|
|
return(empty);
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::FlushWindowTicks(LabError &err)
|
|
{
|
|
err.Reset();
|
|
if(!m_window_ticks_pending || m_window_ticks_seq == 0)
|
|
return(true);
|
|
if(m_session_id == NULL || m_session_id == "")
|
|
return(false);
|
|
const string dir = "RequestLatencyLab\\" + m_settings.campaign_id + "\\" + m_session_id;
|
|
if(!FolderCreate(dir))
|
|
return(false);
|
|
string wid = "";
|
|
for(int i = m_window_count - 1; i >= 0; i--)
|
|
if(m_windows[i].sequence == m_window_ticks_seq)
|
|
{
|
|
wid = m_windows[i].window_id;
|
|
break;
|
|
}
|
|
if(StringLen(wid) == 0)
|
|
{
|
|
//--- окно не сохранилось в протоколе — данные выгрузить нельзя
|
|
m_window_ticks_pending = false;
|
|
return(true);
|
|
}
|
|
//--- R10-B4: заголовок пишется при создании/пустом файле (append-only
|
|
//--- журнал окна); точность bid/ask %.17g, разделитель - реальный CRLF.
|
|
if(TickFileEmpty(dir + "\\market_ticks.csv"))
|
|
{
|
|
LabError herr;
|
|
if(!AppendDurableLine(dir + "\\market_ticks.csv", MarketTickHeader(), herr))
|
|
return(false);
|
|
}
|
|
const int n = ArraySize(m_window_ticks);
|
|
for(int i = 0; i < n; i++)
|
|
{
|
|
const MqlTick t = m_window_ticks[i];
|
|
const string line = FormatMarketTickLine(m_session_id, wid,
|
|
m_window_ticks_seq, t);
|
|
if(!AppendDurableLine(dir + "\\market_ticks.csv", line, err))
|
|
return(false);
|
|
}
|
|
m_window_ticks_pending = false;
|
|
m_window_ticks_seq = 0;
|
|
return(true);
|
|
}
|
|
//--- R10-B4: гарантировать наличие market_ticks.csv (заголовок) для сессии
|
|
//--- - контракт экспорта (пилот LOCAL_ONLY проверяет файл); строки тиков
|
|
//--- окна дописываются FlushWindowTicks. Fail-closed при сбое.
|
|
bool CAppController::EnsureMarketTicksFile(void)
|
|
{
|
|
if(m_session_id == NULL || m_session_id == "")
|
|
return(false);
|
|
const string dir = "RequestLatencyLab\\" + m_settings.campaign_id + "\\" + m_session_id;
|
|
if(!FolderCreate(dir))
|
|
return(false);
|
|
const string path = dir + "\\market_ticks.csv";
|
|
if(!TickFileEmpty(path))
|
|
return(true);
|
|
LabError herr;
|
|
return(AppendDurableLine(path, MarketTickHeader(), herr));
|
|
}
|
|
//--- R5-B1/R7-B3: единый владелец перехода после завершения цикла записи:
|
|
//--- WARMUP пока не исчерпана квота прогрева, иначе MAIN. Публично для
|
|
//--- регрессионных тестов перехода (RequestLatencyTests WARM-02).
|
|
void CAppController::ReturnToSchedule(void)
|
|
{
|
|
if(!m_runner.WarmupQuotaMet())
|
|
{
|
|
if(m_runner.WarmupExhausted())
|
|
TerminateRun("WARMUP_INSUFFICIENT_DATA");
|
|
else
|
|
m_state = LAB_STATE_WARMUP;
|
|
return;
|
|
}
|
|
m_state = LAB_STATE_MAIN;
|
|
}
|
|
//--- R8-S2/R9-S2: устарел ли контекст E3 для отправки. Два условия:
|
|
//--- (1) возраст последнего НОВОГО наблюдения рабочего символа
|
|
//--- (монотонный time_msc, NoteQuoteObserved); (2) возраст ЯКОРЯ
|
|
//--- выбранного окна относительно текущей котировки — проверяется
|
|
//--- именно тот контекст, который будет закреплён за sample.
|
|
bool CAppController::E3ContextStale(const ulong now_us)
|
|
{
|
|
if(m_settings.feature_max_age_ms <= 0)
|
|
return(false);
|
|
const ulong max_age = (ulong)m_settings.feature_max_age_ms * 1000;
|
|
if(m_last_quote_observed_us > 0 && now_us > m_last_quote_observed_us &&
|
|
now_us - m_last_quote_observed_us > max_age)
|
|
return(true);
|
|
if(m_last_anchor_msc > 0)
|
|
{
|
|
MqlTick qw;
|
|
if(SymbolInfoTick(m_settings.symbol, qw) && qw.time_msc > m_last_anchor_msc &&
|
|
(ulong)(qw.time_msc - m_last_anchor_msc) * 1000 > max_age)
|
|
return(true);
|
|
}
|
|
return(false);
|
|
}
|
|
//--- следующий разрешённый шаг (прогрев -> основные -> сбор)
|
|
bool CAppController::DriveNext(LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
//--- B10: потеря журнала -> BLOCKED, торговля прекращается
|
|
if(!CheckIntegrity(error))
|
|
return(false);
|
|
if(m_state == LAB_STATE_WARMUP)
|
|
{
|
|
//--- R5-B1: прогрев строго последователен — одна операция в работе,
|
|
//--- пауза и незавершённый cleanup блокируют следующий прогревочный
|
|
//--- шаг; гарантируется 10 завершённых прогревочных вызовов на
|
|
//--- условие до перехода в MAIN.
|
|
if(m_active_sequence != 0 || m_cleanup_in_flight)
|
|
return(true);
|
|
if(m_clock.NowUs() < m_pause_until)
|
|
return(true);
|
|
RequestPlan warm;
|
|
if(m_runner.NextWarmup(warm))
|
|
{
|
|
ulong seq = 0;
|
|
return(PrepareAndDispatch(warm, true, seq, error));
|
|
}
|
|
if(!m_runner.WarmupExhausted())
|
|
return(true); // R6-B4: впереди прогревочные слоты (requeue)
|
|
//--- R7-B3: единая проверка квоты прогрева на ВСЕХ переходах в MAIN
|
|
//--- (ReturnToSchedule); недобор фактических отправок при исчерпанном
|
|
//--- расписании — терминальный WARMUP_INSUFFICIENT_DATA.
|
|
ReturnToSchedule();
|
|
return(true);
|
|
}
|
|
if(m_state == LAB_STATE_MAIN)
|
|
{
|
|
if(m_cleanup_in_flight)
|
|
return(true); // R5-B1: cleanup не подтверждён
|
|
if(m_active_sequence != 0)
|
|
return(true); // одна основная операция в работе
|
|
if(m_clock.NowUs() < m_pause_until)
|
|
return(true);
|
|
//--- R4-B7: прежде чем начинать новый слот, проверяем, что при
|
|
//--- исчерпанном расписании квота фактических send выполнена.
|
|
if(m_runner.MainComplete() && m_runner.MainDispatched() <
|
|
m_runner.RequiredDispatched())
|
|
{
|
|
TerminateRun("INSUFFICIENT_DATA");
|
|
return(true);
|
|
}
|
|
RequestPlan plan;
|
|
bool have = false;
|
|
if(m_settings.experiment_id == LAB_EXP_E3)
|
|
{
|
|
//--- аудит-2 S5: калибровка обязательна для E3
|
|
if(!m_calibration_loaded)
|
|
{
|
|
TerminateRun("E3_CALIBRATION_MISSING");
|
|
return(true);
|
|
}
|
|
if(m_e3_start_us == 0)
|
|
m_e3_start_us = m_clock.NowUs();
|
|
//--- окно рынка + классификация; отправка только при совпадении режима
|
|
ENUM_LAB_MARKET_REGIME regime = LAB_REGIME_UNKNOWN;
|
|
UpdateMarketRegime(regime);
|
|
//--- R8-B1: терминальный/блокирующий переход ВНУТРИ
|
|
//--- UpdateMarketRegime (ёмкость/ошибка хранения окна) —
|
|
//--- отправка в этом же проходе запрещена.
|
|
if(m_state != LAB_STATE_MAIN)
|
|
return(true);
|
|
if(regime == LAB_REGIME_QUIET || regime == LAB_REGIME_FAST)
|
|
have = m_runner.TryDispatchByRegime(regime, plan);
|
|
if(!have)
|
|
{
|
|
//--- R4-B7: квоты обоих условий заполнены фактическими отправками
|
|
if(m_runner.E3QuotaComplete())
|
|
{
|
|
m_active_sequence = 0;
|
|
m_state = LAB_STATE_DRAINING; // серия E3 завершена штатно
|
|
return(true);
|
|
}
|
|
//--- формальное завершение E3 при недоборе режима/слота
|
|
const ulong limit_us = (ulong)m_settings.e3_campaign_limit_min * 60UL * 1000000UL;
|
|
if(m_settings.e3_campaign_limit_min > 0 &&
|
|
m_clock.NowUs() - m_e3_start_us >= limit_us)
|
|
TerminateRun("E3_INSUFFICIENT_DATA");
|
|
return(true); // ждём подходящего режима, слот не расходуем
|
|
}
|
|
}
|
|
if(!have)
|
|
have = m_runner.NextMain(plan);
|
|
if(have)
|
|
{
|
|
ulong seq = 0;
|
|
//--- B9: статус/счётчик precheck-отказа фиксируются внутри
|
|
//--- PrepareAndDispatch; здесь повторный счёт не выполняется.
|
|
if(!PrepareAndDispatch(plan, true, seq, error))
|
|
return(true); // отказ не заменяется успешной попыткой
|
|
}
|
|
else
|
|
m_state = LAB_STATE_DRAINING;
|
|
return(true);
|
|
}
|
|
return(true);
|
|
}
|
|
//--- пауза после цикла/очистки/сохранения (ТЗ §5)
|
|
void CAppController::RequirePause(void)
|
|
{
|
|
m_pause_until = m_clock.NowUs() + (ulong)m_settings.pause_ms * 1000;
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::PauseElapsed() const
|
|
{
|
|
return(m_clock.NowUs() >= m_pause_until);
|
|
}
|
|
//--- инициализация буферов и плана (OnInit)
|
|
bool CAppController::InitRun(LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(m_session_id == NULL || m_session_id == "")
|
|
m_session_id = NewSessionId();
|
|
if(!m_journal.Init(m_settings.event_capacity, m_session_id))
|
|
{
|
|
error.code = 40;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "journal init failed";
|
|
return(false);
|
|
}
|
|
if(!m_tracker.Init(m_settings.sample_capacity, m_session_id))
|
|
{
|
|
error.code = 41;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "tracker init failed";
|
|
return(false);
|
|
}
|
|
if(!m_runner.Init(m_settings, error))
|
|
return(false);
|
|
return(true);
|
|
}
|
|
//--- завершение цикла основной записи: очистка при известном состоянии
|
|
bool CAppController::EndMainCycle(LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(m_active_sequence == 0)
|
|
{
|
|
//--- R5-B1: cleanup в работе — остаёмся в DRAINING; единственный
|
|
//--- владелец перехода в MAIN — DriveCleanup после подтверждения
|
|
//--- нулевой остаточной экспозиции.
|
|
if(m_cleanup_in_flight)
|
|
return(true);
|
|
ReturnToSchedule();
|
|
return(true);
|
|
}
|
|
RequestMetadata r;
|
|
if(!m_tracker.GetMetadata(m_active_sequence, r))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "active record lost";
|
|
return(false);
|
|
}
|
|
//--- аудит-2 B3: cleanup только после закрытия коллекции
|
|
if(!r.collection_closed)
|
|
{
|
|
const ulong now = m_clock.NowUs();
|
|
if(m_tracker.GraceExpired(m_active_sequence, now))
|
|
{
|
|
ulong cid = 0;
|
|
LabEvent ce;
|
|
ce.Zero();
|
|
ce.kind = LAB_KIND_COLLECTION_CLOSE;
|
|
ce.local_time_us = now;
|
|
ce.origin_sequence = m_active_sequence;
|
|
ce.payload_json = "{\"reason\":\"grace_expired\"}";
|
|
m_journal.Add(ce, cid);
|
|
m_tracker.CloseCollection(m_active_sequence, now, cid);
|
|
if(!ConfirmDurableSendReturns())
|
|
return(false); // R9-B3: сбой flush - BLOCKED
|
|
}
|
|
else
|
|
return(false); // ждём grace / поздние события
|
|
}
|
|
if(!r.trading_state_known)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "trading state unknown after close";
|
|
return(false);
|
|
}
|
|
bool need_cleanup = false;
|
|
MqlTradeRequest req;
|
|
ZeroMemory(req);
|
|
if(r.plan.operation == LAB_OP_PENDING_CREATE)
|
|
{
|
|
//--- R6-B5: UNEXPECTED_ACTIVATION оставляет ПОЗИЦИЮ, а не ордер;
|
|
//--- остаточная экспозиция обрабатывается адресно по её типу.
|
|
ulong pos_ticket = r.position_ticket;
|
|
if(pos_ticket == 0)
|
|
{
|
|
ResetLastError();
|
|
for(int i = PositionsTotal() - 1; i >= 0; i--)
|
|
{
|
|
const ulong t = PositionGetTicket(i);
|
|
if(t == 0 || PositionGetString(POSITION_SYMBOL) != r.plan.symbol)
|
|
continue;
|
|
if((ulong)PositionGetInteger(POSITION_MAGIC) == r.plan.magic)
|
|
{
|
|
pos_ticket = t;
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
if(pos_ticket != 0 && PositionSelectByTicket(pos_ticket))
|
|
{
|
|
req.action = TRADE_ACTION_DEAL;
|
|
req.symbol = r.plan.symbol;
|
|
req.position = (ulong)pos_ticket;
|
|
req.volume = PositionGetDouble(POSITION_VOLUME);
|
|
const long ptype = PositionGetInteger(POSITION_TYPE);
|
|
req.type = (ptype == POSITION_TYPE_BUY ? ORDER_TYPE_SELL : ORDER_TYPE_BUY);
|
|
req.type_filling = SelectFilling(r.plan.symbol, false);
|
|
req.deviation = m_settings.deviation_points;
|
|
need_cleanup = true;
|
|
}
|
|
else
|
|
if(r.order_ticket != 0)
|
|
{
|
|
req.action = TRADE_ACTION_REMOVE;
|
|
req.symbol = r.plan.symbol;
|
|
req.order = (ulong)r.order_ticket;
|
|
req.magic = (long)r.plan.magic;
|
|
req.type_filling = SelectFilling(r.plan.symbol, true);
|
|
need_cleanup = true;
|
|
}
|
|
}
|
|
else
|
|
if(r.plan.operation == LAB_OP_MARKET_OPEN)
|
|
{
|
|
ulong pos_ticket = r.position_ticket;
|
|
if(pos_ticket == 0)
|
|
{
|
|
ResetLastError();
|
|
for(int i = PositionsTotal() - 1; i >= 0; i--)
|
|
{
|
|
const ulong t = PositionGetTicket(i);
|
|
if(t == 0 || PositionGetString(POSITION_SYMBOL) != r.plan.symbol)
|
|
continue;
|
|
if((ulong)PositionGetInteger(POSITION_MAGIC) == r.plan.magic)
|
|
{
|
|
pos_ticket = t;
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
if(pos_ticket != 0 && PositionSelectByTicket(pos_ticket))
|
|
{
|
|
req.action = TRADE_ACTION_DEAL;
|
|
req.symbol = r.plan.symbol;
|
|
req.position = (ulong)pos_ticket;
|
|
req.volume = PositionGetDouble(POSITION_VOLUME);
|
|
const long ptype = PositionGetInteger(POSITION_TYPE);
|
|
req.type = (ptype == POSITION_TYPE_BUY ? ORDER_TYPE_SELL : ORDER_TYPE_BUY);
|
|
//--- S6: cleanup через те же торговые правила (filling/deviation)
|
|
req.type_filling = SelectFilling(r.plan.symbol, false);
|
|
req.deviation = m_settings.deviation_points;
|
|
need_cleanup = true;
|
|
}
|
|
}
|
|
if(need_cleanup)
|
|
{
|
|
const ENUM_LAB_OPERATION cop = (req.action == TRADE_ACTION_REMOVE ?
|
|
LAB_OP_PENDING_DELETE : LAB_OP_POSITION_CLOSE);
|
|
const ulong cseq = SendCleanup(req, m_active_sequence, cop, error);
|
|
if(cseq == 0)
|
|
{
|
|
//--- неопределённая/оставшаяся экспозиция -> BLOCKED (B1)
|
|
if(m_state != LAB_STATE_BLOCKED)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cleanup dispatch failed / precheck rejected";
|
|
}
|
|
return(false);
|
|
}
|
|
//--- cleanup в работе: держим DRAINING до подтверждения нулевой
|
|
//--- экспозиции; R9-B3: терминальный переход внутри SendCleanup
|
|
//--- (BLOCKED) не перезаписывается DRAINING.
|
|
if(m_state != LAB_STATE_BLOCKED)
|
|
{
|
|
m_active_sequence = 0;
|
|
m_state = LAB_STATE_DRAINING;
|
|
}
|
|
return(true);
|
|
}
|
|
m_active_sequence = 0;
|
|
//--- R5-B1: следующий режим зависит от фазы расписания (WARMUP/MAIN)
|
|
ReturnToSchedule();
|
|
//--- без cleanup (нет остатка) — следующий шаг после паузы
|
|
RequirePause();
|
|
return(true);
|
|
}
|
|
//--- B1: завершение lifecycle cleanup. Возвращает true, когда чистка
|
|
//--- подтверждена (нулевая остаточная экспозиция) и можно следующий MAIN.
|
|
bool CAppController::DriveCleanup(LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(!m_cleanup_in_flight || m_cleanup_sequence == 0)
|
|
{
|
|
m_cleanup_in_flight = false;
|
|
return(true);
|
|
}
|
|
RequestMetadata cm;
|
|
if(!m_tracker.GetMetadata(m_cleanup_sequence, cm))
|
|
{
|
|
//--- R4-B1: потеря записи cleanup / journal failure — BLOCKED.
|
|
//--- Неизвестное состояние торговой очистки нельзя считать безопасным.
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cleanup record lost";
|
|
return(false);
|
|
}
|
|
//--- R10-B3: решение о закрытии коллекции cleanup принимается по СВЕЖЕЙ
|
|
//--- сверке: полнота DEAL-callback (T4/T5/callback_coverage_complete)
|
|
//--- переносится в запись ДО проверки условия закрытия.
|
|
{
|
|
const ulong now0 = m_clock.NowUs();
|
|
ConfirmationSnapshot snap0;
|
|
CStateReader::BuildSnapshot(m_cleanup_sequence, cm.plan.operation, cm.order_ticket,
|
|
cm.plan.volume / m_rules.volume_step,
|
|
m_rules.volume_step, m_settings.volume_units_tolerance,
|
|
m_clock, cm.position_ticket,
|
|
m_settings.history_scan_budget, snap0);
|
|
snap0.requested_volume_units = cm.plan.volume / m_rules.volume_step;
|
|
snap0.order_state_after = snap0.order_state_before;
|
|
SyncRecoveredDeals(m_cleanup_sequence, snap0);
|
|
m_tracker.SyncReconcile(m_cleanup_sequence, snap0, m_rules.volume_step,
|
|
m_settings.volume_units_tolerance);
|
|
CompletedCheck(m_cleanup_sequence, snap0, error);
|
|
if(m_state == LAB_STATE_BLOCKED)
|
|
return(false); // R9-B3: сбой flush внутри CompletedCheck
|
|
if(!m_tracker.GetMetadata(m_cleanup_sequence, cm))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cleanup record lost after reconcile";
|
|
return(false);
|
|
}
|
|
}
|
|
if(!cm.collection_closed)
|
|
{
|
|
//--- R5-B1: закрытие collection для cleanup не зависит от наличия T6:
|
|
//--- раннее закрытие после T2 (поздний REQUEST сохраняется), либо
|
|
//--- закрытие по grace при истечении (T6 уже есть — grace всё равно
|
|
//--- должен закрыть коллекцию, иначе запись зависнет).
|
|
const ulong now = m_clock.NowUs();
|
|
const bool grace_expired = m_tracker.GraceExpired(m_cleanup_sequence, now);
|
|
//--- R9-B2: ЕДИНАЯ политика полноты точек закрытия (как в CompletedCheck
|
|
//--- /OnTimer): T2 и T3 по операции (REJECTED без ордера освобождает T3);
|
|
//--- T2 без T3 НЕ завершает наблюдение cleanup раньше времени —
|
|
//--- неполноту закрывает grace.
|
|
const bool points_complete = MeasurementPointsComplete(cm);
|
|
if((points_complete && CleanupCollectionReady(cm)) || grace_expired)
|
|
{
|
|
ulong cid = 0;
|
|
LabEvent ce;
|
|
ce.Zero();
|
|
ce.kind = LAB_KIND_COLLECTION_CLOSE;
|
|
ce.local_time_us = now;
|
|
ce.origin_sequence = m_cleanup_sequence;
|
|
ce.payload_json = StringFormat("{\"reason\":\"%s\"}",
|
|
(grace_expired ? "grace_expired" : "points_complete"));
|
|
m_journal.Add(ce, cid);
|
|
m_tracker.CloseCollection(m_cleanup_sequence, now, cid);
|
|
if(!ConfirmDurableSendReturns())
|
|
return(false); // R9-B3: сбой flush - BLOCKED
|
|
}
|
|
else
|
|
return(false); // ждём REQUEST либо grace
|
|
}
|
|
if((cm.present_mask & (uint)LAB_MASK_T6) == 0)
|
|
{
|
|
//--- сверка (полный reconcile в таймере разрешён)
|
|
ConfirmationSnapshot snap;
|
|
CStateReader::BuildSnapshot(m_cleanup_sequence, cm.plan.operation, cm.order_ticket,
|
|
cm.plan.volume / m_rules.volume_step,
|
|
m_rules.volume_step, m_settings.volume_units_tolerance,
|
|
m_clock, cm.position_ticket,
|
|
m_settings.history_scan_budget, snap);
|
|
snap.requested_volume_units = cm.plan.volume / m_rules.volume_step;
|
|
snap.order_state_after = snap.order_state_before;
|
|
//--- R5-B4: сделки из истории регистрируются в реестре (deals.csv)
|
|
SyncRecoveredDeals(m_cleanup_sequence, snap);
|
|
m_tracker.SyncReconcile(m_cleanup_sequence, snap, m_rules.volume_step,
|
|
m_settings.volume_units_tolerance);
|
|
CompletedCheck(m_cleanup_sequence, snap, error);
|
|
if(m_state == LAB_STATE_BLOCKED)
|
|
return(false); // R9-B3: сбой flush внутри CompletedCheck
|
|
//--- повторное чтение после возможного T6
|
|
RequestMetadata c2;
|
|
if(!m_tracker.GetMetadata(m_cleanup_sequence, c2))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cleanup record lost after check";
|
|
return(false);
|
|
}
|
|
if((c2.present_mask & (uint)LAB_MASK_T6) == 0)
|
|
{
|
|
//--- R6-B5: grace истёк, торговое состояние неопределённо —
|
|
//--- конечный BLOCKED (не бесконечное DRAINING).
|
|
const bool gexp = m_tracker.GraceExpired(m_cleanup_sequence,
|
|
m_clock.NowUs());
|
|
if(c2.collection_closed && !c2.trading_state_known && gexp)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cleanup trading state unknown after grace";
|
|
return(false);
|
|
}
|
|
return(false); // ещё не завершён
|
|
}
|
|
}
|
|
//--- подтверждение нулевой остаточной экспозиции (B1)
|
|
//--- R4-B1: при остатке — BLOCKED (следующий MAIN не разрешается).
|
|
if(!ZeroResidualExposure(cm.plan.operation, cm.plan.symbol, cm.plan.magic))
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "cleanup residual exposure unconfirmed (attempts=" +
|
|
IntegerToString(m_cleanup_attempt) + "/" +
|
|
IntegerToString(m_settings.cleanup_max_attempts) + ")";
|
|
return(false);
|
|
}
|
|
//--- чистка подтверждена: возврат к MAIN, счётчик не дублируем
|
|
m_cleanup_in_flight = false;
|
|
m_cleanup_sequence = 0;
|
|
m_cleanup_parent = 0;
|
|
ReturnToSchedule();
|
|
RequirePause();
|
|
return(true);
|
|
}
|
|
//--- B1: подтверждение, что остатка (позиции/ордера) по символу/magic нет
|
|
static bool CAppController::ZeroResidualExposure(const ENUM_LAB_OPERATION op, const string symbol,
|
|
const ulong magic)
|
|
{
|
|
for(int i = PositionsTotal() - 1; i >= 0; i--)
|
|
{
|
|
const ulong tt = PositionGetTicket(i);
|
|
if(tt == 0)
|
|
continue;
|
|
if(PositionGetString(POSITION_SYMBOL) != symbol)
|
|
continue;
|
|
if(magic != 0 && (ulong)PositionGetInteger(POSITION_MAGIC) != magic)
|
|
continue;
|
|
return(false); // осталась позиция
|
|
}
|
|
for(int i = OrdersTotal() - 1; i >= 0; i--)
|
|
{
|
|
const ulong tt = OrderGetTicket(i);
|
|
if(tt == 0)
|
|
continue;
|
|
if(OrderGetString(ORDER_SYMBOL) != symbol)
|
|
continue;
|
|
if(magic != 0 && (ulong)OrderGetInteger(ORDER_MAGIC) != magic)
|
|
continue;
|
|
return(false); // остался отложенный/активный ордер
|
|
}
|
|
return(true);
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
CRequestTracker *CAppController::Tracker() { return(m_tracker); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
CExperimentRunner *CAppController::Runner() { return(m_runner); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
SymbolRules CAppController::Rules() { return(m_rules); }
|
|
//--- B10 (audit-3): потеря протокольных событий — блокирующее условие
|
|
bool CAppController::CheckIntegrity(LabError &error)
|
|
{
|
|
if(m_journal != NULL && m_journal.LostCount() > 0)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_JOURNAL;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.code = 60;
|
|
error.message = "journal overflow: records lost=" +
|
|
IntegerToString(m_journal.LostCount());
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "journal overflow";
|
|
return(false);
|
|
}
|
|
//--- R6-B6: неподтверждённая сохранность checkpoint/durable-журнала
|
|
//--- останавливает новые отправки (BLOCKED), экспорт остаётся доступен.
|
|
if(!m_checkpoint_ok)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CSV;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.code = 62;
|
|
error.message = "checkpoint write failed";
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "checkpoint write failed";
|
|
return(false);
|
|
}
|
|
if(!m_durable_ok)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_JOURNAL;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.code = 64;
|
|
error.message = "durable intent log write failed";
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "durable intent log write failed";
|
|
return(false);
|
|
}
|
|
//--- R7-S4: потеря первичных фактов сделок (callback/history) —
|
|
//--- fail-closed для дальнейших отправок; итоговая RUN-DEAL-02
|
|
//--- не заменяет своевременную блокировку.
|
|
if(m_tracker != NULL && m_tracker.DealLostCount() > 0)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_TRACKER;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.code = 68;
|
|
error.message = "deal registry overflow: lost=" +
|
|
IntegerToString(m_tracker.DealLostCount());
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "deal registry overflow";
|
|
return(false);
|
|
}
|
|
return(true);
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
ENUM_LAB_STATE CAppController::State() const { return(m_state); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
string CAppController::SessionId() const { return(m_session_id); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
string CAppController::BlockReason() const { return(m_block_reason); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
ulong CAppController::ActiveSequence() const { return(m_active_sequence); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
int CAppController::MarketWindowCount() const { return(m_window_count); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::GetMarketWindow(const int index, MarketWindow &out) const
|
|
{
|
|
if(index < 0 || index >= m_window_count)
|
|
return(false);
|
|
out = m_windows[index];
|
|
return(true);
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
string CAppController::StateName() const
|
|
{
|
|
switch(m_state)
|
|
{
|
|
case LAB_STATE_IDLE:
|
|
return("IDLE");
|
|
case LAB_STATE_PREFLIGHT:
|
|
return("PREFLIGHT");
|
|
case LAB_STATE_READY:
|
|
return("READY");
|
|
case LAB_STATE_WARMUP:
|
|
return("WARMUP");
|
|
case LAB_STATE_MAIN:
|
|
return("MAIN");
|
|
case LAB_STATE_DRAINING:
|
|
return("DRAINING");
|
|
case LAB_STATE_EXPORTING:
|
|
return("EXPORTING");
|
|
case LAB_STATE_FINISHED:
|
|
return("FINISHED");
|
|
case LAB_STATE_PAUSED:
|
|
return("PAUSED");
|
|
case LAB_STATE_BLOCKED:
|
|
return("BLOCKED");
|
|
}
|
|
return("?");
|
|
}
|
|
//--- предварительная проверка среды (без торговли)
|
|
bool CAppController::Preflight(LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(!CConfiguration::Validate(m_settings, error))
|
|
return(false);
|
|
if(!CConfiguration::BuildSymbolRules(m_settings.symbol, m_rules, error))
|
|
return(false);
|
|
//--- B8 (audit-3): E4 НЕ создаёт новых торговых попыток — это анализ
|
|
//--- пяти ASYNC-серий E1 (reuse). Торговый runtime запрещён для E4.
|
|
if(m_settings.experiment_id == LAB_EXP_E4 && m_settings.run_mode != LAB_RUN_LOCAL_ONLY)
|
|
{
|
|
error.code = 14;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "E4 is reuse-analysis of E1; trading runtime not allowed";
|
|
return(false);
|
|
}
|
|
//--- торговые режимы разрешены только на демо
|
|
if(m_settings.run_mode != LAB_RUN_LOCAL_ONLY)
|
|
{
|
|
if(AccountInfoInteger(ACCOUNT_TRADE_MODE) != ACCOUNT_TRADE_MODE_DEMO)
|
|
{
|
|
error.code = 10;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "trading run requires demo account";
|
|
return(false);
|
|
}
|
|
if(!TerminalInfoInteger(TERMINAL_TRADE_ALLOWED) ||
|
|
!MQLInfoInteger(MQL_TRADE_ALLOWED) ||
|
|
!AccountInfoInteger(ACCOUNT_TRADE_EXPERT) ||
|
|
!AccountInfoInteger(ACCOUNT_TRADE_ALLOWED))
|
|
{
|
|
error.code = 11;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "trading permissions missing";
|
|
return(false);
|
|
}
|
|
}
|
|
if(m_settings.run_mode != LAB_RUN_LOCAL_ONLY)
|
|
{
|
|
//--- аудит-2 S6: символ не должен содержать посторонних позиций/ордеров
|
|
for(int i = PositionsTotal() - 1; i >= 0; i--)
|
|
{
|
|
const ulong tt = PositionGetTicket(i);
|
|
if(tt == 0)
|
|
continue;
|
|
if(PositionGetString(POSITION_SYMBOL) != m_settings.symbol)
|
|
continue;
|
|
error.code = 12;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "symbol polluted: open position exists";
|
|
return(false);
|
|
}
|
|
for(int i = OrdersTotal() - 1; i >= 0; i--)
|
|
{
|
|
const ulong tt = OrderGetTicket(i);
|
|
if(tt == 0)
|
|
continue;
|
|
if(OrderGetString(ORDER_SYMBOL) != m_settings.symbol)
|
|
continue;
|
|
error.code = 13;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "symbol polluted: open order exists";
|
|
return(false);
|
|
}
|
|
}
|
|
m_settings.configuration_id = CConfiguration::ConfigurationId(m_settings, m_rules);
|
|
m_state = LAB_STATE_READY;
|
|
return(true);
|
|
}
|
|
//--- явный Start оператора
|
|
bool CAppController::StartRun(LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(m_state != LAB_STATE_READY && m_state != LAB_STATE_IDLE)
|
|
{
|
|
error.code = 20;
|
|
error.message = "start not allowed in state " + StateName();
|
|
return(false);
|
|
}
|
|
if(m_settings.run_mode == LAB_RUN_LOCAL_ONLY)
|
|
{
|
|
//--- синтетический прогон без реальной торговли
|
|
ResetSessionWindows(); // R13-B3: оконный контекст новой сессии чист
|
|
if(!RunSynthetic(error))
|
|
{
|
|
//--- R13-B1: отказ durable-записи финализирует сессию с
|
|
//--- сохранением контекста (FINISHED + termination_reason);
|
|
//--- true направляет вызывающий код в экспорт частичного отчёта.
|
|
if(m_state == LAB_STATE_FINISHED)
|
|
return(true);
|
|
return(false);
|
|
}
|
|
m_state = LAB_STATE_FINISHED;
|
|
return(true);
|
|
}
|
|
if(m_settings.experiment_id == LAB_EXP_E4)
|
|
{
|
|
error.code = 15;
|
|
error.message = "E4 reuse-analysis requires LOCAL_ONLY";
|
|
return(false);
|
|
}
|
|
m_operator_started = true;
|
|
m_session_id = NewSessionId();
|
|
if(m_session_id == "")
|
|
{
|
|
m_operator_started = false;
|
|
error.code = 46;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "cannot allocate unique session id";
|
|
return(false);
|
|
}
|
|
//--- R4-S5: отказы инициализации торговых буферов не игнорируются
|
|
if(!m_journal.Init(m_settings.event_capacity, m_session_id))
|
|
{
|
|
m_operator_started = false;
|
|
error.code = 42;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "journal init failed at start";
|
|
return(false);
|
|
}
|
|
if(!m_tracker.Init(m_settings.sample_capacity, m_session_id))
|
|
{
|
|
m_operator_started = false;
|
|
error.code = 43;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "tracker init failed at start";
|
|
return(false);
|
|
}
|
|
if(!m_runner.Init(m_settings, error))
|
|
{
|
|
m_operator_started = false;
|
|
return(false);
|
|
}
|
|
m_termination_reason = "";
|
|
ResetSessionWindows(); // R13-B3: вместо набора отдельных сбросов
|
|
m_cleanup_req_count = 0;
|
|
m_checkpoint_ok = true; // R6-B6
|
|
m_durable_ok = true; // R6-B6
|
|
m_last_quote_observed_us = 0; // R7-B5
|
|
m_last_quote_msc = 0; // R7-B5
|
|
m_pending_sr_count = 0; // R8-B3
|
|
m_state = LAB_STATE_WARMUP;
|
|
return(true);
|
|
}
|
|
//--- формальный статус завершения (аудит-2 S5: INSUFFICIENT_DATA)
|
|
void CAppController::TerminateRun(const string reason)
|
|
{
|
|
if(m_journal != NULL)
|
|
{
|
|
LabEvent e;
|
|
e.Zero();
|
|
e.kind = LAB_KIND_RUN_END;
|
|
e.local_time_us = m_clock.NowUs();
|
|
e.payload_json = StringFormat("{\"stop_reason\":\"%s\"}", reason);
|
|
ulong eid = 0;
|
|
m_journal.Add(e, eid);
|
|
}
|
|
m_termination_reason = reason;
|
|
m_operator_started = false;
|
|
m_state = LAB_STATE_FINISHED;
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
string CAppController::TerminationReason() const { return(m_termination_reason); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::CleanupInFlight() const { return(m_cleanup_in_flight); } // R5-B1/B2
|
|
string CAppController::ConfigurationId() const { return(m_settings.configuration_id); } // R5-S2
|
|
void CAppController::MarkFinished() { m_operator_started = false; m_state = LAB_STATE_FINISHED; } // R5-B2
|
|
//--- R10-B1: повторный запуск оператора после финализации - контроллер
|
|
//--- переводится из FINISHED в READY (новую сессию создаёт StartRun).
|
|
//--- R13-B3: единая очистка сессионного оконного контекста E3 после
|
|
//--- выгрузки прежней сессии — старые окна НЕ переносятся в новый CSV
|
|
//--- (иначе BuildMarketWindows пометил бы их новой session_id), счётчик
|
|
//--- окон и недописанные оконные ссылки обнуляются.
|
|
void CAppController::ResetSessionWindows(void)
|
|
{
|
|
m_window_count = 0;
|
|
if(ArraySize(m_windows) > 0)
|
|
ArrayResize(m_windows, 0);
|
|
m_e3_start_us = 0;
|
|
m_e3_window_seq = 0;
|
|
m_last_window_seq = 0;
|
|
m_last_anchor_msc = 0;
|
|
m_window_ticks_seq = 0;
|
|
m_window_ticks_pending = false;
|
|
if(ArraySize(m_window_ticks) > 0)
|
|
ArrayResize(m_window_ticks, 0);
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::ResetForRestart(void)
|
|
{
|
|
if(m_state != LAB_STATE_FINISHED)
|
|
return(false);
|
|
m_session_id = "";
|
|
m_termination_reason = "";
|
|
m_operator_started = false;
|
|
ResetSessionWindows(); // R13-B3: окна прежней сессии не переносятся
|
|
m_state = LAB_STATE_READY;
|
|
return(true);
|
|
}
|
|
//--- R13-S2/R14-S1: открыто ли измерительное окно. Это активная
|
|
//--- последовательность ЛИБО незавершённый служебный cleanup (тоже
|
|
//--- рабочая отправка в открытом окне): нулевая основная sequence после
|
|
//--- EndMainCycle сама по себе окно не закрывает. Панель не должна
|
|
//--- принудительно перерисовывать график внутри такого окна.
|
|
bool CAppController::InMeasurementWindow() const
|
|
{
|
|
return(m_active_sequence != 0 || (m_cleanup_in_flight && m_cleanup_sequence != 0));
|
|
}
|
|
//--- R13-B1/B3: служебные слоты ТОЛЬКО для регрессий — инъекция отказа
|
|
//--- durable-записи SyntheticDispatch и имитация сохранённого окна E3.
|
|
void CAppController::TestSetFailIntentAfter(const int n) { m_test_fail_intent_counter = (n > 0 ? n : 0); }
|
|
//--- R14-S1: ТОЛЬКО для регрессий — имитация открытого окна cleanup.
|
|
void CAppController::TestSetCleanupWindow(const bool open)
|
|
{
|
|
if(open)
|
|
{
|
|
m_cleanup_in_flight = true;
|
|
m_cleanup_sequence = 1;
|
|
}
|
|
else
|
|
{
|
|
m_cleanup_in_flight = false;
|
|
m_cleanup_sequence = 0;
|
|
}
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
void CAppController::TestAddMarketWindow(const MarketWindow &w)
|
|
{
|
|
if(m_window_count >= 4096)
|
|
return;
|
|
if(m_window_count >= ArraySize(m_windows) &&
|
|
ArrayResize(m_windows, m_window_count + 1) != m_window_count + 1)
|
|
return;
|
|
MarketWindow c = w;
|
|
c.session_id = m_session_id;
|
|
m_windows[m_window_count] = c;
|
|
m_window_count++;
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
void CAppController::StopRun(void)
|
|
{
|
|
if(m_state == LAB_STATE_MAIN || m_state == LAB_STATE_WARMUP)
|
|
{
|
|
m_state = LAB_STATE_DRAINING;
|
|
}
|
|
m_operator_started = false;
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::IsRunning() const
|
|
{
|
|
return (m_state == LAB_STATE_WARMUP || m_state == LAB_STATE_MAIN ||
|
|
m_state == LAB_STATE_DRAINING || m_state == LAB_STATE_EXPORTING);
|
|
}
|
|
//--- R6-B6: безопасная фаза для тяжёлого контрольного сохранения —
|
|
//--- нет открытого окна измерения и незавершённой очистки.
|
|
bool CAppController::CanCheckpoint() const
|
|
{
|
|
return(m_active_sequence == 0 && !m_cleanup_in_flight &&
|
|
(m_state == LAB_STATE_WARMUP || m_state == LAB_STATE_MAIN ||
|
|
m_state == LAB_STATE_DRAINING || m_state == LAB_STATE_FINISHED));
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::CheckpointOk() const { return(m_checkpoint_ok); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
void CAppController::SetCheckpointState(const bool ok) { m_checkpoint_ok = ok; }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
void CAppController::BlockRun(const string reason)
|
|
{
|
|
if(m_state != LAB_STATE_FINISHED)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = reason;
|
|
}
|
|
}
|
|
//--- R6-S2/R7-B5: фиксация факта НОВОГО наблюдения рабочего символа.
|
|
//--- Считается только монотонное увеличение серверной метки тика
|
|
//--- (реально новое обновление), а не повторное чтение прежней
|
|
//--- котировки; иначе застой рынка маскируется регулярным опросом.
|
|
void CAppController::NoteQuoteObserved(const long tick_msc)
|
|
{
|
|
if(tick_msc > m_last_quote_msc)
|
|
{
|
|
m_last_quote_msc = tick_msc;
|
|
m_last_quote_observed_us = m_clock.NowUs();
|
|
}
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
ulong CAppController::LastQuoteObservedUs() const { return(m_last_quote_observed_us); }
|
|
//--- R6-B6/R7-S3/R8-B4: чтение неопределённых попыток из durable-журнала
|
|
//--- сессии (INTENT без SEND_RETURN). Торговое состояние не меняется —
|
|
//--- это учёт качества для manifest/checks после сбоев и перезапусков.
|
|
//--- Семантика результата:
|
|
//--- 0 — журнала нет/пуст: неопределённых отправок нет;
|
|
//--- ULONG_MAX — журнал ЕСТЬ, но недоступен/частично прочитан
|
|
//--- (неопределённость НЕ сужается до нуля, sentinel);
|
|
//--- иначе — число INTENT-строк без парного SEND_RETURN.
|
|
//--- Подтверждённые ключи (sequence) набираются ДИНАМИЧЕСКИ, без
|
|
//--- фиксированного предела 256 (штатный парный запуск с cleanup на
|
|
//--- 5x100 превышает его — иначе получаются ложные неопределённости).
|
|
ulong CAppController::CountUncertainDurableIntents(const string session_id)
|
|
{
|
|
m_durable_uncertain = 0;
|
|
const string path = StringFormat("RequestLatencyLab\\%s\\%s\\intents.csv",
|
|
m_settings.campaign_id, session_id);
|
|
const int h = FileOpen(path, FILE_READ | FILE_BIN);
|
|
if(h == INVALID_HANDLE)
|
|
{
|
|
//--- R8-B4: отсутствие файла НЕ равно недоступности существующего
|
|
if(FileIsExist(path))
|
|
return(ULONG_MAX);
|
|
return(0); // журнала нет — неопределённых отправок нет
|
|
}
|
|
const int total = (int)FileSize(h);
|
|
if(total <= 0)
|
|
{
|
|
FileClose(h);
|
|
return(0);
|
|
}
|
|
uchar b[];
|
|
if(ArrayResize(b, total) != total)
|
|
{
|
|
FileClose(h);
|
|
return(ULONG_MAX);
|
|
}
|
|
const uint got = FileReadArray(h, b, 0, total);
|
|
FileClose(h);
|
|
if(got != (uint)total)
|
|
return(ULONG_MAX); // короткое чтение — нецелый журнал
|
|
string content = CharArrayToString(b, 0, total, CP_UTF8);
|
|
string lines[];
|
|
const int n = StringSplit(content, '\n', lines);
|
|
ulong ret_seq[]; // подтверждённые SEND_RETURN (sequence)
|
|
int rc = 0;
|
|
ulong intent_seq[]; // все INTENT-строки
|
|
ulong cancel_seq[]; // R9-S2: отменённые INTENT (до T0)
|
|
bool cancel_used[]; // R10-B2: одна отмена на одну попытку
|
|
int cc = 0;
|
|
int ic = 0;
|
|
for(int i = 0; i < n; i++)
|
|
{
|
|
string s = lines[i];
|
|
StringReplace(s, "\r", "");
|
|
if(StringLen(s) == 0)
|
|
continue;
|
|
string cols[];
|
|
const int c = StringSplit(s, ',', cols);
|
|
if(c < 4)
|
|
{
|
|
//--- R7-S3/R8-B4: повреждённая строка — неопределённый факт отправки
|
|
m_durable_uncertain++;
|
|
continue;
|
|
}
|
|
const ulong seq = (ulong)StringToInteger(cols[1]);
|
|
if(cols[2] == "SEND_RETURN")
|
|
{
|
|
if(rc >= ArraySize(ret_seq) &&
|
|
ArrayResize(ret_seq, rc + 16) != (rc + 16))
|
|
return(ULONG_MAX); // ёмкость не подтверждена — sentinel
|
|
ret_seq[rc] = seq;
|
|
rc++;
|
|
}
|
|
else
|
|
if(cols[2] == "INTENT")
|
|
{
|
|
if(ic >= ArraySize(intent_seq) &&
|
|
ArrayResize(intent_seq, ic + 16) != (ic + 16))
|
|
return(ULONG_MAX);
|
|
intent_seq[ic] = seq;
|
|
ic++;
|
|
}
|
|
else
|
|
if(cols[2] == "INTENT_CANCELLED")
|
|
{
|
|
//--- R9-S2: явная отмена неотправленного намерения (pre-T0
|
|
//--- отказ после записи durable-INTENT): INTENT с парной
|
|
//--- отменой НЕ считается неопределённой отправкой.
|
|
//--- R10-B2: одна INTENT_CANCELLED снимает неопределённость
|
|
//--- РОВНО с одной попытки того же seq (хронологическая пара).
|
|
if(cc >= ArraySize(cancel_seq) &&
|
|
ArrayResize(cancel_seq, cc + 16) != (cc + 16))
|
|
return(ULONG_MAX);
|
|
if(cc >= ArraySize(cancel_used) &&
|
|
ArrayResize(cancel_used, cc + 16) != (cc + 16))
|
|
return(ULONG_MAX);
|
|
cancel_seq[cc] = seq;
|
|
cancel_used[cc] = false;
|
|
cc++;
|
|
}
|
|
}
|
|
for(int k = 0; k < ic; k++)
|
|
{
|
|
const ulong seq = intent_seq[k];
|
|
bool has = false;
|
|
for(int r = 0; r < rc && !has; r++)
|
|
if(ret_seq[r] == seq)
|
|
has = true;
|
|
if(!has)
|
|
{
|
|
//--- R10-B2: INTENT_CANCELLED снимает неопределённость РОВНО с
|
|
//--- одной попытки того же seq (хронологическая пара): повторный
|
|
//--- INTENT после отмены без SEND_RETURN остаётся неопределённым.
|
|
bool has_cancel = false;
|
|
for(int c2 = 0; c2 < cc && !has_cancel; c2++)
|
|
if(cancel_seq[c2] == seq && !cancel_used[c2])
|
|
{
|
|
cancel_used[c2] = true;
|
|
has_cancel = true;
|
|
}
|
|
if(!has_cancel)
|
|
m_durable_uncertain++;
|
|
}
|
|
}
|
|
return(m_durable_uncertain);
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::LoadDurableContent(const string path, string &content)
|
|
{
|
|
content = "";
|
|
const int h = FileOpen(path, FILE_READ | FILE_BIN);
|
|
if(h == INVALID_HANDLE)
|
|
return(false);
|
|
const int total = (int)FileSize(h);
|
|
if(total > 0)
|
|
{
|
|
uchar b[];
|
|
if(ArrayResize(b, total) != total)
|
|
{
|
|
FileClose(h);
|
|
return(false);
|
|
}
|
|
const uint got = FileReadArray(h, b, 0, total);
|
|
FileClose(h);
|
|
if(got != (uint)total)
|
|
return(false); // R8-B4: короткое чтение — отказ
|
|
content = CharArrayToString(b, 0, total, CP_UTF8);
|
|
}
|
|
else
|
|
FileClose(h);
|
|
return(true);
|
|
}
|
|
//--- обработка OnTradeTransaction (единственная торговая связь)
|
|
void CAppController::OnTransaction(TransactionObservation &obs, LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(m_journal == NULL || m_tracker == NULL)
|
|
return;
|
|
//--- корреляция до журнала: origin без догадок (аудит-2 B2)
|
|
const bool is_request = (obs.trans.type == TRADE_TRANSACTION_REQUEST);
|
|
const ulong seq = m_tracker.ResolveEventSequence(is_request, obs.result.request_id,
|
|
obs.trans.order, obs.trans.deal);
|
|
obs.target_sequence = seq;
|
|
//--- неизвестный REQUEST сохраняется как unresolved; LatencySample не меняется
|
|
if(seq == 0 && is_request && obs.result.request_id != 0 &&
|
|
!IsCleanupRequest(obs.result.request_id))
|
|
m_tracker.PushUnresolved(obs.result.request_id);
|
|
//--- B2 (audit-3): ранний ORDER_ADD/DEAL_ADD без точной привязки
|
|
//--- сохраняется с исходной меткой для последующей корреляции.
|
|
if(seq == 0 && !is_request &&
|
|
(obs.trans.type == TRADE_TRANSACTION_ORDER_ADD ||
|
|
obs.trans.type == TRADE_TRANSACTION_DEAL_ADD))
|
|
m_tracker.PushDeferred(obs.entry_us, (long)obs.trans.type,
|
|
obs.trans.order, obs.trans.deal,
|
|
obs.trans.position, obs.trans.volume);
|
|
//--- journal: типизированная запись (вход, затем выход)
|
|
LabEvent ev;
|
|
ev.Zero();
|
|
ev.kind = LAB_KIND_TRADE_TRANSACTION;
|
|
ev.local_time_us = obs.entry_us;
|
|
ev.origin_sequence = seq; // 0 = без привязки к записи
|
|
ev.transaction_type = EnumToString(obs.trans.type);
|
|
ev.handler_enter_us = obs.entry_us;
|
|
ev.order_ticket = obs.trans.order;
|
|
ev.deal_ticket = obs.trans.deal;
|
|
ev.position_ticket = obs.trans.position;
|
|
ev.request_id = obs.result.request_id;
|
|
ev.payload_json = StringFormat("{\"type\":\"%s\"}", EnumToString(obs.trans.type));
|
|
ev.quality_code = "OK";
|
|
ulong eseq = 0;
|
|
m_journal.Add(ev, eseq);
|
|
obs.trade_event_id = eseq;
|
|
if(seq == 0)
|
|
return;
|
|
RequestMetadata vr;
|
|
if(!m_tracker.GetMetadata(seq, vr))
|
|
return;
|
|
const bool verbose = (vr.plan.logging_mode == LAB_LOG_VERBOSE);
|
|
const bool is_request_event = (obs.trans.type == TRADE_TRANSACTION_REQUEST);
|
|
//--- R4-B3: поздний REQUEST сохраняет исходную метку (T2), даже если
|
|
//--- торговый исход уже известен и коллекция закрыта: REQUEST является
|
|
//--- измерительной точкой и не должен теряться после T6.
|
|
if(is_request_event)
|
|
{
|
|
m_tracker.MarkRequest(seq, obs.entry_us, obs.result.request_id,
|
|
obs.result.retcode, obs.result.retcode_external,
|
|
obs.result.order, obs.result.deal);
|
|
if(verbose)
|
|
PrintFormat("[LAB][%s] callback REQUEST seq=%I64u request_id=%u retcode=%u",
|
|
m_session_id, seq, obs.result.request_id, obs.result.retcode);
|
|
}
|
|
//--- аудит-2 B3: коллекция закрыта — события хвоста (кроме REQUEST) не меняют запись
|
|
if(vr.collection_closed && !is_request_event)
|
|
return;
|
|
if(!is_request_event)
|
|
{
|
|
switch(obs.trans.type)
|
|
{
|
|
case TRADE_TRANSACTION_ORDER_ADD:
|
|
{
|
|
const ulong ticket = obs.trans.order;
|
|
RequestMetadata r;
|
|
if(m_tracker.GetMetadata(seq, r))
|
|
{
|
|
if(r.order_ticket == 0 || r.order_ticket == ticket)
|
|
m_tracker.MarkOrderAdd(seq, obs.entry_us, ticket, "TRANS_ORDER");
|
|
}
|
|
if(verbose)
|
|
PrintFormat("[LAB][%s] callback ORDER_ADD seq=%I64u order=%I64u",
|
|
m_session_id, seq, ticket);
|
|
break;
|
|
}
|
|
case TRADE_TRANSACTION_DEAL_ADD:
|
|
{
|
|
const bool deal_ok = m_tracker.MarkDealAdd(seq, obs.entry_us,
|
|
obs.trans.deal,
|
|
obs.trans.volume > 0.0 ? obs.trans.volume : 0.0);
|
|
//--- R7-S4: потеря первичного факта сделки (переполнение
|
|
//--- реестра) — fail-closed сразу, а не только RUN-DEAL-02
|
|
//--- в конце серии: дальнейшие отправки запрещены.
|
|
if(!deal_ok && m_tracker.DealLostCount() > 0 &&
|
|
m_state != LAB_STATE_BLOCKED)
|
|
{
|
|
m_state = LAB_STATE_BLOCKED;
|
|
m_block_reason = "deal registry overflow at callback";
|
|
}
|
|
if(obs.trans.position != 0)
|
|
m_tracker.SetPositionTicket(seq, obs.trans.position);
|
|
if(verbose)
|
|
PrintFormat("[LAB][%s] callback DEAL_ADD seq=%I64u deal=%I64u",
|
|
m_session_id, seq, obs.trans.deal);
|
|
break;
|
|
}
|
|
case TRADE_TRANSACTION_HISTORY_ADD:
|
|
case TRADE_TRANSACTION_ORDER_DELETE:
|
|
//--- событийная сверка выполняется общей точкой ниже (B4)
|
|
break;
|
|
default:
|
|
break;
|
|
}
|
|
}
|
|
//--- B2: после появления точного моста (REQUEST/result.order) повторно
|
|
//--- связываем буферизованные ранние ORDER_ADD/DEAL_ADD исходными метками
|
|
//--- R4-B8: каждая первая привязка фиксируется как CORRELATION_LINK.
|
|
const int bound = m_tracker.ReconcileDeferred(seq);
|
|
if(bound > 0)
|
|
{
|
|
LabEvent ce;
|
|
ce.Zero();
|
|
ce.kind = LAB_KIND_CORRELATION_LINK;
|
|
ce.local_time_us = m_clock.NowUs();
|
|
ce.origin_sequence = seq;
|
|
ce.payload_json = StringFormat("{\"bound\":%d,\"via\":\"order_ticket\"}", bound);
|
|
ulong ceid = 0;
|
|
m_journal.Add(ce, ceid);
|
|
}
|
|
//--- аудит-2 B4 + B5: событийная адресная сверка (light, без скан истории)
|
|
const bool tracked_main = (vr.present_mask & (uint)LAB_MASK_T0) != 0;
|
|
if(tracked_main && !vr.collection_closed)
|
|
CheckTriggered(seq, error);
|
|
}
|
|
//--- таймер: сроки, сверка (watchdog), разрешённые шаги; событийная
|
|
//--- сверка выполняется в OnTransaction (аудит-2 B4)
|
|
bool CAppController::OnTimer(const ulong timer_last_call_us, LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(m_state != LAB_STATE_MAIN && m_state != LAB_STATE_WARMUP &&
|
|
m_state != LAB_STATE_DRAINING)
|
|
return(false);
|
|
//--- B1: в DRAINING без активной основной операции может идти cleanup
|
|
if(m_state == LAB_STATE_DRAINING && m_active_sequence == 0 && m_cleanup_in_flight)
|
|
{
|
|
const bool done = DriveCleanup(error);
|
|
return(done); // true = возврат к MAIN (следующий шаг в DriveNext)
|
|
}
|
|
if(m_active_sequence == 0)
|
|
return(false);
|
|
RequestMetadata r;
|
|
if(!m_tracker.GetMetadata(m_active_sequence, r))
|
|
return(false);
|
|
const ulong now = m_clock.NowUs();
|
|
//--- запись фактического интервала таймера
|
|
m_tracker.SetTimerInterval(m_active_sequence, timer_last_call_us);
|
|
//--- deadline exceeded: нет T6 к основному сроку
|
|
if(!r.deadline_exceeded && r.observation_deadline_us > 0 &&
|
|
now >= r.observation_deadline_us &&
|
|
(r.present_mask & (uint)LAB_MASK_T6) == 0)
|
|
m_tracker.SetDeadlineExceeded(m_active_sequence);
|
|
//--- аудит-2 B3: истёк late-grace без T6 — коллекция закрывается,
|
|
//--- исход остаётся неопределённым (completed=false), переход к очистке
|
|
if(!r.collection_closed && r.collection_deadline_us > 0 &&
|
|
now >= r.collection_deadline_us &&
|
|
(r.present_mask & (uint)LAB_MASK_T6) == 0)
|
|
{
|
|
ulong cid = 0;
|
|
LabEvent ce;
|
|
ce.Zero();
|
|
ce.kind = LAB_KIND_COLLECTION_CLOSE;
|
|
ce.local_time_us = now;
|
|
ce.origin_sequence = m_active_sequence;
|
|
ce.payload_json = "{\"reason\":\"grace_expired_no_t6\"}";
|
|
m_journal.Add(ce, cid);
|
|
m_tracker.CloseCollection(m_active_sequence, now, cid);
|
|
if(!ConfirmDurableSendReturns())
|
|
return(false); // R9-B3: сбой flush - BLOCKED
|
|
m_state = LAB_STATE_DRAINING;
|
|
return(true);
|
|
}
|
|
//--- сверка по таймеру (резерв)
|
|
const bool have_t6 = (r.present_mask & (uint)LAB_MASK_T6) != 0;
|
|
if(!have_t6)
|
|
{
|
|
ConfirmationSnapshot snap;
|
|
//--- полный reconcile в таймере (не в callback, B5)
|
|
CStateReader::BuildSnapshot(m_active_sequence, r.plan.operation, r.order_ticket,
|
|
r.plan.volume / m_rules.volume_step,
|
|
m_rules.volume_step, m_settings.volume_units_tolerance,
|
|
m_clock, r.position_ticket,
|
|
m_settings.history_scan_budget, snap);
|
|
snap.requested_volume_units = r.plan.volume / m_rules.volume_step;
|
|
snap.order_state_after = snap.order_state_before;
|
|
//--- R5-B4: сделки из истории — в реестр (deals.csv) с происхождением
|
|
SyncRecoveredDeals(m_active_sequence, snap);
|
|
//--- R4-B8: полная сверка синхронизирует объёмы записи (как синтетика)
|
|
m_tracker.SyncReconcile(m_active_sequence, snap, m_rules.volume_step,
|
|
m_settings.volume_units_tolerance);
|
|
if(snap.position_ticket != 0 && r.position_ticket == 0)
|
|
m_tracker.SetPositionTicket(m_active_sequence, snap.position_ticket);
|
|
CompletedCheck(m_active_sequence, snap, error);
|
|
}
|
|
else
|
|
if(!r.callback_coverage_complete)
|
|
{
|
|
//--- R6-B1: даже после T6 полная сверка переносит callback-покрытие
|
|
//--- (T5) в запись; T6/исход не переписываются.
|
|
ConfirmationSnapshot snap;
|
|
CStateReader::BuildSnapshot(m_active_sequence, r.plan.operation, r.order_ticket,
|
|
r.plan.volume / m_rules.volume_step,
|
|
m_rules.volume_step, m_settings.volume_units_tolerance,
|
|
m_clock, r.position_ticket,
|
|
m_settings.history_scan_budget, snap);
|
|
snap.requested_volume_units = r.plan.volume / m_rules.volume_step;
|
|
snap.order_state_after = snap.order_state_before;
|
|
SyncRecoveredDeals(m_active_sequence, snap);
|
|
m_tracker.SyncReconcile(m_active_sequence, snap, m_rules.volume_step,
|
|
m_settings.volume_units_tolerance);
|
|
//--- R7-B2: полнота callback-меток доказана полным сканом —
|
|
//--- коллекцию закрываем сразу (не ждём late-grace); T6/исход
|
|
//--- не переписываются, запоздавшие метки уже внесены.
|
|
RequestMetadata r3;
|
|
//--- R8-B2: ЕДИНАЯ политика полноты измерительных точек (T2/T3 по
|
|
//--- операции) применяется и в пост-T6 ветви OnTimer: одинаковое
|
|
//--- событие «покрытие теперь полное» закрывает коллекцию только
|
|
//--- при собранных метках; иначе коллекцию закрывает grace.
|
|
if(m_tracker.GetMetadata(m_active_sequence, r3) &&
|
|
r3.callback_coverage_complete && !r3.collection_closed &&
|
|
MeasurementPointsComplete(r3))
|
|
{
|
|
ulong cid3 = 0;
|
|
LabEvent ce3;
|
|
ce3.Zero();
|
|
ce3.kind = LAB_KIND_COLLECTION_CLOSE;
|
|
ce3.local_time_us = m_clock.NowUs();
|
|
ce3.origin_sequence = m_active_sequence;
|
|
ce3.payload_json = "{\"reason\":\"callback_coverage_complete\"}";
|
|
m_journal.Add(ce3, cid3);
|
|
m_tracker.CloseCollection(m_active_sequence, m_clock.NowUs(), cid3);
|
|
if(!ConfirmDurableSendReturns())
|
|
return(false); // R9-B3: сбой flush - BLOCKED
|
|
m_state = LAB_STATE_DRAINING;
|
|
}
|
|
}
|
|
return(true);
|
|
}
|
|
//--- проверка завершения (адресная для record_sequence)
|
|
bool CAppController::CompletedCheck(const ulong record_sequence, ConfirmationSnapshot &snap, LabError &error)
|
|
{
|
|
error.Reset();
|
|
snap.sequence = record_sequence;
|
|
if(snap.check_end_us == 0)
|
|
snap.check_end_us = m_clock.NowUs();
|
|
snap.trigger = (snap.trigger == LAB_FINAL_NONE ? LAB_FINAL_TIMER_CHECK : snap.trigger);
|
|
//--- B7: reject-retcode из записи (sync-return/REQUEST) для однозначного отказа
|
|
RequestMetadata crm;
|
|
bool have_crm = false;
|
|
if(m_tracker != NULL && m_tracker.GetMetadata(record_sequence, crm))
|
|
{
|
|
have_crm = true;
|
|
uint rr = 0;
|
|
if(crm.has_retcode)
|
|
rr = crm.sample.retcode;
|
|
else
|
|
if(crm.send_retcode != 0)
|
|
rr = crm.send_retcode;
|
|
if(rr != 0 && snap.reject_retcode == 0)
|
|
snap.reject_retcode = rr;
|
|
//--- R4-S2: ожидаемый контракт запроса для сверки параметров
|
|
//--- (pending-create). Допуск цены - половина тика символа.
|
|
snap.expected_symbol = crm.plan.symbol;
|
|
snap.expected_magic = crm.plan.magic;
|
|
snap.expected_volume_lots = crm.plan.volume;
|
|
snap.expected_price = crm.plan.request_price;
|
|
snap.price_tolerance = 0.5 * m_rules.tick_size;
|
|
snap.expected_order_type = crm.plan.actual_type; // R5-B8
|
|
}
|
|
CompletionDecision d;
|
|
d.Zero();
|
|
CCompletionPolicy::Evaluate(snap.expected_operation, snap,
|
|
m_settings.volume_units_tolerance, d);
|
|
//--- R6-B5: доказанный исход без права установить T6 (например,
|
|
//--- UNEXPECTED_ACTIVATION) сохраняется отдельно; T6 не ставится.
|
|
if(d.confirmed && !d.t6_assignable)
|
|
{
|
|
ulong eid2 = 0;
|
|
LabEvent e2;
|
|
e2.Zero();
|
|
e2.kind = LAB_KIND_CONFIRMATION_SNAPSHOT;
|
|
e2.local_time_us = snap.check_end_us;
|
|
e2.origin_sequence = record_sequence;
|
|
e2.payload_json = StringFormat(
|
|
"{\"outcome\":\"%s\",\"t6_assignable\":false,\"reason\":\"%s\","
|
|
"\"order_ticket\":%I64u}", LabOutcomeName(d.outcome), d.reason_code,
|
|
snap.order_ticket);
|
|
m_journal.Add(e2, eid2);
|
|
m_tracker.RecordOutcome(record_sequence, d.outcome);
|
|
//--- известное торговое состояние переносится независимо от T6;
|
|
//--- final rejection без ордера = доказанное отсутствие экспозиции.
|
|
const bool ts_known2 = (snap.trading_state_known || d.outcome == LAB_OUT_REJECTED);
|
|
m_tracker.SetTradingStateKnown(record_sequence, ts_known2);
|
|
return(false);
|
|
}
|
|
//--- T6 = check_end_us первой успешной проверки
|
|
if(d.confirmed && d.t6_assignable)
|
|
{
|
|
ulong event_id = 0;
|
|
LabEvent e;
|
|
e.Zero();
|
|
e.kind = LAB_KIND_CONFIRMATION_SNAPSHOT;
|
|
e.local_time_us = snap.check_end_us;
|
|
e.origin_sequence = record_sequence;
|
|
//--- R4-B8: T6 evidence — проверяемое основание решения
|
|
e.payload_json = StringFormat(
|
|
"{\"outcome\":\"%s\",\"check_end_us\":%I64u,\"reason\":\"%s\","
|
|
"\"order_ticket\":%I64u,\"order_state\":%d,\"requested_units\":%.10g,"
|
|
"\"executed_units\":%.10g,\"remaining_units\":%.10g,"
|
|
"\"active_status\":\"%s\",\"history_scan_failed\":%s,"
|
|
"\"coverage_complete\":%s,\"deal_count\":%d,\"position_ticket\":%I64u,"
|
|
"\"trigger\":\"%s\"}",
|
|
LabOutcomeName(d.outcome), snap.check_end_us, d.reason_code,
|
|
snap.order_ticket, snap.order_state_before,
|
|
snap.requested_volume_units, snap.executed_volume_units,
|
|
snap.remaining_volume_units,
|
|
LabReadStatusName((ENUM_LAB_READ_STATUS)snap.active_select_status),
|
|
(snap.history_scan_failed ? "true" : "false"),
|
|
(snap.coverage_complete ? "true" : "false"),
|
|
snap.deal_count, snap.position_ticket,
|
|
LabFinalSourceName(snap.trigger));
|
|
m_journal.Add(e, event_id);
|
|
m_tracker.SetFinal(record_sequence, snap.check_end_us, d.outcome,
|
|
snap.trigger, event_id);
|
|
//--- R6-B5: определённый REJECT без ордера также доказывает
|
|
//--- отсутствие остаточной экспозиции (ордер не был создан).
|
|
const bool ts_known = (snap.trading_state_known || d.outcome == LAB_OUT_REJECTED);
|
|
m_tracker.SetTradingStateKnown(record_sequence, ts_known);
|
|
//--- R5-B4: результаты проверки переносятся в metadata (executed_volume,
|
|
//--- торговое покрытие и callback-покрытие); иначе событийный full-fill
|
|
//--- остаётся с нулевым executed_volume и T5 получает MISSING.
|
|
m_tracker.SyncReconcile(record_sequence, snap, m_rules.volume_step,
|
|
m_settings.volume_units_tolerance);
|
|
//--- R7-B2: для раннего закрытия нужна СВЕЖАЯ полнота callback-меток
|
|
//--- ПОСЛЕ переноса покрытия (SyncReconcile выше).
|
|
RequestMetadata crm2;
|
|
if(!m_tracker.GetMetadata(record_sequence, crm2))
|
|
crm2 = crm;
|
|
//--- аудит-2 B3: T6 и закрытие коллекции разделены. Если торговое
|
|
//--- состояние известно и покрытие сделок подтверждено — закрываем
|
|
//--- сразу; иначе late-grace закроет коллекцию в OnTimer/EndMainCycle.
|
|
const bool deal_expected = (snap.expected_operation == LAB_OP_MARKET_OPEN ||
|
|
snap.expected_operation == LAB_OP_POSITION_CLOSE);
|
|
//--- B3 + R4-B3: раннее закрытие допустимо только при T6 + собранном
|
|
//--- ожидаемом REQUEST (T2) и (для market) подтверждённом покрытии.
|
|
//--- Для ВСЕХ операций, где REQUEST — измерительная точка, закрытие
|
|
//--- обязано ждать T2 либо истечения late-grace (OnTimer/EndMainCycle).
|
|
//--- Поздний REQUEST после T6 не должен терять T2-T4.
|
|
//--- R8-B2: ЕДИНАЯ политика полноты измерительных точек для всех
|
|
//--- путей закрытия (CompletedCheck/OnTimer/cleanup) — REQUEST (T2)
|
|
//--- и ORDER_ADD (T3) по операции; при неполноте закрывает grace.
|
|
const bool close_allowed = (!have_crm ||
|
|
MeasurementPointsComplete(crm2));
|
|
//--- R7-B2: раннее закрытие ТОЛЬКО при подтверждённой полноте
|
|
//--- ожидаемых callback-точек (T4/T5) ИЛИ отсутствии сделок в
|
|
//--- контракте операции. snap.coverage_complete подтверждает объём
|
|
//--- по истории, но не полноту локальных callback-меток; снимок с
|
|
//--- ошибкой/обрезанием скана полноту T5 не доказывает. Недостающие
|
|
//--- метки принимаются до grace (коллекция остаётся открытой).
|
|
const bool deal_callback_expected = (!have_crm ||
|
|
crm.plan.operation == LAB_OP_MARKET_OPEN ||
|
|
crm.plan.operation == LAB_OP_POSITION_CLOSE ||
|
|
crm.plan.operation == LAB_OP_PENDING_CREATE);
|
|
const bool callback_proven = (!deal_callback_expected ||
|
|
crm2.callback_coverage_complete);
|
|
if(snap.trading_state_known && close_allowed && callback_proven &&
|
|
(!deal_expected || snap.coverage_complete))
|
|
{
|
|
ulong cid = 0;
|
|
LabEvent ce;
|
|
ce.Zero();
|
|
ce.kind = LAB_KIND_COLLECTION_CLOSE;
|
|
ce.local_time_us = snap.check_end_us;
|
|
ce.origin_sequence = record_sequence;
|
|
ce.payload_json = "{\"reason\":\"state_known\"}";
|
|
m_journal.Add(ce, cid);
|
|
m_tracker.CloseCollection(record_sequence, snap.check_end_us, cid);
|
|
if(!ConfirmDurableSendReturns())
|
|
return(false); // R9-B3: сбой flush - BLOCKED
|
|
}
|
|
//--- состояние DRAINING: cleanup разрешён только после collection_closed
|
|
m_state = LAB_STATE_DRAINING;
|
|
return(true);
|
|
}
|
|
return(false);
|
|
}
|
|
//--- E3: загрузка калибровки (формируется MarketCalibration.mq5)
|
|
bool CAppController::SetCalibration(const Calibration &cal)
|
|
{
|
|
//--- S5: калибровка должна соответствовать символу эксперимента
|
|
if(cal.valid && cal.calibration_id != "" &&
|
|
StringLen(cal.symbol) > 0 && m_settings.symbol != cal.symbol)
|
|
{
|
|
m_calibration_loaded = false;
|
|
return(false);
|
|
}
|
|
m_calibration = cal;
|
|
m_calibration_loaded = (cal.valid && cal.calibration_id != "");
|
|
return(m_calibration_loaded);
|
|
}
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
bool CAppController::CalibrationLoaded() const { return(m_calibration_loaded); }
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
string CAppController::CalibrationId() const { return(m_calibration.calibration_id); }
|
|
//--- E3: вычислить окно рынка и классифицировать режим
|
|
void CAppController::UpdateMarketRegime(ENUM_LAB_MARKET_REGIME ®ime)
|
|
{
|
|
regime = LAB_REGIME_UNKNOWN;
|
|
if(!m_calibration_loaded)
|
|
return;
|
|
//--- R4-B5: опорная котировка (time_msc) снимается ПЕРЕД классификацией;
|
|
//--- окно строится относительно реально использованной котировки,
|
|
//--- а не TimeCurrent() (нарушало причинный контракт окна).
|
|
MqlTick quote;
|
|
if(!SymbolInfoTick(m_settings.symbol, quote) || quote.time_msc <= 0)
|
|
return;
|
|
//--- R6-S2/R7-B5/R8-S2: возраст последнего НОВОГО обновления рабочего
|
|
//--- символа. СНАЧАЛА фиксируется факт реально нового тика (монотон-
|
|
//--- ный quote.time_msc [R7-B5]) — он освежает отметку наблюдения;
|
|
//--- ТОЛЬКО затем проверяется возраст. Так polling-путь (SymbolInfoTick
|
|
//--- в UpdateMarketRegime) восстанавливает контекст после паузы, а
|
|
//--- повторное чтение прежней котировки застой не маскирует.
|
|
{
|
|
const ulong obs_us = m_clock.NowUs();
|
|
if(quote.time_msc > m_last_quote_msc)
|
|
{
|
|
m_last_quote_msc = quote.time_msc;
|
|
m_last_quote_observed_us = obs_us;
|
|
}
|
|
if(m_settings.feature_max_age_ms > 0 && m_last_quote_observed_us > 0 &&
|
|
obs_us > m_last_quote_observed_us &&
|
|
obs_us - m_last_quote_observed_us > (ulong)m_settings.feature_max_age_ms * 1000)
|
|
return; // застой котировок: режим не определяется
|
|
}
|
|
const long anchor_msc = quote.time_msc;
|
|
//--- R5-B6: возраст признаков — на ОДНОЙ календарной шкале с якорем
|
|
//--- окна (quote.time_msc и TimeCurrent() — серверные метки времени;
|
|
//--- GetMicrosecondCount не смешивается с календарной шкалой).
|
|
const ulong now_us = m_clock.NowUs();
|
|
const long server_now_msc = (long)TimeCurrent() * 1000;
|
|
const ulong age_us = (anchor_msc > 0 && server_now_msc > anchor_msc ?
|
|
(ulong)(server_now_msc - anchor_msc) * 1000 : 0);
|
|
if(m_settings.feature_max_age_ms > 0 &&
|
|
age_us > (ulong)m_settings.feature_max_age_ms * 1000)
|
|
return; // котировка устарела: режим не определяется
|
|
MarketWindow w;
|
|
w.Zero();
|
|
LabError werr;
|
|
if(!CMarketRegime::ComputeWindowRawTicks(m_settings.symbol, anchor_msc, now_us, w, werr,
|
|
m_window_ticks))
|
|
return;
|
|
m_last_anchor_msc = anchor_msc; // R9-S2: якорь выбранного окна
|
|
w.age_at_send_us = age_us;
|
|
//--- R5-B6: задержка CopyTicksRange не должна устаревать якорь окна
|
|
if(m_settings.feature_max_age_ms > 0)
|
|
{
|
|
MqlTick q2;
|
|
if(SymbolInfoTick(m_settings.symbol, q2) && q2.time_msc > anchor_msc &&
|
|
(ulong)(q2.time_msc - anchor_msc) * 1000 >
|
|
(ulong)m_settings.feature_max_age_ms * 1000)
|
|
return; // якорь устарел во время вычисления окна
|
|
}
|
|
//--- R4-B5: трассировка window/sequence/calibration/dataset
|
|
m_e3_window_seq++;
|
|
w.sequence = m_e3_window_seq;
|
|
w.calibration_id = m_calibration.calibration_id;
|
|
w.dataset_id = m_calibration.calibration_id;
|
|
w.session_id = m_session_id;
|
|
bool valid = false;
|
|
CMarketRegime::Classify(w, m_calibration, regime, valid);
|
|
if(!valid || (regime != LAB_REGIME_QUIET && regime != LAB_REGIME_FAST))
|
|
{
|
|
regime = LAB_REGIME_UNKNOWN;
|
|
return;
|
|
}
|
|
m_last_regime = regime;
|
|
m_last_window_id = w.window_id;
|
|
m_last_window_seq = m_e3_window_seq; // R8-S2
|
|
m_last_regime_us = m_clock.NowUs();
|
|
//--- R9-S3: сырые тики выбранного окна удерживаются в памяти до
|
|
//--- safe-фазы закрытия sample (запись в файл ПОСЛЕ закрытия окна;
|
|
//--- синхронная запись внутри измерительной фазы запрещена R8-B3).
|
|
//--- Окно, не закреплённое за sample, перезаписывается следующим.
|
|
m_window_ticks_seq = m_e3_window_seq;
|
|
m_window_ticks_pending = true;
|
|
//--- сохранить окно для market_windows.csv (не более 4096)
|
|
w.regime = regime;
|
|
w.session_id = m_session_id;
|
|
if(m_window_count < 4096)
|
|
{
|
|
if(m_window_count >= ArraySize(m_windows) &&
|
|
ArrayResize(m_windows, m_window_count + 1) != m_window_count + 1)
|
|
{
|
|
//--- R8-B1: сбой хранения окна — режим НЕ выводится и торговля
|
|
//--- останавливается (иначе соответствие sample->window теряется
|
|
//--- молча, а отправка в этом же проходе невозможна).
|
|
regime = LAB_REGIME_UNKNOWN;
|
|
m_last_regime = LAB_REGIME_UNKNOWN;
|
|
TerminateRun("E3_WINDOW_STORAGE_FAILED");
|
|
return;
|
|
}
|
|
m_windows[m_window_count] = w;
|
|
m_window_count++;
|
|
}
|
|
else
|
|
{
|
|
//--- R7-S3/R8-B1: исчерпание протокола окон E3 — обязательная
|
|
//--- остановка; выходной regime сбрасывается, чтобы DriveNext не
|
|
//--- продолжил отправку в этом же проходе после записи причины.
|
|
regime = LAB_REGIME_UNKNOWN;
|
|
m_last_regime = LAB_REGIME_UNKNOWN;
|
|
TerminateRun("E3_WINDOW_CAPACITY_EXCEEDED");
|
|
}
|
|
}
|
|
//--- п.11: финализация OnTradeTransaction (длительность, выход журнала, режим)
|
|
void CAppController::FinalizeTransaction(TransactionObservation &obs, LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
if(m_tracker == NULL || m_journal == NULL)
|
|
return;
|
|
//--- аудит-2 B2: без точной привязки запись не трогаем
|
|
const ulong seq = obs.target_sequence;
|
|
//--- R7-S4: ЕДИНОЕ итоговое exit_us пишется в события и metadata ОДИН
|
|
//--- раз. R8-S3: итоговый штамп снимается ДО хвоста финализации —
|
|
//--- линейный поиск события в журнале (SetExitTime) и SetHandlerDuration
|
|
//--- в измеряемую длительность обработчика не входят. Это ЯВНАЯ
|
|
//--- инструментальная граница: длительность = от входа до подготовки
|
|
//--- финальной записи журнала (комментарий отражает фактическое поле).
|
|
const ulong exit_us = m_clock.NowUs();
|
|
if(obs.trade_event_id != 0)
|
|
m_journal.SetExitTime(obs.trade_event_id, exit_us);
|
|
obs.handler_exit_us = exit_us;
|
|
//--- R4-S4: длительность в первичном журнале фиксируется независимо
|
|
//--- от результата корреляции (для seq==0 тоже).
|
|
if(seq == 0)
|
|
return;
|
|
if(exit_us >= obs.entry_us && exit_us > 0)
|
|
m_tracker.SetHandlerDuration(seq, exit_us - obs.entry_us);
|
|
//--- R5-B6: режим/окно закреплены при отправке (PrepareAndDispatch);
|
|
//--- поздний callback не должен перезаписывать рыночный контекст.
|
|
}
|
|
//--- локальный прогон: полная серия на синтетических данных
|
|
bool CAppController::RunSynthetic(LabError &error)
|
|
{
|
|
error.Reset();
|
|
error.component = LAB_COMP_CONTROLLER;
|
|
m_session_id = NewSessionId();
|
|
if(m_session_id == "")
|
|
{
|
|
//--- R14-B1: уникальный каталог сессии не выдан (перебор кандидатов
|
|
//--- исчерпан) — ошибка ДО любой записи/торговли.
|
|
error.code = 45;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "cannot allocate unique session id";
|
|
return(false);
|
|
}
|
|
m_termination_reason = ""; // R11-S1: свежий статус для нового прогона
|
|
if(!m_clock.IsSynthetic())
|
|
{
|
|
error.code = 30;
|
|
error.message = "synthetic run requires synthetic clock";
|
|
m_session_id = ""; // R13-B1: сессия не начата
|
|
return(false);
|
|
}
|
|
m_journal.Init(m_settings.event_capacity, m_session_id);
|
|
m_tracker.Init(m_settings.sample_capacity, m_session_id);
|
|
if(!m_runner.Init(m_settings, error))
|
|
{
|
|
m_session_id = ""; // R13-B1: данные ещё не записаны
|
|
return(false);
|
|
}
|
|
LabEvent start;
|
|
start.Zero();
|
|
start.kind = LAB_KIND_RUN_START;
|
|
start.local_time_us = m_clock.NowUs();
|
|
start.payload_json = StringFormat("{\"mode\":\"LOCAL_ONLY\",\"config\":\"%s\"}",
|
|
m_settings.configuration_id);
|
|
ulong es = 0;
|
|
m_journal.Add(start, es);
|
|
//--- прогревочные планы
|
|
RequestPlan warm;
|
|
while(m_runner.NextWarmup(warm))
|
|
{
|
|
if(!SyntheticDispatch(warm))
|
|
{
|
|
error.code = 44;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "synthetic durable write failed (warmup)";
|
|
//--- R13-B1: контекст НЕЗАВЕРШЁННОЙ сессии СОХРАНЯЕТСЯ (id,
|
|
//--- tracker, журнал) — частичный диагностический отчёт
|
|
//--- финализируется как INCOMPLETE, а не стирается молча.
|
|
TerminateRun("SYNTHETIC_DURABLE_WRITE_FAILED");
|
|
return(false);
|
|
}
|
|
}
|
|
//--- основные планы
|
|
RequestPlan main;
|
|
while(m_runner.NextMain(main))
|
|
{
|
|
if(!SyntheticDispatch(main))
|
|
{
|
|
error.code = 44;
|
|
error.severity = LAB_SEV_BLOCKER;
|
|
error.message = "synthetic durable write failed (main)";
|
|
//--- R13-B1: см. ветвь warmup — контекст сохраняется до финализации.
|
|
TerminateRun("SYNTHETIC_DURABLE_WRITE_FAILED");
|
|
return(false);
|
|
}
|
|
}
|
|
m_state = LAB_STATE_DRAINING;
|
|
return(true);
|
|
}
|
|
//--- синтетическая отправка: транспорт-счётчик, никакой реальной торговли
|
|
bool CAppController::SyntheticDispatch(const RequestPlan &plan)
|
|
{
|
|
ulong seq = plan.sequence;
|
|
RequestPlan p = plan;
|
|
if(!m_tracker.Reserve(p, seq))
|
|
return(false);
|
|
//--- R13-B1: инъекция отказа ДО записи INTENT (первый INTENT);
|
|
//--- резерв до T0 снимается той же ветвью, что и штатный отказ.
|
|
if(m_test_fail_intent_counter > 0)
|
|
{
|
|
m_test_fail_intent_counter--;
|
|
if(m_test_fail_intent_counter == 0)
|
|
{
|
|
m_tracker.ReleaseReservation(seq);
|
|
return(false);
|
|
}
|
|
}
|
|
//--- синтетические T0/T1
|
|
//--- R10-S1: локальный прогон пишет устойчивый журнал намерений
|
|
//--- (INTENT до T0) - intents.csv сессии существует,
|
|
//--- CountUncertainDurableIntents работает и на LOCAL_ONLY.
|
|
if(!AppendDurableIntent(seq, "INTENT", m_clock.NowUs()))
|
|
{
|
|
//--- R13-B1: отказ ДО T0 — подготовительный резерв НЕ является
|
|
//--- выполненным наблюдением; снимаем фантомный MAIN, чтобы он не
|
|
//--- попадал в samples и повтор той же sequence был возможен.
|
|
m_tracker.ReleaseReservation(seq);
|
|
return(false);
|
|
}
|
|
const ulong t0 = m_clock.NowUs();
|
|
m_tracker.MarkSendStart(seq, t0);
|
|
m_clock.AdvanceUs(40);
|
|
const ulong t1 = m_clock.NowUs();
|
|
SendObservation send;
|
|
send.Zero();
|
|
send.t0_us = t0;
|
|
send.t1_us = t1;
|
|
send.send_ok = true;
|
|
send.result.retcode = TRADE_RETCODE_DONE;
|
|
send.result.request_id = (uint)seq;
|
|
send.result.order = (ulong)(100000 + seq);
|
|
send.result_known = true;
|
|
m_tracker.MarkSendReturn(seq, send);
|
|
//--- R13-B1: инъекция отказа на записи SEND_RETURN (после T1):
|
|
//--- запись остаётся DISPATCH_UNCERTAIN (INTENT без SEND_RETURN).
|
|
if(m_test_fail_intent_counter > 0)
|
|
{
|
|
m_test_fail_intent_counter--;
|
|
if(m_test_fail_intent_counter == 0)
|
|
return(false);
|
|
}
|
|
//--- R10-S1: SEND_RETURN (после T1) - пара для INTENT в durable-журнале.
|
|
if(!AppendDurableIntent(seq, "SEND_RETURN", t1))
|
|
return(false);
|
|
m_clock.AdvanceUs(500);
|
|
const ulong t2 = m_clock.NowUs();
|
|
m_tracker.MarkRequest(seq, t2, send.result.request_id, TRADE_RETCODE_DONE, 0,
|
|
send.result.order, 0);
|
|
m_clock.AdvanceUs(800);
|
|
const ulong t3 = m_clock.NowUs();
|
|
m_tracker.MarkOrderAdd(seq, t3, send.result.order, "SYNTHETIC");
|
|
m_clock.AdvanceUs(1200);
|
|
const ulong t4 = m_clock.NowUs();
|
|
//--- R10-S1: согласованный объём DEAL-callback - единственная сделка
|
|
m_tracker.MarkDealAdd(seq, t4, (ulong)(200000 + seq), p.volume);
|
|
m_tracker.SetDealCoverageComplete(seq);
|
|
//--- финал: синтетическая сверка подтверждает
|
|
m_clock.AdvanceUs(900);
|
|
const ulong t6 = m_clock.NowUs();
|
|
ulong eid = 0;
|
|
LabEvent e;
|
|
e.Zero();
|
|
e.kind = LAB_KIND_CONFIRMATION_SNAPSHOT;
|
|
e.local_time_us = t6;
|
|
e.origin_sequence = seq;
|
|
m_journal.Add(e, eid);
|
|
m_tracker.SetFinal(seq, t6, LAB_OUT_FILLED, LAB_FINAL_CALLBACK_CHECK, eid);
|
|
m_tracker.SetTradingStateKnown(seq, true);
|
|
m_tracker.CloseCollection(seq, t6, eid);
|
|
m_tracker.SetDeadlinesFromSample(seq, (ulong)m_settings.outcome_timeout_ms * 1000,
|
|
(ulong)m_settings.late_grace_ms * 1000);
|
|
m_runner.CountDispatched(plan.role);
|
|
m_clock.AdvanceUs(m_settings.pause_ms * 1000);
|
|
return(true);
|
|
}
|
|
//--- R14-B1: идентификатор запуска ОТДЕЛЁН от измерительной временной
|
|
//--- шкалы: базовый суффикс времени дополняется монотонным счётчиком
|
|
//--- сессий контроллера, а кандидат дополнительно проверяется по
|
|
//--- существованию каталога на диске (защита от коллизии после
|
|
//--- ПЕРЕУСТАНОВКИ EA, когда счётчик начинается заново). Каталог
|
|
//--- занят, если в нём есть любой известный файл данных сессии
|
|
//--- (intents/manifest/checkpoint.marker/samples/events/…) — перезапись
|
|
//--- прежнего отчёта исключена (см. SessionDirUsed).
|
|
//--- При невозможности выдать отдельный каталог возвращается "" и
|
|
//--- вызывающий код завершается ошибкой до новой записи/торговли.
|
|
//--- R15-B2: признак занятости каталога сессии по ЛЮБОМУ известному
|
|
//--- файлу отчёта (не только трём служебным именам): после прерванного
|
|
//--- экспорта в каталоге могут остаться ТОЛЬКО samples/events (marker
|
|
//--- уже снят, manifest ещё не записан) — такой каталог тоже занят.
|
|
bool CAppController::SessionDirUsed(const string dir)
|
|
{
|
|
string names[] = {"intents.csv", "manifest.csv", "checkpoint.marker",
|
|
"samples.csv", "events.csv", "schedule.csv", "deals.csv",
|
|
"summary.csv", "histogram.csv", "offset_signs.csv",
|
|
"comparisons.csv", "market_windows.csv", "market_ticks.csv",
|
|
"checks.csv", "experiment_manifest.csv"
|
|
};
|
|
for(int i = 0; i < ArraySize(names); i++)
|
|
{
|
|
if(FileIsExist(dir + "\\" + names[i]))
|
|
return(true);
|
|
}
|
|
return(false);
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
string CAppController::NewSessionId(void)
|
|
{
|
|
const string base = StringFormat("%s_%s_S%02d", m_settings.campaign_id,
|
|
LabExperimentName(m_settings.experiment_id),
|
|
m_settings.series_id);
|
|
const ulong t0 = m_clock.NowUs();
|
|
//--- R14-B1/R15-B2: монотонный курсор сессий — счётчик растёт ВНУТРИ
|
|
//--- цикла и НЕ сбрасывается между вызовами (иначе первый свободный
|
|
//--- каталог при фиксированном наборе занятых слотов детерминирован
|
|
//--- и повтор сессии снова коллизирует). Каждый кандидат проверяется
|
|
//--- по ЛЮБОМУ известному файлу отчёта (SessionDirUsed: intents/manifest/
|
|
//--- marker/samples/events/schedule/deals/summary/histogram/…): это
|
|
//--- защищает и случай ПЕРЕУСТАНОВКИ EA, когда счётчик начинается
|
|
//--- заново, и частичный остаток прерванного экспорта — прежний
|
|
//--- каталог пропускается, перезапись отчёта исключена.
|
|
for(int tries = 0; tries < 4096; tries++)
|
|
{
|
|
m_session_counter++;
|
|
const string cand = StringFormat("%s_%08X", base,
|
|
(uint)((t0 + m_session_counter) & 0xFFFFFFFF));
|
|
const string dir = StringFormat("RequestLatencyLab\\%s\\%s",
|
|
m_settings.campaign_id, cand);
|
|
if(!SessionDirUsed(dir))
|
|
return(cand);
|
|
}
|
|
return("");
|
|
}
|
|
|
|
//+------------------------------------------------------------------+
|
|
//| |
|
|
//+------------------------------------------------------------------+
|
|
void CAppController::SetSettings(const LabSettings &s) { m_settings = s; }
|
|
LabSettings CAppController::Settings() { return(m_settings); }
|
|
|
|
|
|
#endif // REQUEST_LATENCY_LAB_APP_CONTROLLER_MQH
|
|
//+------------------------------------------------------------------+ |