MetaTrader 4 + Python を ZeroMQ で接続する:安全で検証可能なブリッジ・プロトコル(2026)

2026年メンテナンス注記。 2019年の記事は実際に起きた二つのインストール問題を記録していますが、ブリッジの仕様や動作するコードはありませんでした。本版はその短い記録と題名を基に新しいチュートリアルを構成したものです。末尾に元のエクスポート全体を逐語保存し、Visual C++ 2015 という記述は特定ビルドの要件として扱います。2026年の普遍的な前提ではありません。

ZeroMQ は MetaTrader 4 Expert Advisor(EA)と Python の間でメッセージを運べます。しかし、これはあくまでトランスポートです。プロトコルを定義せず、信頼できないプロセスを安全にせず、注文の成否不明状態を解消せず、戦略を収益化もしません。以下はローカルの health 要求だけを実装します。通信障害を学ぶための例であり、ライブ取引システムではありません。

1. 対象範囲と安全契約

最初はオフライン端末またはデモ口座で試し、トランスポートの検証中は自動売買を無効にします。このサンプルは:

  • IPv4 ループバック tcp://127.0.0.1:5557 だけに bind する;
  • 許可リストにある health 操作だけを受け入れる;
  • 認証情報、口座識別子、価格、注文を送らない;
  • メッセージ、待機、キュー、再試行、終了に上限を設ける;
  • 同じ要求 ID と同じシリアライズ済みバイトを再送し、サーバーで重複排除する;
  • リターンや約定品質を約束しない。

Allow DLL imports を有効にすると、ネイティブコードに端末プロセスの権限を与えます。ソース、ビルド、アーキテクチャ、依存関係、チェックサムを確認した DLL だけを読み込んでください。コピーしてきた DLL は単なる接続部品ではなく、コード実行の信頼境界です。

2. アーキテクチャと信頼境界

保守的な分担では取引権限を MT4 に残します:

market / broker
      |
      v
MT4 terminal -> EA policy gate -> reviewed MQL4-to-libzmq adapter
                                    || loopback, fixed protocol
                                    /
                              Python worker
                         analytics / health only

MT4 が端末状態、tick 時刻、口座モード、リスク上限、将来追加され得る注文送信を管理します。Python は計算やヘルスチェックへの応答を行えますが、その返信は EA のローカルポリシーへの入力にすぎません。「ZeroMQ の返信を受信した」を「安全に実行できる」に読み替えてはいけません。

境界には、MQL4 宣言、ネイティブ・アダプター/libzmq バイナリ、Python/pyzmq という三つの独立したバージョン要素があります。Python が正しくても、版や ABI の不一致は端末をクラッシュさせ得ます。MQL4 コンパイラーは外部関数の引数を完全には検証できず、MetaQuotes は DLL 呼び出しが呼出元モジュールのスレッドで動くと説明しています。

3. ランタイムを推測せず MT4 DLL を設定する

MetaTrader 4 Expert Advisor の DLL インポート設定

保存された記事は Microsoft Visual C++ 2015 Redistributable を求めています。当時の libzmq.dll には正しかった可能性がありますが、本当の規則は「監査したそのバイナリが要求するランタイムを入れる」です。現在のビルドは別のランタイムを使うか、静的リンクかもしれません。

EA を読み込む前に:

  1. 端末と DLL のアーキテクチャを一致させる。ファイル名変更で不一致を隠さない。
  2. 監査済みライブラリを端末文書の MQL4/Libraries 検索パスに置き、依存 DLL もすべて確認する。
  3. MQL4 の各 #import シグネチャと呼出規約をアダプター ABI に一致させる。広い C API の直接公開より、小さく監査されたアダプターを優先する。
  4. Allow DLL imports はこの EA にだけ許可する。実行時に IsDllsAllowed() が false なら fail closed にする。
  5. ExpertsJournal を調べる。MetaQuotes によれば DLL が欠落または禁止されると EA は再初期化まで停止する。

ファイル名が本稿と一致するという理由だけで、任意の「MT4 ZeroMQ bridge」をダウンロードしないでください。

4. OnTick ではなく OnTimer で駆動する

OnTick() は新しい相場更新のためのもので、信頼できるメッセージポンプではありません。MetaQuotes は、先の OnTick() が実行中なら新しい tick は無視されると説明しています。そこで受信をブロックすると価格が古くなり、EA のほかの処理も止まります。

OnInit() でタイマーを一つ作り、OnTimer() では小さく有界またはノンブロッキングな一工程だけを行い、OnDeinit() で解除します。MQL4 プログラムごとにタイマーは一つで、タイマーイベントが待機中または実行中なら次は追加されません。これはバックプレッシャーになりますが、各ハンドラーが短時間で終了することも必要です。イベント処理とワーカースレッドで ZeroMQ socket を共有しないでください。libzmq は通常の REQ/REP socket がスレッドセーフではないと明記しています。

5. 厳密な multipart プロトコルを定義する

ZeroMQ は multipart メッセージの原子性を保ちますが、アプリケーション schema は作りません。本稿ではアプリから見えるフレームをちょうど二つにします:

フレームバイト数規則
07ASCII リテラル LZMT4/1
1最大 65,536UTF-8 JSON の厳密なオブジェクト。重複キー、NaN、無限大は禁止

要求 schema(追加キーは禁止):

フィールド型と制約
schemaリテラル lazying.mt4.bridge.request
version整数 1
request_idUUID 文字列。再試行で変更しない
client_idASCII の英数字と ._- で 1〜64 文字
session_idUUID 文字列。意図的に開始した各クライアント・セッションで新規作成
sequence1 から一つずつ増える正整数
sent_at_utcZ で終わる RFC 3339 UTC 時刻
max_age_ms100〜60,000 の整数
operationこの版ではリテラル health
payload空の JSON オブジェクト

返信のキーも schemaversion、三つの相関値(request_idsession_idsequence)、statusprocessed_at_utcresulterror に限定します。成功は status: "ok"、オブジェクトの resulterror: null。拒否は status: "error"result: null、安定した code と安全な message を持つエラーオブジェクトです。相関値が null になれるのは要求をデコードできない場合だけです。

この境界では send_pyobj()recv_pyobj() を使わないでください。PyZMQ の文書によればこれらは Python pickle を使い、信頼できない pickle のデコードは任意コードを実行し得ます。

6. 価格と時刻の鮮度を別々に定義する

少なくとも三種類の時計があります:

  • MqlTick.time は一つの銘柄について最新既知 tick に付いたサーバー時刻。
  • OnTimer() 内の TimeCurrent() は Market Watch 内のいずれかの銘柄に関する最終既知サーバー時刻であり、対象銘柄が直前に動いた証明ではない。
  • monotonic clock はローカルの経過時間とタイムアウトに使えるが、別ホストのイベントを時刻付けできない。

将来の価格メッセージには symbolbidasktick_server_timeterminal_received_at_utc、有効期間を持たせるべきです。EA は動作の直前に対象銘柄の最新 tick と端末の現在状態を比較します。Python の wall clock は価格が新しい証拠になりません。許容した小さな clock skew より未来の時刻と max_age_ms より古い要求を拒否し、別ホストなら時計も同期します。

ヘルスチェック例は UTC で要求の age を調べ、単調時計で再試行時間を測ります。価格は含みません。

7. REQ/REP の障害状態を理解する

REQ/REP は lock-step 状態機械です。REQ は送信後に受信、REP は受信後に送信します。厳格な REQ がタイムアウトした後、同じ socket で再度送信すると状態機械に違反します。socket を作り直して再接続し、同じ request_id を持つ同じシリアライズ済み要求を再送します。

難しいものの正常なケース:

  1. Python が要求を送る。
  2. サーバーは処理を終える。
  3. 返信が遅延または失われる。
  4. Python はタイムアウトし、処理されたか判断できない。

そのため、トランスポート再試行にはアプリ層の冪等性が必要です。下のサーバーは (client_id, session_id, request_id) で成功返信をキャッシュし、同一 ID が異なるバイトで再利用されたら拒否します。シーケンス検査より先にキャッシュを調べるので、本当の再送は以前の返信を受け取れます。このメモリーキャッシュは health デモには十分ですが、注文には不十分です。実運用の executor には永続台帳と MT4 の権威ある状態との照合が必要です。

対応する libzmq ビルドには relaxed/correlated REQ オプションもあります。公式文書は relaxed モードが以前の返信を破棄し、古い返信と新しい要求の取り違えを防ぐため correlation と組み合わせるべきだと注意しています。本例は障害を明示するため、厳格な REQ と socket の再作成を選びます。

8. キュー、待機、終了を有界にする

接続関連の socket オプションは connect() または bind() より前に設定します:

制御例の方針理由
RCVTIMEOSNDTIMEOクライアント 600 ms無限ブロックを避け EAGAIN/zmq.Again を処理
LINGER使い捨て要求 socket は 0作り直しと終了を有界化し、未送信メッセージを破棄
SNDHWMRCVHWM10 メッセージトランスポート・キューを制限。実効動作は transport/socket に依存
MAXMSGSIZE65,536 バイトlibzmq で巨大な受信フレームを拒否し、アプリでも検査
IMMEDIATEクライアントは 1接続成立前に送信データをキューへ入れない
RECONNECT_IVLRECONNECT_IVL_MAX100/1,000 ms再接続の間隔を制限しバックオフ

デフォルトの linger と timeout は無限です。EA のライフサイクルでそのまま使わないでください。high-water mark はバイト数でなくメッセージ数であり、65,536 バイトのアプリ上限の代用でもありません。

9. 認証と暗号化

ループバックと固定ポートは最も安全な開始範囲ですが、別のローカルプロセスは接続できる可能性があります。OS ファイアウォール規則を加え、最小権限アカウントで動かします。

ネットワークを越える場合、ZeroMQ のセキュリティ機構は同等ではありません:

  • ZMTP NULL は認証も機密性も提供しない。
  • PLAIN はユーザー名とパスワードを暗号化せず送るため、単独では信頼できないネットワークに不適切。
  • CURVE は ZeroMQ で認証と機密性を提供するための機構。

正確な libzmq ビルドと MQL4 アダプターの両方が互換 CURVE 設定を公開し、鍵配布とサーバー鍵 pinning を確認した場合だけ CURVE を使います。そうでなければ endpoint を非公開に保ち、別に認証された暗号化トンネルを加えます。生の MT4 bridge ポートをインターネットへ公開しないでください。メタデータと鍵も保護し、秘密をログへ書かないことが必要です。

10. 冪等性と順序

次の値は別概念です:

  • request_id:一つの論理要求の識別子。すべての再試行で一定。
  • session_id:意図した一回のクライアント実行の識別子。新セッションでのみ変更。
  • sequence:セッション内の期待順序。暗黙にリセットしない。

返信と、正規化要求バイトのハッシュをキャッシュします。同じ ID が異なるバイトで来れば拒否します。キャッシュ済み要求に一致しない sequence の飛びや巻き戻しも拒否します。サンプルは返信キャッシュとセッション別 sequence map の両方を 1,024 項・60 秒に制限します。期限切れまたは追い出し後、継続するクライアントは新しい session_idを生成し、その新セッションを sequence 1 から始めます。旧セッションの続き位置を推測してはいけません。これはクライアント側の契約です。有界なデモサーバーは期限切れ ID を永久には記憶できず、そのままでは古い ID の sequence 1 も受理し得ます。

health は本質的に無害ですが、取引命令は違います。追加前に永続的なコマンド台帳、明確な状態機械(receivedvalidatedsubmittedconfirmedrejectedunknown)、プラットフォームの ticket/history との照合を設計します。unknown の送信を確実な失敗だと見なして再試行してはいけません。

11. 再接続、watchdog、circuit breaker

トランスポートの再接続はアプリの復旧ではありません。有用な watchdog は、最後の正しく相関した返信、連続 timeout、socket 再作成、protocol error、銘柄ごとの最新 tick を追跡します。少数回の連続失敗で circuit breaker を開きます:

  • 新しい外部コマンドを受け付けない;
  • bridge 障害だけを理由にポジションを決済または反転しない;
  • MT4 内部の観測を続ける;
  • 一定期間の正常性を確認し、可能なら手動で再 arm する。

Python、EA、端末、コンピューターが再起動したら新セッションを作り、コマンド受付前に永続状態を照合します。TCP 再接続だけからアプリ状態を推測してはいけません。

12. 最小ローカル・ヘルスチェック

空の一時ディレクトリに次の三ファイルを作ります。注文操作は意図的に含めません。

protocol.py

from __future__ import annotations

import json
import re
import uuid
from datetime import datetime, timezone
from typing import Any

PROTOCOL = b"LZMT4/1"
MAX_FRAME_BYTES = 65_536
REQUEST_SCHEMA = "lazying.mt4.bridge.request"
REPLY_SCHEMA = "lazying.mt4.bridge.reply"
REQUEST_KEYS = {
    "schema", "version", "request_id", "client_id", "session_id",
    "sequence", "sent_at_utc", "max_age_ms", "operation", "payload",
}
REPLY_KEYS = {
    "schema", "version", "request_id", "session_id", "sequence",
    "status", "processed_at_utc", "result", "error",
}
CLIENT_ID = re.compile(r"[A-Za-z0-9_.-]{1,64}Z")
ERROR_CODE = re.compile(r"[A-Z0-9_]{1,64}Z")


class ProtocolError(ValueError):
    pass


def utc_now() -> str:
    return (
        datetime.now(timezone.utc)
        .isoformat(timespec="milliseconds")
        .replace("+00:00", "Z")
    )


def parse_utc(value: str) -> datetime:
    if not isinstance(value, str) or not value.endswith("Z"):
        raise ProtocolError("timestamp must be an RFC 3339 UTC string ending in Z")
    try:
        parsed = datetime.fromisoformat(value[:-1] + "+00:00")
    except ValueError as exc:
        raise ProtocolError("timestamp is invalid") from exc
    if parsed.tzinfo is None or parsed.utcoffset() != timezone.utc.utcoffset(parsed):
        raise ProtocolError("timestamp must use UTC")
    return parsed


def _no_duplicate_keys(pairs: list[tuple[str, Any]]) -> dict[str, Any]:
    result: dict[str, Any] = {}
    for key, value in pairs:
        if key in result:
            raise ProtocolError(f"duplicate JSON key: {key}")
        result[key] = value
    return result


def _reject_constant(value: str) -> None:
    raise ProtocolError(f"non-finite JSON number: {value}")


def encode_object(value: dict[str, Any]) -> bytes:
    try:
        raw = json.dumps(
            value,
            ensure_ascii=True,
            allow_nan=False,
            sort_keys=True,
            separators=(",", ":"),
        ).encode("utf-8")
    except (TypeError, ValueError) as exc:
        raise ProtocolError("object is not strict JSON") from exc
    if len(raw) > MAX_FRAME_BYTES:
        raise ProtocolError("JSON frame is too large")
    return raw


def decode_object(raw: bytes) -> dict[str, Any]:
    if len(raw) > MAX_FRAME_BYTES:
        raise ProtocolError("JSON frame is too large")
    try:
        value = json.loads(
            raw.decode("utf-8"),
            object_pairs_hook=_no_duplicate_keys,
            parse_constant=_reject_constant,
        )
    except (UnicodeDecodeError, json.JSONDecodeError) as exc:
        raise ProtocolError("frame is not strict UTF-8 JSON") from exc
    if not isinstance(value, dict):
        raise ProtocolError("JSON root must be an object")
    return value


def _uuid(value: Any, field: str) -> None:
    if not isinstance(value, str):
        raise ProtocolError(f"{field} must be a UUID string")
    try:
        uuid.UUID(value)
    except ValueError as exc:
        raise ProtocolError(f"{field} must be a UUID string") from exc


def decode_request(parts: list[bytes]) -> tuple[dict[str, Any], bytes]:
    if len(parts) != 2 or parts[0] != PROTOCOL:
        raise ProtocolError("expected exactly two frames with protocol LZMT4/1")
    request = decode_object(parts[1])
    if set(request) != REQUEST_KEYS:
        raise ProtocolError("request keys do not match version 1")
    if request["schema"] != REQUEST_SCHEMA or request["version"] != 1:
        raise ProtocolError("unsupported request schema or version")
    _uuid(request["request_id"], "request_id")
    _uuid(request["session_id"], "session_id")
    if not isinstance(request["client_id"], str) or not CLIENT_ID.fullmatch(request["client_id"]):
        raise ProtocolError("client_id has an invalid format")
    if type(request["sequence"]) is not int or request["sequence"] < 1:
        raise ProtocolError("sequence must be a positive integer")
    if type(request["max_age_ms"]) is not int or not 100 <= request["max_age_ms"] <= 60_000:
        raise ProtocolError("max_age_ms is out of range")
    parse_utc(request["sent_at_utc"])
    if request["operation"] != "health" or request["payload"] != {}:
        raise ProtocolError("only health with an empty payload is allowed")
    return request, parts[1]


def make_reply(
    request: dict[str, Any] | None,
    *,
    status: str,
    result: dict[str, Any] | None,
    error: dict[str, str] | None,
) -> dict[str, Any]:
    return {
        "schema": REPLY_SCHEMA,
        "version": 1,
        "request_id": request.get("request_id") if request else None,
        "session_id": request.get("session_id") if request else None,
        "sequence": request.get("sequence") if request else None,
        "status": status,
        "processed_at_utc": utc_now(),
        "result": result,
        "error": error,
    }


def decode_reply(parts: list[bytes]) -> dict[str, Any]:
    if len(parts) != 2 or parts[0] != PROTOCOL:
        raise ProtocolError("reply framing is invalid")
    reply = decode_object(parts[1])
    if set(reply) != REPLY_KEYS:
        raise ProtocolError("reply keys do not match version 1")
    if reply["schema"] != REPLY_SCHEMA or reply["version"] != 1:
        raise ProtocolError("unsupported reply schema or version")
    if reply["status"] not in {"ok", "error"}:
        raise ProtocolError("reply status is invalid")
    for field in ("request_id", "session_id"):
        if reply[field] is not None:
            _uuid(reply[field], field)
    if reply["sequence"] is not None and (
        type(reply["sequence"]) is not int or reply["sequence"] < 1
    ):
        raise ProtocolError("reply sequence is invalid")
    parse_utc(reply["processed_at_utc"])
    if reply["status"] == "ok":
        if not isinstance(reply["result"], dict) or reply["error"] is not None:
            raise ProtocolError("success reply shape is invalid")
    else:
        error = reply["error"]
        if (
            reply["result"] is not None
            or not isinstance(error, dict)
            or set(error) != {"code", "message"}
            or not isinstance(error.get("code"), str)
            or not ERROR_CODE.fullmatch(error["code"])
            or not isinstance(error.get("message"), str)
            or not 1 <= len(error["message"]) <= 256
        ):
            raise ProtocolError("error reply shape is invalid")
    return reply

health_server.py

from __future__ import annotations

import argparse
import hashlib
import json
import time
from collections import OrderedDict
from datetime import datetime, timezone

import zmq

from protocol import (
    MAX_FRAME_BYTES,
    PROTOCOL,
    ProtocolError,
    decode_request,
    encode_object,
    make_reply,
    parse_utc,
)

ENDPOINT = "tcp://127.0.0.1:5557"
CACHE_LIMIT = 1_024
CACHE_TTL_SECONDS = 60
STREAM_LIMIT = 1_024
STREAM_TTL_SECONDS = 60
FUTURE_SKEW_MS = 5_000


def log(event: str, **fields: object) -> None:
    print(json.dumps({"event": event, **fields}, sort_keys=True), flush=True)


def configure(socket: zmq.Socket) -> None:
    socket.setsockopt(zmq.LINGER, 0)
    socket.setsockopt(zmq.RCVTIMEO, 1_000)
    socket.setsockopt(zmq.SNDTIMEO, 1_000)
    socket.setsockopt(zmq.RCVHWM, 10)
    socket.setsockopt(zmq.SNDHWM, 10)
    socket.setsockopt(zmq.MAXMSGSIZE, MAX_FRAME_BYTES)


def prune_expired(mapping: OrderedDict, now_monotonic: float, ttl_seconds: float) -> None:
    for expired_key in [
        item_key
        for item_key, item in mapping.items()
        if now_monotonic - item[-1] > ttl_seconds
    ]:
        del mapping[expired_key]


def trim_oldest(mapping: OrderedDict, limit: int) -> None:
    while len(mapping) > limit:
        mapping.popitem(last=False)


def main() -> int:
    parser = argparse.ArgumentParser()
    parser.add_argument("--delay-first-ms", type=int, default=0)
    args = parser.parse_args()
    if not 0 <= args.delay_first_ms <= 5_000:
        parser.error("--delay-first-ms must be between 0 and 5000")

    cache: OrderedDict[
        tuple[str, str, str], tuple[str, dict[str, object], float]
    ] = OrderedDict()
    stream_state: OrderedDict[tuple[str, str], tuple[int, float]] = OrderedDict()
    delayed = False
    context = zmq.Context()
    socket = context.socket(zmq.REP)
    configure(socket)
    socket.bind(ENDPOINT)
    log("ready", endpoint=ENDPOINT, pyzmq=zmq.pyzmq_version(), libzmq=zmq.zmq_version())

    try:
        while True:
            try:
                parts = socket.recv_multipart()
            except zmq.Again:
                continue

            started = time.monotonic()
            request = None
            duplicate = False
            try:
                request, raw = decode_request(parts)
                now = datetime.now(timezone.utc)
                age_ms = (now - parse_utc(request["sent_at_utc"])).total_seconds() * 1_000
                if age_ms < -FUTURE_SKEW_MS:
                    raise ProtocolError("request timestamp is too far in the future")
                if age_ms > request["max_age_ms"]:
                    raise ProtocolError("request is stale")

                key = (request["client_id"], request["session_id"], request["request_id"])
                stream = (request["client_id"], request["session_id"])
                fingerprint = hashlib.sha256(raw).hexdigest()
                now_monotonic = time.monotonic()
                prune_expired(cache, now_monotonic, CACHE_TTL_SECONDS)
                prune_expired(stream_state, now_monotonic, STREAM_TTL_SECONDS)
                cached = cache.get(key)
                if cached is not None:
                    if cached[0] != fingerprint:
                        raise ProtocolError("request_id was reused with different bytes")
                    reply = cached[1]
                    duplicate = True
                else:
                    expected = stream_state.get(stream, (0, now_monotonic))[0] + 1
                    if request["sequence"] != expected:
                        raise ProtocolError(f"expected sequence {expected}")
                    reply = make_reply(
                        request,
                        status="ok",
                        result={"service": "python-local-health", "protocol": "LZMT4/1"},
                        error=None,
                    )
                    cache[key] = (fingerprint, reply, now_monotonic)
                    stream_state[stream] = (request["sequence"], now_monotonic)
                    stream_state.move_to_end(stream)
                    trim_oldest(cache, CACHE_LIMIT)
                    trim_oldest(stream_state, STREAM_LIMIT)
            except ProtocolError as exc:
                reply = make_reply(
                    request,
                    status="error",
                    result=None,
                    error={"code": "PROTOCOL_REJECTED", "message": str(exc)},
                )

            if args.delay_first_ms and not delayed and reply["status"] == "ok":
                delayed = True
                time.sleep(args.delay_first_ms / 1_000)

            try:
                socket.send_multipart([PROTOCOL, encode_object(reply)])
            except zmq.Again:
                log("send_timeout", action="exit_for_supervisor_restart")
                return 2
            log(
                "reply",
                request_id=reply["request_id"],
                sequence=reply["sequence"],
                status=reply["status"],
                duplicate=duplicate,
                elapsed_ms=round((time.monotonic() - started) * 1_000, 1),
            )
    except KeyboardInterrupt:
        log("stopping")
        return 0
    finally:
        socket.close(linger=0)
        context.term()


if __name__ == "__main__":
    raise SystemExit(main())

health_client.py

from __future__ import annotations

import json
import time
import uuid

import zmq

from protocol import MAX_FRAME_BYTES, PROTOCOL, ProtocolError, decode_reply, encode_object, utc_now

ENDPOINT = "tcp://127.0.0.1:5557"
TIMEOUT_MS = 600
MAX_ATTEMPTS = 3


def open_request_socket(context: zmq.Context) -> zmq.Socket:
    socket = context.socket(zmq.REQ)
    socket.setsockopt(zmq.LINGER, 0)
    socket.setsockopt(zmq.RCVTIMEO, TIMEOUT_MS)
    socket.setsockopt(zmq.SNDTIMEO, TIMEOUT_MS)
    socket.setsockopt(zmq.RCVHWM, 10)
    socket.setsockopt(zmq.SNDHWM, 10)
    socket.setsockopt(zmq.MAXMSGSIZE, MAX_FRAME_BYTES)
    socket.setsockopt(zmq.IMMEDIATE, 1)
    socket.setsockopt(zmq.RECONNECT_IVL, 100)
    socket.setsockopt(zmq.RECONNECT_IVL_MAX, 1_000)
    socket.connect(ENDPOINT)
    return socket


def main() -> int:
    request = {
        "schema": "lazying.mt4.bridge.request",
        "version": 1,
        "request_id": str(uuid.uuid4()),
        "client_id": "local-smoke-test",
        "session_id": str(uuid.uuid4()),
        "sequence": 1,
        "sent_at_utc": utc_now(),
        "max_age_ms": 10_000,
        "operation": "health",
        "payload": {},
    }
    frames = [PROTOCOL, encode_object(request)]
    context = zmq.Context()
    started = time.monotonic()
    try:
        for attempt in range(1, MAX_ATTEMPTS + 1):
            socket = open_request_socket(context)
            try:
                socket.send_multipart(frames)
                reply = decode_reply(socket.recv_multipart())
            except zmq.Again:
                print(json.dumps({"event": "timeout", "attempt": attempt}))
            except ProtocolError as exc:
                print(json.dumps({"event": "protocol_error", "message": str(exc)}))
                return 2
            else:
                if (
                    reply["request_id"] != request["request_id"]
                    or reply["session_id"] != request["session_id"]
                    or reply["sequence"] != request["sequence"]
                ):
                    print(json.dumps({"event": "correlation_error"}))
                    return 2
                print(json.dumps({
                    "event": "reply",
                    "attempt": attempt,
                    "elapsed_ms": round((time.monotonic() - started) * 1_000, 1),
                    "reply": reply,
                }, sort_keys=True))
                return 0 if reply["status"] == "ok" else 2
            finally:
                socket.close(linger=0)
            time.sleep(0.15)
        print(json.dumps({"event": "unavailable", "attempts": MAX_ATTEMPTS}))
        return 1
    finally:
        context.term()


if __name__ == "__main__":
    raise SystemExit(main())

各 socket は同一スレッドで作成、使用、終了します。シリアライズ済み frames は再試行ループ外で一度だけ構築します。

13. MT4 より先にローカルで試す

現在の Python 環境と、それに適した公式 PyZMQ リリースを使います:

python -m venv .venv
. .venv/bin/activate                  # Windows PowerShell: .venvScriptsActivate.ps1
python -m pip install --upgrade pip pyzmq
python health_server.py

二つ目のターミナルで:

python health_client.py

次にサーバーを止め、最初の返信を意図的に遅らせて起動します:

python health_server.py --delay-first-ms 900

クライアントは最初の timeout を報告し、厳格な REQ socket を再作成し、同じバイトを再送して、後の試行でキャッシュ済み返信を受け取るはずです。サーバーは後の要求を duplicate: true と記録します。サーバーなし、不正 JSON、余分なフレーム、内容を変えた同一 ID、sequence の飛び、期限切れ timestamp、Ctrl-C による正常終了も試してください。

このテストが証明するのはローカル framing、検証、timeout、再試行だけです。MQL4 アダプターの検証ではありません。その境界には双方で byte-for-byte fixture テストを追加し、使い捨て端末で bitness/ABI を確認し、通常プロファイルで DLL を許可する前に MQL4 が生成したフレームと Python fixture を比較します。

14. 拡張前に必要な MT4 側の安全ゲート

次のすべてを明文化して個別テストするまで、操作は health のみにします:

  • デモ口座ゲート(IsDemo())と明示的な手動 arm/disarm;
  • DLL 許可と、端末自身の取引許可/取引コンテキスト検査;
  • 操作と銘柄の allow-list。任意の MQL 関数名は受け付けない;
  • 動作直前に SymbolInfoTick() で対象銘柄の現在 tick の鮮度を検査;
  • 最大サイズ、exposure、未処理要求、spread/slippage 方針、日次損失;
  • MT4 の権威ある equity に基づき、手動再 arm が必要な latched daily-loss stop;
  • 送信済みまたは結果不明コマンドごとの永続冪等性と照合;
  • malformed、late、duplicate、out-of-order、unavailable への fail-closed。

Python に同じリスク制限があっても EA 内で強制してください。古い、侵害された、または切断された Python が端末側上限を緩められてはいけません。テストはデモだけで行い、デモの挙動がライブの流動性、slippage、fill を保証しないことも忘れないでください。

15. 可観測性と運用チェックリスト

UTC 時刻と monotonic duration を持つ構造化イベントを記録します。有用なフィールドは、プロトコル版、EA/adapter build、pyzmq/libzmq 版、request ID、session ID、sequence、operation、age、latency、結果、duplicate flag、reject code、連続 miss、socket rebuild 回数です。口座情報を秘匿し、鍵をログに残さず、将来戦略情報を含む可能性がある payload も丸ごと記録しません。

各テストの前に:

  1. バイナリの checksum、bitness、ABI 宣言、依存元、DLL 許可を確認する。
  2. loopback bind と firewall 範囲を確認し、予期しないプロセスがポートを所有していないか調べる。
  3. 両端の protocol/version、厳格 framing、schema fixture、サイズ上限を確認する。
  4. timeout、HWM、linger、retry 回数、cache 容量、cache expiry が有界か確認する。
  5. watchdog が circuit を開き、黙って再 arm できないことを確認する。
  6. Python の停止、低速化、不正入力、再起動でも MT4 が応答することを確認する。
  7. 秘匿化したログ、バージョン、結果を保存する。

16. 主要な公式文書

17. 2019年の元エクスポート(逐語)

以下は歴史的証拠として元のエクスポート全体を byte-for-byte で保存したものです。相対画像参照はアーカイブの一部なので、意図的に書き換えていません。

---
id: 1968
title: 'MetaTrader 4 + Python: ZeroMQ'
slug: 'metatrader-4-python-zeromq'
date: '2019-07-02T14:05:13'
modified: '2019-07-02T14:19:47'
status: 'publish'
link: 'https://blog.lazying.art/en/html/securities-forex/metatrader/1968/metatrader-4-python-zeromq.html'
author: 'Lachlan Chen'
categories:
  - 'MetaTrader'
---

Issues:

- *Microsoft Visual C++ 2015 Redistributable* need to be installed
- *Allow DLL imports *should be checked when an EA is loaded: Cannot call ‘libzmq.dll::zmq_ctx_new’, DLL is not allowed

![](images/image-1024x498.png)

Leave a Reply