mql5-execution-microstructu.../Include/RequestLatencyLab/AppController.mqh

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 &regime);
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 &regime)
{
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
//+------------------------------------------------------------------+