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