تاریخی مارکیٹ کا ڈیٹا عام طور پر ایک مکمل ڈیٹاسیٹ کے طور پر آتا ہے۔ یہ تجزیہ کے لیے آسان ہے، لیکن یہ اس سے بہت مختلف ہے کہ کس طرح ٹریڈنگ سافٹ ویئر حقیقی مارکیٹوں کا تجربہ کرتا ہے۔ پیداوار میں، واقعات ایک وقت میں ہوتے ہیں، مستقبل نامعلوم ہے، اور تمام فیصلے صرف اس بات پر منحصر ہیں کہ اب تک کیا ہوا ہے۔
اس ٹیوٹوریل میں، ہم اس تجربے کو دوبارہ بنانے کے لیے تاریخی ٹک ڈیٹا استعمال کریں گے۔ ہم EODHD سے ایک پورا AAPL ٹریڈنگ سیشن لیں گے، ایک ملین سے زیادہ ٹریڈز کو ایک ڈیٹرمنسٹک ایونٹ ٹیپ میں معمول بنائیں گے، اور اسے قابل کنٹرول مارکیٹ کلاک کے ذریعے اس کے اصل وقت کے مطابق دوبارہ چلائیں گے۔
راستے میں، ہم پلے بیک کی رفتار کو ایڈجسٹ کریں گے، کنٹرولز کو روکیں گے اور دوبارہ شروع کریں گے، تلاش کریں گے، اور ایک FastAPI سروس جو WebSockets پر ٹرانزیکشنز کو اسٹریم کرتے ہوئے REST پر کنٹرول کو ظاہر کرتی ہے۔
ہم ایک علیحدہ صارف بھی بنائیں گے جو صرف موصول ہونے والے واقعات پر رولنگ VWAP اور مارکیٹ کی حالت کا حساب لگاتا ہے۔ آخر کار، ہمارے پاس ایک مکمل طور پر مقامی ری پلے سسٹم ہوگا جو پہلے سے مکمل ہونے والے تجارتی دنوں کو ایک ٹائم سٹریم کے طور پر ایونٹ سے چلنے والے سافٹ ویئر کو فیڈ کر سکتا ہے، جبکہ بازیافت کے بعد نیچے کی دھارے کی حالت کو درست طریقے سے دوبارہ بناتا ہے اور خودکار جانچ کے ذریعے نتائج کی توثیق کرتا ہے۔
انڈیکس
شرائط
شروع کرنے سے پہلے، درج ذیل کو چیک کریں:
-
Python 3.10 یا اس سے زیادہ انسٹال ہے۔
-
تاریخی ٹک ڈیٹا اینڈ پوائنٹ تک رسائی کے لیے EODHD API کلید۔ آپ EODHD قیمتوں کے صفحہ پر ایک ڈویلپر اکاؤنٹ بنا سکتے ہیں۔
-
ٹرمینل اور کوڈ ایڈیٹر۔
-
بنیادی Python علم، بشمول فنکشنز، کلاسز، لغات اور پیکجز کے ساتھ کام کرنا۔
-
HTTP اور WebSockets کا بنیادی علم کوئی سابقہ FastAPI تجربہ درکار نہیں ہے۔
-
ڈاؤن لوڈ کردہ خام ٹک ڈیٹا اور پروسیس شدہ پلے بیک ٹیپس کو ذخیرہ کرنے کے لیے کافی مقامی ڈسک کی جگہ۔ اس ٹیوٹوریل میں استعمال ہونے والے مکمل AAPL سیشن میں 10 لاکھ سے زیادہ لین دین کے ریکارڈ شامل ہیں۔
اس ٹیوٹوریل میں شیل کمانڈز یونکس طرز کا نحو استعمال کرتی ہیں، لہذا وہ براہ راست میک او ایس اور لینکس پر کام کرتی ہیں۔ ونڈوز پر، آپ اسے WSL، Git Bash کے ذریعے چلا سکتے ہیں، یا مساوی PowerShell کمانڈز استعمال کر سکتے ہیں۔
ہم کیا بنا رہے ہیں
کسی کوڈ کو چھونے سے پہلے پورے سسٹم پر ایک نظر ڈالنا مفید ہے۔ پلے بیک انجن EODHD سے تاریخی لین دین لیتا ہے، انہیں ایک مستقل اندرونی شکل میں تبدیل کرتا ہے، ٹائمنگ کو بحال کرتا ہے، اور تجارتی دن کے دوبارہ کھلتے ہی انہیں صارفین کے لیے الگ کرتا ہے۔
مکمل بہاؤ مندرجہ ذیل ہے:
ہر پرت کا ایک کام ہوتا ہے۔ لوڈر خام تاریخی سیشنز کو بازیافت اور محفوظ کرتا ہے۔ نارملائزر ریکارڈ کی توثیق کرتا ہے اور اسے ڈیٹرمنسٹک پلے بیک ٹیپ میں تبدیل کرتا ہے۔ گھڑیاں نقشہ ریکارڈنگ ٹائم اسٹیمپ کو وال کلاک ٹائم کے مطابق کرتی ہیں، جبکہ پلے بیک سیشنز کنٹرولز کو شامل کرتے ہیں جیسے کہ اسٹارٹ، موقوف، دوبارہ شروع، رفتار تبدیل کرنا، تلاش کرنا اور رکنا۔
فاسٹ اے پی آئی اس پلے بیک انجن کے ارد گرد بیٹھتا ہے۔ REST اینڈ پوائنٹس کنٹرول ہوائی جہاز بناتے ہیں، اور WebSockets اصل لین دین اور پلے بیک کنٹرول ایونٹس کو لے جاتے ہیں۔ دوسری طرف، صارفین اپنی رولنگ سٹیٹ کو صرف اس سٹریم کے ذریعے پہنچنے والی چیزوں سے برقرار رکھتے ہیں۔
ہم ان ذمہ داریوں کو پراجیکٹ کے ڈھانچے میں الگ رکھیں گے۔
market-time-machine/
├── data/
│ ├── raw/
│ └── processed/
├── replay/
│ ├── __init__.py
│ ├── config.py
│ ├── loader.py
│ ├── events.py
│ ├── clock.py
│ └── session.py
├── api/
│ ├── __init__.py
│ ├── server.py
│ └── run.py
├── consumer/
│ ├── __init__.py
│ └── consumer.py
├── tests/
│ ├── __init__.py
│ ├── conftest.py
│ └── test_replay.py
├── .env
├── .gitignore
└── pytest.ini
پوری تعمیر کے لیے اہم اصول آسان ہے: صارفین کو صرف یہ جاننے کی ضرورت ہے کہ پلے بیک اسٹریم کے ذریعے پہلے سے کیا پہنچ چکا ہے۔ تاریخ کے ٹیپ کو پہلے سے نہیں پڑھنا چاہیے۔ یہ رکاوٹیں ہیں جو پوسٹ نیویگیشن ٹائمنگ، موقوف/دوبارہ شروع کرنے کے رویے، اور ریاست کی تشکیل نو کو صحیح طریقے سے لاگو کرنے کے قابل بناتی ہیں۔
ازگر پروجیکٹ کی ترتیبات
ایک پروجیکٹ ڈائرکٹری بنا کر اور ان پیکیجز کو انسٹال کرکے شروع کریں جنہیں آپ ڈیٹا کی بازیافت، پلے بیک ٹائمنگ، API لیئر، ویب ساکٹ کمیونیکیشن، اور ٹیسٹنگ کے لیے استعمال کریں گے۔
mkdir -p market-time-machine/data/raw
mkdir -p market-time-machine/data/processed
mkdir -p market-time-machine/replay
mkdir -p market-time-machine/api
mkdir -p market-time-machine/consumer
mkdir -p market-time-machine/tests
cd market-time-machine
pip install requests fastapi "uvicorn[standard]" websockets httpx python-dotenv numpy pytest pytest-asyncio
ایک ڈبہ بنانا __init__.py اندر فائل replay, api, consumerاور tests لہذا، ازگر ہر ڈائریکٹری کو ایک پیکج کے طور پر دیکھتا ہے۔
replay/__init__.py
api/__init__.py
consumer/__init__.py
tests/__init__.py
چونکہ ہم ماضی کے لین دین کو EODHD سے درآمد کریں گے۔ .env فائل کو اپنے پروجیکٹ کی جڑ میں محفوظ کریں اور اپنی API کلید کو وہاں محفوظ کریں۔
چونکہ ڈاؤن لوڈ کردہ سیشنز بھی کافی بڑے ہیں، اس لیے آپ کو کسی بھی اسناد یا مقامی مارکیٹ ڈیٹا فائلوں کا ارتکاب نہیں کرنا چاہیے۔ بنانا .gitignore:
.env
data/
__pycache__/
*.pyc
.ipynb_checkpoints/
میمو: اگر آپ کے پاس EODHD API کلید نہیں ہے، تو آپ EODHD ڈویلپر اکاؤنٹ کھول کر آسانی سے حاصل کر سکتے ہیں۔
اس سیٹ اپ کے بعد، آپ کا پروجیکٹ اس طرح نظر آنا چاہیے:
market-time-machine/
├── data/
│ ├── raw/
│ └── processed/
├── replay/
│ └── __init__.py
├── api/
│ └── __init__.py
├── consumer/
│ └── __init__.py
├── tests/
│ └── __init__.py
├── .env
└── .gitignore
کہ raw/ ڈائریکٹری EODHD سے موصول ہونے والے جوابات کو محفوظ رکھتی ہے، لیکن processed/ ہم اپنے بنائے ہوئے عام پلے بیک ٹیپس کو رکھیں گے۔
EODHD سے مکمل تجارتی سیشن ڈاؤن لوڈ کریں۔
پلے بیک انجن کے لیے وقت کا احساس بحال کرنے کے لیے ایک مکمل تجارتی سیشن درکار ہے۔ ہم 15 جولائی 2026 کے لیے AAPL ٹرانزیکشنز کو بازیافت کرنے کے لیے EODHD کا تاریخی ٹک API استعمال کرتے ہیں، لیکن تلاش کی پرت کو پلے بیک سے متعلق کسی بھی چیز سے الگ رکھتے ہیں۔
دو فائلیں پروجیکٹ کے اس حصے کو سنبھالتی ہیں۔
market-time-machine/
└── replay/
├── __init__.py
├── config.py
└── loader.py
config.py مشترکہ APIs، روٹس، اور مارکیٹ سیشن کی ترتیبات کو ایک جگہ پر رکھیں۔ loader.py ہم ان ترتیبات کو سیشن کو بازیافت کرنے اور نیچے دیئے گئے خام جواب کو محفوظ کرنے کے لیے استعمال کرتے ہیں۔ data/raw/.
بنانا replay/config.py
درج ذیل شامل کریں:
import os
from pathlib import Path
from dotenv import load_dotenv
ROOT = Path(__file__).resolve().parent.parent
load_dotenv(ROOT / ".env")
TOKEN = os.environ.get("EODHD_API_TOKEN")
TICKS_URL = "https://eodhd.com/api/ticks/"
RAW = ROOT / "data" / "raw"
PROCESSED = ROOT / "data" / "processed"
MARKET_TZ = "America/New_York"
OPEN = "09:30:00"
CLOSE = "16:00:00"
MAX_LIMIT = 10_000
MIN_WINDOW_S = 1
CLOSE_GRACE_S = 5
FIELDS = ("mkt", "price", "seq", "shares", "sl", "sub_mkt", "ts")
NON_LAST_SALE = frozenset("IWVT47")
def token():
if not TOKEN:
raise RuntimeError("EODHD_API_TOKEN not set")
return TOKEN
def redact(text):
return str(text).replace(TOKEN, "") if TOKEN else str(text)
ایک باقاعدہ امریکی اسٹاک سیشن کی وضاحت اس طرح کی گئی ہے: America/New_York یہ ضروری ہے کیونکہ UTC 09:30 کے مساوی یو ٹی سی ٹائم اسٹیمپ کے بجائے دن کی روشنی کی بچت کے وقت کے ساتھ تبدیل ہوتا ہے۔
ہم درخواست کی مدت کو 16:00 سے 5 سیکنڈ تک بڑھاتے ہیں۔ CLOSE_GRACE_S. اس ٹیوٹوریل میں استعمال ہونے والے سیشن میں ایسی سرگرمی ہوتی ہے جو 16:00:00 کے فوراً بعد ختم ہو جاتی ہے، اس لیے رعایتی مدت اس تاریخ کو ڈاؤن لوڈ کے اندر برقرار رکھتی ہے۔
بنانا replay/loader.py
ایک بڑی درخواست ٹک ڈیٹا کے گھنے سیشن کو بازیافت کرنے کا محفوظ طریقہ نہیں ہے۔ سرگرمی پورے دن میں نمایاں طور پر تبدیل ہوتی ہے، اور کوئی بھی درخواست جو ترتیب شدہ تک پہنچ جاتی ہے۔ 10,000-ریکارڈ کی حدیں کٹے ہوئے وقفوں کی نشاندہی کر سکتی ہیں۔
اس کے بجائے، لوڈر پچھلے جوابات کی کثافت کی بنیاد پر اپنی درخواست کی ونڈو کو ایڈجسٹ کرتا ہے۔
بنانا replay/loader.py:
import json, time
from datetime import datetime
from zoneinfo import ZoneInfo
import requests
from . import config
def fetch(symbol, frm, to, limit=None):
limit = limit or config.MAX_LIMIT
r = requests.get(config.TICKS_URL, timeout=180, params={
"s": symbol,
"from": frm,
"to": to,
"limit": limit,
"api_token": config.token(),
"fmt": "json"
})
if r.status_code != 200:
raise RuntimeError(
f"HTTP {r.status_code} {config.redact(r.text[:200])}"
)
return r.json()
def bounds(date_str, grace=None):
grace = config.CLOSE_GRACE_S if grace is None else grace
tz = ZoneInfo(config.MARKET_TZ)
d = datetime.strptime(date_str, "%Y-%m-%d").date()
def at(hms):
h, m, s = map(int, hms.split(":"))
return datetime(
d.year, d.month, d.day, h, m, s, tzinfo=tz
).timestamp()
return int(at(config.OPEN)), int(at(config.CLOSE)) + grace
def fetch_session(symbol, date_str, tag="session", window=None,
force=False, verbose=True):
raw = config.RAW / f"{symbol}_{date_str}_{tag}.jsonl"
man = config.RAW / f"{symbol}_{date_str}_{tag}.manifest.json"
if raw.exists() and man.exists() and not force:
m = json.loads(man.read_text())
print(f"cached {raw.name}: {m['ticks']:,} ticks")
return m, raw
start, end = window or bounds(date_str)
cursor, win = start, 30
total = pages = retries = 0
first_ts = last_ts = None
seen_fields = set()
t0 = time.perf_counter()
with open(raw, "w") as fh:
while cursor < end:
b = min(cursor + win, end)
span = b - cursor
payload = fetch(symbol, cursor, b)
n = len(payload.get("ts", []))
if n >= config.MAX_LIMIT:
if span <= config.MIN_WINDOW_S:
raise RuntimeError(
f"second {cursor} has >= {config.MAX_LIMIT} ticks "
"and cannot be paginated"
)
win = max(1, span // 2)
retries += 1
continue
if n:
seen_fields.update(payload.keys())
if first_ts is None:
first_ts = payload["ts"][0]
last_ts = payload["ts"][-1]
fh.write(json.dumps({
"from": cursor,
"to": b,
"payload": payload
}) + "n")
total += n
pages += 1
cursor = b
density = n / span if span else 0
win = int(min(
1800,
max(1, config.MAX_LIMIT * 0.75 / max(density, 0.01))
))
if verbose and pages % 20 == 0:
pct = 100 * (cursor - start) / (end - start)
print(f"{pct:5.1f}% {total:,} ticks")
m = {
"symbol": symbol,
"date": date_str,
"tag": tag,
"ticks": total,
"pages": pages,
"retries": retries,
"window_from_utc": start,
"window_to_utc": end,
"first_timestamp_ms": first_ts,
"last_timestamp_ms": last_ts,
"fields": sorted(seen_fields),
"elapsed_s": round(time.perf_counter() - t0, 1),
"api_calls": pages * 10
}
man.write_text(json.dumps(m, indent=2))
return m, raw
def read_pages(path):
with open(path) as fh:
for line in fh:
if line.strip():
yield json.loads(line)
لوڈر 30 سیکنڈ کی مدت کے ساتھ شروع ہوتا ہے۔ اگر وہ وقفہ ریکارڈ کی بالائی حد تک پہنچ جاتا ہے، تو یہ ممکنہ طور پر نامکمل جواب کو قبول کرنے کے بجائے ایک چھوٹے وقفے پر دوبارہ کوشش کرتا ہے۔ پرسکون ادوار کے لیے، یہ ونڈو 30 منٹ تک بڑھ سکتی ہے۔
ہر قبول شدہ جواب کو معمول پر لانے سے پہلے براہ راست JSONL کو لکھا جاتا ہے۔ مینی فیسٹ سیشن کی حدود، ٹک کاؤنٹ، ٹائم اسٹیمپ، مشاہدہ شدہ فیلڈز اور تلاش کے اعدادوشمار کے ساتھ محفوظ کیا جاتا ہے۔
اب مکمل AAPL سیشن حاصل کریں۔
from replay.loader import fetch_session
SYMBOL = "AAPL"
DATE = "2026-07-15"
print("=== fullday ===")
manifest, raw_path = fetch_session(
SYMBOL,
DATE,
tag="fullday"
)
print(
f" window {manifest['window_from_utc']}..{manifest['window_to_utc']} | "
f"{manifest['ticks']:,} ticks, {manifest['pages']} pages, "
f"{manifest['retries']} retries | "
f"{manifest['elapsed_s']}s, {manifest['api_calls']} metered api calls"
)
print(
f" first_ts {manifest['first_timestamp_ms']} "
f"last_ts {manifest['last_timestamp_ms']}"
)
print(" fields:", manifest["fields"])
فائنل کلین اپ رن نے پہلے سے ڈاؤن لوڈ کردہ سیشنز کو دوبارہ استعمال کیا

ہم اب ہیں 1,032,411 خام لین دین کے ریکارڈ جو پورے ریگولر سیشن اور رعایتی مدت کا احاطہ کرتے ہیں۔ کہ 145 وائٹ لسٹ شدہ صفحات اور 16 دوبارہ کوششیں یہ بھی ظاہر کرتی ہیں کہ اس طرح کے گھنے ٹک ڈیٹا کے لیے ایک مقررہ درخواست کی مدت کیوں کمزور مفروضہ تھی۔
تاہم، یہ ریکارڈ بالکل اسی طرح محفوظ کیے جاتے ہیں جیسے وہ EODHD سے حاصل کیے گئے تھے۔ اس سے پہلے کہ کوئی پلے بیک انجن اسے استعمال کر سکے، اسے پہلے واقعات کی ایک تعییناتی داخلی ترتیب بننا چاہیے۔
ٹک ڈیٹا کو پلے بیک ٹیپ میں معمول بنائیں
لوڈر پورے سیشن میں کام کرتا ہے، لیکن پلے بیک انجن کو براہ راست EODHD کے خام رسپانس فارمیٹ کے ساتھ کام نہیں کرنا چاہیے۔ ٹک اینڈ پوائنٹ ایک متوازی صف میں ٹائم اسٹیمپ، قیمت، سائز، ترتیب نمبر، اور مارکیٹ کوڈ جیسے فیلڈز واپس کرتا ہے۔
کھیلنے سے پہلے، آپ کو اس بات کو یقینی بنانا ہوگا کہ صف کو ترتیب دیا گیا ہے، واقعات کی ایک تعییناتی ترتیب قائم کریں، ڈپلیکیٹس کو ہٹا دیں، اور نتائج کو ایک داخلی شکل میں تبدیل کریں۔
منطق مندرجہ ذیل ہے: replay/events.py:
market-time-machine/
└── replay/
├── config.py
├── loader.py
└── events.py
ہم یہاں دو اشیاء استعمال کریں گے۔ TradeEvent بالآخر، یہ ویب ساکٹ پر منتقل ہونے کی صورت میں ایک ہی لین دین کی نمائندگی کرتا ہے۔ TradeTape کالمی NumPy صفوں میں پورے سیشنز کو مؤثر طریقے سے اسٹور کرتا ہے اور انفرادی سیشنز کو عملی شکل دیتا ہے۔ TradeEvent ضرورت کے وقت ہی اشیاء بنائیں۔
بنانا replay/events.py
بنانا replay/events.py کے ساتھ:
from dataclasses import dataclass
import numpy as np
from . import config
from .loader import read_pages
@dataclass(frozen=True)
class TradeEvent:
symbol: str
timestamp_ms: int
price: float
size: int
sequence: int
market: str
sub_market: str
sale_condition: str
source: str = "replay"
def to_wire(self):
sl = self.sale_condition
return {
"type": "trade",
"symbol": self.symbol,
"timestamp_ms": self.timestamp_ms,
"price": self.price,
"size": self.size,
"sequence": self.sequence,
"source": self.source,
"metadata": {
"market": self.market,
"sub_market": self.sub_market or None,
"sale_condition": sl,
"odd_lot": "I" in sl,
"zero_size": self.size == 0,
"last_sale_eligible": not (
set(sl) & config.NON_LAST_SALE
)
}
}
class TradeTape:
def __init__(self, symbol, ts, price, size, seq, mkt, sub, sl):
self.symbol = symbol
self.ts = ts
self.price = price
self.size = size
self.seq = seq
self.mkt = mkt
self.sub = sub
self.sl = sl
def __len__(self):
return len(self.ts)
def __getitem__(self, i):
return TradeEvent(
self.symbol,
int(self.ts[i]),
float(self.price[i]),
int(self.size[i]),
int(self.seq[i]),
str(self.mkt[i]),
str(self.sub[i]),
str(self.sl[i])
)
def index_at(self, ts_ms):
return int(np.searchsorted(self.ts, ts_ms, side="left"))
def span(self):
if not len(self):
return None, None
return int(self.ts[0]), int(self.ts[-1])
def save(self, path):
np.savez_compressed(
path,
ts=self.ts,
price=self.price,
size=self.size,
seq=self.seq,
mkt=self.mkt,
sub=self.sub,
sl=self.sl,
symbol=np.array([self.symbol])
)
@classmethod
def load(cls, path):
z = np.load(path, allow_pickle=False)
return cls(
str(z["symbol"][0]),
z["ts"],
z["price"],
z["size"],
z["seq"],
z["mkt"],
z["sub"],
z["sl"]
)
def normalize(raw_path, symbol, verbose=True):
cols = {k: [] for k in config.FIELDS}
pages = 0
for page in read_pages(raw_path):
pages += 1
p = page["payload"]
lens = {k: len(p.get(k, [])) for k in config.FIELDS}
if len(set(lens.values())) != 1:
raise ValueError(
f"ragged page {page['from']}: {lens}"
)
for k in config.FIELDS:
cols[k].extend(p[k])
ts = np.asarray(cols["ts"], dtype=np.int64)
price = np.asarray(cols["price"], dtype=np.float64)
size = np.asarray(cols["shares"], dtype=np.int64)
seq = np.asarray(cols["seq"], dtype=np.int64)
mkt = np.asarray(cols["mkt"], dtype=str)
sub = np.asarray(cols["sub_mkt"], dtype=str)
sl = np.asarray(cols["sl"], dtype=str)
raw_n = len(ts)
def arrays(mask):
return tuple(
a[mask]
for a in (ts, price, size, seq, mkt, sub, sl)
)
keep = (
(ts > 0)
& np.isfinite(price)
& (price > 0)
& (size >= 0)
)
ts, price, size, seq, mkt, sub, sl = arrays(keep)
order = np.lexsort((seq, ts))
ts, price, size, seq, mkt, sub, sl = arrays(order)
dup = np.zeros(len(ts), dtype=bool)
if len(ts) > 1:
dup[1:] = (
(ts[1:] == ts[:-1])
& (seq[1:] == seq[:-1])
)
ts, price, size, seq, mkt, sub, sl = arrays(~dup)
tape = TradeTape(
symbol,
ts,
price,
size,
seq,
mkt,
sub,
sl
)
odd = sum("I" in str(s) for s in sl)
elig = sum(
not (set(str(s)) & config.NON_LAST_SALE)
for s in sl
)
rep = {
"pages": pages,
"raw": raw_n,
"kept": len(ts),
"dropped": raw_n - len(ts) - int(dup.sum()),
"dupes": int(dup.sum()),
"zero_size": int((size == 0).sum()),
"odd_lot": int(odd),
"last_sale_eligible": int(elig),
"seq_strict": bool(
np.all(seq[1:] > seq[:-1])
) if len(seq) > 1 else True,
"span": tape.span()
}
if verbose:
n = max(1, len(ts))
print(
f"{rep['raw']:,} raw -> {rep['kept']:,} kept "
f"({rep['dupes']} dupes, {rep['dropped']} invalid)"
)
print(
f"zero-size {100*rep['zero_size']/n:.1f}% | "
f"odd-lot {100*odd/n:.1f}% | "
f"last-sale-eligible {100*elig/n:.1f}%"
)
print(
f"seq strictly increasing: {rep['seq_strict']}"
)
return tape, rep
پہلی توثیق ٹرانزیکشن کی تشکیل سے پہلے ہوتی ہے۔ سورس فیلڈز متوازی صفوں میں آتے ہیں، لہذا صفحہ پر موجود تمام فیلڈز میں مشاہدات کی ایک ہی تعداد ہونی چاہیے۔ بصورت دیگر، ان کو یکجا کرنے سے ایک ٹرانزیکشن کی قیمت خود بخود دوسرے ٹرانزیکشن کے ٹائم اسٹیمپ کے ساتھ منسلک ہو سکتی ہے۔
اس کے بعد، صف کو NumPy میں تبدیل کر دیا جاتا ہے، بنیادی غلط ریکارڈز کو ہٹا دیا جاتا ہے، اور لین دین کو اس طرح ترتیب دیا جاتا ہے: (timestamp, sequence). ٹائم اسٹیمپ تاریخ کی ترتیب فراہم کرتے ہیں، جب کہ ترتیب نمبر جب ایک سے زیادہ لین دین ایک ہی ملی سیکنڈ کا اشتراک کرتے ہیں تو تعییناتی ترتیب فراہم کرتے ہیں۔
ایک ہی ٹائم اسٹیمپ اور ترتیب کے ساتھ صحیح ڈپلیکیٹس کو پھر ہٹا دیا جاتا ہے۔ TradeTape ایک ملین سے زیادہ مستقل پائتھون اشیاء مختص کرنے کے بجائے، میں نتیجے میں آنے والے کالموں کو صفوں کے طور پر رکھتا ہوں، جو میرے سیشنز کو دن بھر میموری پر کافی حد تک ہلکا رکھتا ہے۔
اب خام سیشن کو معمول پر لائیں اور اسے نیچے محفوظ کریں۔ data/processed/:
import json
from replay import config
from replay.events import normalize
tape, report = normalize(raw_path, SYMBOL)
tape.save(
config.PROCESSED / f"{SYMBOL}_{DATE}_fullday.npz"
)
lo, hi = tape.span()
print(
f"span {lo}..{hi} "
f"({(hi-lo)/3_600_000:.2f} market hours)"
)
print("first 3 normalized events:")
for i in range(3):
print(json.dumps(tape[i].to_wire()))
اصل نارملائزیشن رن مندرجہ ذیل ہے:

ایک ملین سے زیادہ خام مشاہدات، صرف دو ڈپلیکیٹ ریکارڈز کھو گئے ہیں، اور کوئی بھی ریکارڈ ڈیفالٹ ٹائم اسٹیمپ، قیمت، یا سائز کی جانچ میں ناکام نہیں ہوتا ہے۔ ری پلے کے لیے، جو چیز زیادہ اہم ہے وہ یہ ہے کہ معمول کے مطابق ترتیب کو سختی سے بڑھایا گیا ہے۔
بینچ مارکنگ اور جانچ کے لیے چھوٹے ٹیپ بنائیں
سارا دن، ٹیپ آپ کے آخری پلے بیک کو طاقت دے گی۔ تاہم، ٹائمنگ بینچ مارکس اور خودکار ٹیسٹوں کے لیے، آپ کو ہر بار پورے 6.5 گھنٹے چلانے کی ضرورت نہیں ہے۔
آئیے 12:00 AM سے 12:15 AM تک 15 منٹ کا ٹکڑا براہ راست ایک عام پورے دن کے ٹیپ سے نکالیں۔
from datetime import datetime
from zoneinfo import ZoneInfo
from replay.events import TradeTape
tz = ZoneInfo(config.MARKET_TZ)
quiet_start = int(
datetime(
2026, 7, 15, 12, 0,
tzinfo=tz
).timestamp() * 1000
)
quiet_end = quiet_start + 15 * 60_000
i = tape.index_at(quiet_start)
j = tape.index_at(quiet_end)
quiet_tape = TradeTape(
tape.symbol,
tape.ts[i:j],
tape.price[i:j],
tape.size[i:j],
tape.seq[i:j],
tape.mkt[i:j],
tape.sub[i:j],
tape.sl[i:j]
)
quiet_tape.save(config.PROCESSED / f"{SYMBOL}_{DATE}_quiet15m.npz")
اب آپ کے پاس دو پروسیس شدہ ٹیپ ہیں۔ ایک اینڈ ٹو اینڈ ری پلے کے لیے ایک مکمل سیشن ہے، اور دوسرا ریپبل ایبل ٹائمنگ اور کنٹرول ٹیسٹنگ کے لیے حقیقی دنیا کی مارکیٹ کا چھوٹا وقفہ ہے۔
ایک ریکارڈ پلے بیک گھڑی بنانا
اب ہمارے پاس ایک تعییناتی تجارتی ترتیب ہے، لیکن اب بھی ایسا کچھ نہیں ہے جو ان تجارتوں کو بازار کے بہاؤ کی طرح برتاؤ کرے۔ بس ٹیپ کو لوپ کریں اور Python سیشن کو اتنی ہی تیزی سے پروسیس کرے گا جتنی آپ کا سسٹم اجازت دیتا ہے۔
ری پلے کلاک تاریخی بازار کے اوقات کو دیوار کی گھڑی کے حقیقی اوقات میں نقشہ بنا کر اس مسئلے کو حل کرتی ہیں۔ آپ اصل ٹائم اسٹیمپ کو تبدیل کیے بغیر پلے بیک کی رفتار بھی تبدیل کر سکتے ہیں۔
سادہ ورژن تمام تجارتی جوڑوں کے درمیان تاریخی خلا پر سو سکتا ہے۔
gap = (next_ts - current_ts) / 1000
await asyncio.sleep(gap / speed)
کو 10x500ms کا تاریخی وقفہ 50ms ہوگا۔ کو 100xیہ 5ms ہو گا۔
مسئلہ یہ ہے۔ asyncio.sleep() یہ صرف اس بات کی ضمانت دیتا ہے کہ پھانسی دوبارہ شروع ہو جائے گی۔ ~ بعد تاخیر کی درخواست کی۔ اگر ہر نیند کے بعد تھوڑی دیر بعد بیداری ہوتی ہے، اور اگلی تاخیر کو دیر سے بیداری کے وقت سے ماپا جاتا ہے، تو یہ خرابیاں طویل پلے بیکس میں جمع ہو سکتی ہیں۔
اس کے بجائے، میں یہاں مکمل ری پلے پن کروں گا: time.monotonic():
historical elapsed time
÷
replay speed
+
wall-clock start
=
target wall-clock time
اس لیے، تمام تقریبات ایک ہی اینکر کی بنیاد پر طے کی جاتی ہیں بجائے اس کے کہ پچھلا ایونٹ مکمل ہونے پر۔
بنانا replay/clock.py
اپنے پلے بیک پیکیج میں گھڑی شامل کریں۔
market-time-machine/
└── replay/
├── config.py
├── loader.py
├── events.py
└── clock.py
بنانا replay/clock.py:
import asyncio, time
import numpy as np
MIN_SLEEP_S = 0.0005
class ReplayClock:
def __init__(self, start_ms, speed=1.0):
self.speed = float(speed)
self._anchor_ms = float(start_ms)
self._anchor_wall = None
self.running = False
self.epoch = 0
def start(self):
self._anchor_wall = time.monotonic()
self.running = True
return self
def now_ms(self, now=None):
if not self.running or self._anchor_wall is None:
return self._anchor_ms
now = now if now is not None else time.monotonic()
return (
self._anchor_ms
+ (now - self._anchor_wall) * 1000 * self.speed
)
def wall_for(self, ms):
return (
self._anchor_wall
+ (ms - self._anchor_ms) / 1000 / self.speed
)
def _reanchor(self, ms):
self._anchor_ms = float(ms)
self._anchor_wall = time.monotonic()
self.epoch += 1
def set_speed(self, speed):
self._reanchor(self.now_ms())
self.speed = float(speed)
def pause(self):
if self.running:
self._anchor_ms = self.now_ms()
self.running = False
def resume(self):
if not self.running:
self._anchor_wall = time.monotonic()
self.running = True
self.epoch += 1
def seek(self, ms):
self._reanchor(ms)
def new_stats(speed):
return {
"emitted": 0,
"batches": 0,
"lateness": [],
"dropped": 0,
"speed": speed,
"wall0": None,
"market0": None,
"market1": None
}
def summarize(st):
if not st["lateness"]:
return {
"emitted": st["emitted"],
"batches": st["batches"]
}
a = np.asarray(st["lateness"])
wall = (
time.monotonic() - st["wall0"]
if st["wall0"] else 0.0
)
mkt = (
(st["market1"] - st["market0"]) / 1000
if st["market0"] is not None else 0.0
)
ok = st["dropped"] == 0 and wall > 0
realized = round(mkt / wall, 2) if ok else None
return {
"emitted": st["emitted"],
"batches": st["batches"],
"mean_batch": round(
st["emitted"] / max(1, st["batches"]), 1
),
"market_s": round(mkt, 3),
"wall_s": round(wall, 3),
"requested_speed": st["speed"],
"realized_speed": realized,
"speed_error_pct": (
round(
100 * (realized - st["speed"]) / st["speed"],
2
)
if ok else None
),
"lateness_p50_ms": round(
float(np.percentile(a, 50)), 2
),
"lateness_p95_ms": round(
float(np.percentile(a, 95)), 2
),
"lateness_max_ms": round(
float(a.max()), 2
),
"reanchor_batches_dropped": st["dropped"]
}
async def replay_batches(
tape,
clock,
start,
stats,
max_batch=4096
):
i, n = start, len(tape)
last_epoch = clock.epoch
if stats["wall0"] is None:
stats["wall0"] = time.monotonic()
stats["market0"] = int(tape.ts[start])
while i < n:
if not clock.running:
await asyncio.sleep(0.005)
continue
now = time.monotonic()
j = min(
int(
np.searchsorted(
tape.ts,
clock.now_ms(now),
side="right"
)
),
n,
i + max_batch
)
if j > i and not clock.running:
continue
if j > i:
if clock.epoch == last_epoch:
targets = clock.wall_for(
tape.ts[i:j].astype(np.float64)
)
stats["lateness"].extend(
((now - targets) * 1000).tolist()
)
else:
stats["dropped"] += 1
last_epoch = clock.epoch
stats["emitted"] += j - i
stats["batches"] += 1
stats["market1"] = int(tape.ts[j - 1])
yield i, j
i = j
continue
wait = clock.wall_for(float(tape.ts[i])) - now
await asyncio.sleep(
wait if wait > MIN_SLEEP_S else 0
)
now_ms() یہ آپ کو بتاتا ہے کہ موجودہ کھیل تاریخی بازار کے اوقات میں کہاں ہے۔ wall_for() یہ ریورس کنورژن کرتا ہے اور مشین کی مونوٹونک کلاک کو بتاتا ہے کہ ہسٹری ٹائم اسٹیمپ کب ختم ہونا چاہیے۔
آپ بنیادی ٹیپ میں ترمیم کیے بغیر اس میپنگ کو روک سکتے ہیں، دوبارہ شروع کر سکتے ہیں، رفتار کو تبدیل کر سکتے ہیں اور دوبارہ منجمد کرنے کی کوشش کر سکتے ہیں۔
ایک اور اہم حصہ بیچ پروسیسنگ ہے۔ اعلی ریفریش ریٹ پر، ہر لین دین کے لیے ایک نیند کا وقت مقرر کرنے سے اہم اوور ہیڈ ہوتا ہے۔ replay_batches() اس کے بجائے، یہ پوچھتا ہے کہ مارکیٹ کا وقت کس حد تک آگے بڑھا ہے اور کنفیگر شدہ بیچ سائز تک تمام پہلے سے بند تجارت جاری کرتا ہے۔
اگر ایونٹ کا لوپ تھوڑا سا پیچھے ہو جاتا ہے، تو اگلا بیچ ایک اور مصنوعی تاخیر متعارف کرانے کے بجائے بڑا ہو گا۔
واچ بینچ مارکنگ کھیلیں
اب پچھلے حصے میں بنائی گئی دوپہر کی ٹیپ لوڈ کریں اور مارکیٹ ٹائم کے پہلے 30 سیکنڈز کی جانچ کریں۔
import numpy as np
from replay import config
from replay.events import TradeTape
from replay.clock import (
ReplayClock,
replay_batches,
new_stats,
summarize
)
tape = TradeTape.load(
config.PROCESSED / "AAPL_2026-07-15_quiet15m.npz"
)
end = int(
np.searchsorted(
tape.ts,
tape.ts[0] + 30_000,
side="right"
)
)
print(
f"{end:,} events in the first "
"30 market seconds of AAPL quiet15mn"
)
async def measure():
print(
f"{'speed':>6} {'market_s':>9} "
f"{'wall_s':>8} {'realized':>9} "
f"{'err_%':>7} {'p50_ms':>7} "
f"{'p95_ms':>7} {'max_ms':>7}"
)
for speed in [1, 10, 50, 100]:
clock = ReplayClock(
tape.ts[0],
speed
).start()
st = new_stats(speed)
async for i, j in replay_batches(
tape,
clock,
0,
st
):
if j >= end:
break
r = summarize(st)
print(
f"{r['requested_speed']:>6} "
f"{r['market_s']:>9} "
f"{r['wall_s']:>8} "
f"{r['realized_speed']:>9} "
f"{r['speed_error_pct']:>7} "
f"{r['lateness_p50_ms']:>7} "
f"{r['lateness_p95_ms']:>7} "
f"{r['lateness_max_ms']:>7}"
)
await measure()
عمل درآمد کے حقیقی نتائج درج ذیل ہیں:

تاریخی بازار کے اوقات میں 30 سیکنڈ لگے۔ 30.001 1x موم بتی، 3.001 10x پر سیکنڈ، 0.6 50x پر سیکنڈ، 0.3 سیکنڈ میں 100x لہذا، حاصل کردہ رفتار ہماری درخواست کے بہت قریب رہی۔
تاخیر کی قیمت آپ کو بتاتی ہے کہ شیڈولر نے انفرادی ایونٹ کی آخری تاریخ کتنی چھوٹ دی ہے۔ مثال کے طور پر، 10x پر اوسط تاخیر کا وقت ہے: 0.36 ms95واں پرسنٹائل ہے۔ 1.12 msاس رن سے بدترین مشاہدہ ہے۔ 11.75 ms.
یہ نمبر خود ریفریش گھڑی کی پیمائش کرتا ہے۔ یہ اینڈ ٹو اینڈ ویب ساکٹ لیٹینسی پیمائش نہیں ہے، جو اب بھی ازگر میں ایونٹ لوپس کے لیے بہترین شیڈولنگ ہے، نہ کہ ایکسچینج گریڈ ٹائمنگ۔
پلے بیک سیشنز میں پلے بیک کنٹرولز شامل کریں۔
ری پلے کلاک کو معلوم ہوتا ہے کہ ڈیل کب ہے، لیکن یہ نہیں جانتی کہ موجودہ ری پلے کہاں ہے یا ری پلے کیا جانا چاہیے۔ ٹیپ کے مالک ہونے، موجودہ کرسر کو ٹریک کرنے، ایونٹ کی قطار کو منظم کرنے، اور کنٹرولز کو مربوط کرنے جیسے کہ شروع، توقف، دوبارہ شروع، رفتار کو تبدیل کرنے، تلاش کرنے اور رکنے کے لیے ایک اور پرت کی ضرورت ہے۔
منطق مندرجہ ذیل ہے: replay/session.py:
market-time-machine/
└── replay/
├── config.py
├── loader.py
├── events.py
├── clock.py
└── session.py
اس فرق کو واضح رکھنا مفید ہے کہ گھڑی وقت کی مالک ہے اور سیشن ریاست کا مالک ہے۔
ایک پلے بیک سیشن ریاستوں کے ایک چھوٹے سیٹ سے گزرتا ہے:

بنانا replay/session.py
بنانا replay/session.py:
import asyncio, collections, contextlib, uuid
from enum import Enum
from .clock import ReplayClock, replay_batches, new_stats, summarize
class State(str, Enum):
CREATED, RUNNING, PAUSED, COMPLETED, STOPPED = (
"created", "running", "paused", "completed", "stopped"
)
class ReplaySession:
PRIORITY = {
"paused", "resumed", "speed_changed",
"replay_reset", "session_stopped"
}
def __init__(self, tape, speed=1.0, warmup_ms=120_000, maxsize=256):
self.id = uuid.uuid4().hex[:12]
self.tape = tape
self.warmup_ms = warmup_ms
self.maxsize = maxsize
self.state = State.CREATED
self.cursor = 0
self.clock = ReplayClock(tape.ts[0], speed)
self.stats = new_stats(speed)
self._q = collections.deque()
self._wake = asyncio.Event()
self._task = None
self._epoch = 0
self._lock = asyncio.Lock()
def info(self):
lo, hi = self.tape.span()
return {
"session_id": self.id,
"symbol": self.tape.symbol,
"state": self.state.value,
"speed": self.clock.speed,
"cursor": self.cursor,
"total_events": len(self.tape),
"market_ts_ms": int(
self.tape.ts[min(self.cursor, len(self.tape)-1)]
),
"session_start_ms": lo,
"session_end_ms": hi,
"queued": len(self._q)
}
def _ctrl(self, kind, **kw):
msg = {
"type": kind,
"session_id": self.id,
"source": "replay",
**kw
}
if kind in self.PRIORITY:
self._q.appendleft(msg)
else:
self._q.append(msg)
self._wake.set()
async def _put(self, msg):
while len(self._q) >= self.maxsize:
self._wake.set()
await asyncio.sleep(0)
self._q.append(msg)
self._wake.set()
async def _kill(self):
t, self._task = self._task, None
if t and not t.done():
t.cancel()
with contextlib.suppress(
asyncio.CancelledError,
Exception
):
await t
async def start(self):
self.clock.start()
self.state = State.RUNNING
self._task = asyncio.create_task(self._run())
self._ctrl(
"session_started",
info=self.info()
)
return self.info()
async def pause(self):
if self.state is State.RUNNING:
async with self._lock:
self.clock.pause()
self.stats["dropped"] += 1
self.state = State.PAUSED
self._ctrl(
"paused",
market_ts_ms=self.info()["market_ts_ms"]
)
return self.info()
async def resume(self):
if self.state is State.PAUSED:
async with self._lock:
self.clock.resume()
self.state = State.RUNNING
if self._task is None or self._task.done():
self._task = asyncio.create_task(self._run())
self._ctrl(
"resumed",
market_ts_ms=self.info()["market_ts_ms"]
)
return self.info()
async def set_speed(self, speed):
async with self._lock:
old = self.clock.speed
self.clock.set_speed(speed)
self.stats["speed"] = speed
self._ctrl(
"speed_changed",
old_speed=old,
new_speed=speed
)
return self.info()
async def seek(self, target_ms):
was = self.state
await self._kill()
async with self._lock:
idx = max(
0,
min(
self.tape.index_at(target_ms),
len(self.tape)-1
)
)
self._epoch += 1
self.cursor = idx
self.state = State.PAUSED
self.clock.pause()
warm = max(
0,
self.tape.index_at(
int(self.tape.ts[idx]) - self.warmup_ms
)
)
purged = sum(
1 for m in self._q
if m.get("type") == "trade"
)
self._q = collections.deque(
m for m in self._q
if m.get("type") != "trade"
)
self._ctrl(
"replay_reset",
reason="seek",
target_timestamp_ms=int(self.tape.ts[idx]),
warmup_from_ms=int(self.tape.ts[warm]),
warmup_events=idx-warm,
purged_stale_events=purged,
epoch=self._epoch
)
for k in range(warm, idx):
await self._put({
**self.tape[k].to_wire(),
"warmup": True
})
self._ctrl(
"warmup_complete",
market_ts_ms=int(self.tape.ts[idx])
)
async with self._lock:
self.clock.seek(float(self.tape.ts[idx]))
if was is State.RUNNING:
self.clock.start()
self.state = State.RUNNING
self._task = asyncio.create_task(self._run())
return self.info()
async def stop(self):
self.state = State.STOPPED
await self._kill()
self._ctrl(
"session_stopped",
info=self.info(),
timing=summarize(self.stats)
)
return self.info()
async def _run(self):
epoch = self._epoch
async for i, j in replay_batches(
self.tape,
self.clock,
self.cursor,
self.stats
):
if self._epoch != epoch or self.state is State.STOPPED:
return
for k in range(i, j):
await self._put(self.tape[k].to_wire())
self.cursor = k+1
if self._epoch == epoch and self.cursor >= len(self.tape):
self.state = State.COMPLETED
self._ctrl(
"session_completed",
info=self.info(),
timing=summarize(self.stats)
)
async def events(self):
while True:
if not self._q:
self._wake.clear()
await self._wake.wait()
continue
m = self._q.popleft()
yield m
if m.get("type") in (
"session_completed",
"session_stopped"
):
return
سیشن ریاست کے اہم حصے ہیں: cursorیہ درج ذیل مقام کی طرف اشارہ کرتا ہے: TradeTape. پروڈیوسروں کی طرف سے استعمال کیا جاتا ہے replay_batches() گھڑی کی پرت ہر ٹیپ کی پوزیشن کو وائر ٹرانزیکشن ایونٹ میں تبدیل کرتی ہے اور اسے سیشن کی قطار میں رکھتی ہے۔
موقوف کرسر کو تبدیل کیے بغیر گھڑی کو منجمد کر دیتا ہے۔ دوبارہ کھلنے پر، گھڑی کو ایک نیا وال کلاک اینکر دیا جائے گا اور اسی تاریخی مقام پر جاری رہے گا۔ رفتار کی تبدیلیاں بھی اسی طرح کام کرتی ہیں۔ گھڑی پہلے موجودہ پلے بیک ٹائم اسٹیمپ پر لاک کرتی ہے اور پھر اس مقام سے نئی رفتار کا اطلاق کرتی ہے۔
قطاروں میں صرف لین دین سے زیادہ شامل ہیں۔ کنٹرول جیسے paused, resumed, speed_changedاور replay_reset اس کا مطلب ہے کہ ڈاؤن اسٹریم صارفین لین دین کے ٹائم اسٹیمپ سے اندازہ لگانے کے بجائے پلے بیک حالت میں ہونے والی تبدیلیوں پر رد عمل ظاہر کر سکتے ہیں۔
seek() یہ سب سے زیادہ متعلقہ کنٹرولز ہیں۔ موجودہ پروڈیوسر کو روکتا ہے اور مطلوبہ مقام کا پتہ لگاتا ہے۔ TradeTape.index_at()پرانے اسٹینڈ بائی لین دین کو ہٹاتا ہے اور پلے بیک کو جاری رکھنے سے پہلے وارم اپ مدت کے لیے تیار کرتا ہے۔ آئیے دیکھتے ہیں کہ ریاستی صارفین کو تعینات کرنے کے بعد انہیں گرم کرنے کی ضرورت کیوں ہے۔
آپریشن میں کوئی الگ ٹرمینل نہیں ہے۔ ReplaySession اس مقام پر۔ ایک بار جب باقی سسٹم منسلک ہو جاتا ہے، تو ہم ان کنٹرولز کو اصل API اور WebSocket اسٹریمز کے ذریعے انجام دیں گے۔
فاسٹ اے پی آئی اور ویب ساکٹ کا استعمال کرتے ہوئے ری پلے کو بے نقاب کرنا
پلے بیک سیشن میں اب وہ سب کچھ ہوتا ہے جس کی آپ کو ریکارڈنگ پلے بیک کو کنٹرول کرنے کی ضرورت ہوتی ہے، لیکن یہ اب بھی صرف Python آبجیکٹ کے طور پر موجود ہے۔ یہ اپنے ارد گرد ایک چھوٹی API پرت رکھتا ہے تاکہ دوسرے پروگرام سیشنز بنا اور کنٹرول کر سکیں اور نتیجے میں لین دین کا سلسلہ حاصل کر سکیں۔
یہ الگ الگ تہوں پر ہیں۔ api/ پیکیج:
market-time-machine/
├── replay/
│ └── ...
└── api/
├── __init__.py
├── server.py
└── run.py
ہم دو مواصلاتی راستے استعمال کرتے ہیں: REST اینڈ پوائنٹس کنٹرول ہوائی جہاز بناتے ہیں اور ایک مستقل WebSocket ایونٹ کے سلسلے کو لے جاتا ہے۔
Control plane
POST /sessions
POST /sessions/{id}/start
POST /sessions/{id}/pause
POST /sessions/{id}/resume
POST /sessions/{id}/speed
POST /sessions/{id}/seek
POST /sessions/{id}/stop
Event stream
WS /sessions/{id}/stream
لہذا، روکنے یا تلاش کرنے جیسی کمانڈز HTTP پر پہنچتی ہیں، جب کہ لین دین اور پلے بیک کنٹرول کے واقعات WebSocket کے ذریعے صارفین تک پہنچتے رہتے ہیں۔
بنانا api/server.py
بنانا api/server.py:
from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect
from pydantic import BaseModel, Field
from replay import config
from replay.events import TradeTape
from replay.session import ReplaySession
from replay.clock import summarize
app = FastAPI(title="Market Time Machine")
SESSIONS = {}
ATTACHED = set()
class Create(BaseModel):
symbol: str = "AAPL"
date: str
tag: str = "fullday"
speed: float = Field(1.0, gt=0)
warmup_ms: int = 120_000
class Speed(BaseModel):
speed: float = Field(..., gt=0)
class Seek(BaseModel):
target_timestamp_ms: int
def get(sid):
if sid not in SESSIONS:
raise HTTPException(404, f"no session {sid}")
return SESSIONS[sid]
@app.post("/sessions")
async def create(b: Create):
path = config.PROCESSED / f"{b.symbol}_{b.date}_{b.tag}.npz"
if not path.exists():
raise HTTPException(404, f"no tape {path.name}")
s = ReplaySession(
TradeTape.load(path),
b.speed,
b.warmup_ms
)
SESSIONS[s.id] = s
return s.info()
@app.get("/sessions/{sid}")
async def info(sid: str):
return get(sid).info()
@app.get("/sessions/{sid}/timing")
async def timing(sid: str):
return summarize(get(sid).stats)
@app.post("/sessions/{sid}/start")
async def start(sid: str):
return await get(sid).start()
@app.post("/sessions/{sid}/pause")
async def pause(sid: str):
return await get(sid).pause()
@app.post("/sessions/{sid}/resume")
async def resume(sid: str):
return await get(sid).resume()
@app.post("/sessions/{sid}/stop")
async def stop(sid: str):
return await get(sid).stop()
@app.post("/sessions/{sid}/speed")
async def speed(sid: str, b: Speed):
return await get(sid).set_speed(b.speed)
@app.post("/sessions/{sid}/seek")
async def seek(sid: str, b: Seek):
return await get(sid).seek(b.target_timestamp_ms)
@app.websocket("/sessions/{sid}/stream")
async def stream(ws: WebSocket, sid: str):
await ws.accept()
if sid not in SESSIONS:
return await ws.close(4004, "unknown session")
if sid in ATTACHED:
return await ws.close(4009, "consumer already attached")
ATTACHED.add(sid)
try:
await ws.send_json({
"type": "attached",
"session_id": sid
})
async for msg in SESSIONS[sid].events():
await ws.send_json(msg)
except (WebSocketDisconnect, Exception):
pass
finally:
ATTACHED.discard(sid)
جب آپ سیشن بناتے ہیں تو پروسیس شدہ مواد لوڈ ہو جاتا ہے۔ .npz اسے ٹیپ سے لپیٹنے کے بعد ReplaySession. اس وقت، API کو EODHD کو کال کرنے یا خام JSONL جواب کو دوبارہ پڑھنے کی ضرورت نہیں ہے۔ عام ٹیپ پر پلے بیک مکمل طور پر فعال ہے۔
REST ہینڈلرز کو جان بوجھ کر پتلا رکھا جاتا ہے۔ /pauseمثال کے طور پر، اس میں اپنی توقف کی منطق نہیں ہے۔
@app.post("/sessions/{sid}/pause")
async def pause(sid: str):
return await get(sid).pause()
یہ صرف حکم دیتا ہے. ReplaySession. یہی پیٹرن دوبارہ شروع کرنے، رفتار کو تبدیل کرنے، تلاش کرنے اور روکنے پر لاگو ہوتا ہے۔ یہ اندرونی طور پر پلے بیک رویے کو محفوظ رکھتا ہے۔ replay/ فاسٹ اے پی آئی کو پابند کرنے کے بجائے۔
WebSocket اینڈ پوائنٹ دوسری سمت کو سنبھالتے ہیں۔ جب کوئی صارف جڑتا ہے، سرور جو کچھ بھی صارف پیدا کرتا ہے اسے آگے بھیج دیتا ہے۔ session.events():
async for msg in SESSIONS[sid].events():
await ws.send_json(msg)
یہ ایک عام لین دین ہو سکتا ہے۔
{
"type": "trade",
"symbol": "AAPL",
"timestamp_ms": 1784122200009,
"price": 317.46,
"size": 3,
"sequence": 61530328,
"source": "replay"
}
یا پلے بیک کنٹرول پیغام:
{
"type": "paused",
"market_ts_ms": 1784122200009
}
نیویگیشن بعد میں ایک اور اہم کنٹرول ایونٹ متعارف کراتی ہے۔
{
"type": "replay_reset",
"reason": "seek",
"target_timestamp_ms": 1784136600030
}
سرور فی پلے بیک سیشن ایک WebSocket صارف کی اجازت دیتا ہے۔ فی الحال، قطار میں لگانا براڈکاسٹ سسٹم کے بجائے ایک FIFO ہینڈ آف ہے، لہذا ایک سے زیادہ صارفین کو ایک ہی سیشن سے جوڑنے سے وہ ایونٹس کو تقسیم کرنے کا سبب بنے گا بجائے اس کے کہ ہر ایک کو مکمل سلسلہ ملے۔
بنانا api/run.py
دوسری API فائل آسانی سے FastAPI ایپلیکیشن چلاتی ہے۔
بنانا api/run.py:
import argparse
import uvicorn
from api.server import app
if __name__ == "__main__":
p = argparse.ArgumentParser()
p.add_argument("--port", type=int, default=8765)
a = p.parse_args()
uvicorn.run(
app,
host="127.0.0.1",
port=a.port,
log_level="warning"
)
پروجیکٹ روٹ میں سروس شروع کریں۔
python -m api.run --port 8765
پلے بیک انجن میں اب ایک بیرونی کنٹرول انٹرفیس اور WebSocket ایونٹ اسٹریم ہے۔ اگلا حصہ اس ندی کے دوسری طرف پروگرام ہے۔ یعنی وہ صارفین جو صرف ان واقعات کے ذریعے مارکیٹ سٹیٹس بناتے ہیں جو انہیں موصول ہوتے ہیں۔
ایک ریاستی WebSocket صارف بنانا
پلے بیک سروسز اب تاریخی لین دین کو سٹریم کر سکتی ہیں، لیکن انہیں ابھی بھی WebSocket کے دوسری طرف سے کسی چیز کی ضرورت ہے تاکہ وہ ایک حقیقی بہاو ایپلی کیشن کی طرح کام کرے۔
صارف کو ریکارڈنگ ٹیپس کو لوڈ نہیں کرنا چاہئے یا EODHD کو براہ راست کال نہیں کرنا چاہئے۔ پلے بیک سٹریم کے ذریعے آنے والے پیغامات سے مارکیٹ کا ایک جامع نقطہ نظر آنا چاہیے۔
ہم اسے الگ پیکج میں رکھیں گے:
market-time-machine/
├── replay/
│ └── ...
├── api/
│ └── ...
└── consumer/
├── __init__.py
└── consumer.py
اس ٹیوٹوریل میں، صارف برقرار رکھتا ہے:
-
تازہ ترین سودے
-
حالیہ آخری فروخت کوالیفائنگ لین دین
-
مجموعی حجم
-
30 سیکنڈ VWAP
-
2 منٹ VWAP
-
طاق اور صفر کی شدت کے فیصد
-
سادہ
SHORT_ABOVE/SHORT_BELOWصورت حال
اختتامی ریاست تجارتی حکمت عملی نہیں ہے۔ ہمیں پلے بیک کنٹرولز کی ضرورت ہے، خاص طور پر نیویگیشن، حقیقی معنوں میں بیان کرنے کے لیے تاکہ ہم بعد میں اس بات کی تصدیق کر سکیں کہ ہم صارفین کو بازاری ریکارڈ کے ساتھ نہیں چھوڑ رہے ہیں۔
بنانا consumer/consumer.py
بنانا consumer/consumer.py:
import argparse, asyncio, collections, json
import websockets
class VWAP:
def __init__(self, window_ms):
self.w = window_ms
self.buf = collections.deque()
self.pv = 0.0
self.vol = 0.0
def add(self, ts, px, sz):
self.buf.append((ts, px, sz))
self.pv += px * sz
self.vol += sz
cut = ts - self.w
while self.buf and self.buf[0][0] < cut:
_, p, s = self.buf.popleft()
self.pv -= p * s
self.vol -= s
if self.vol <= 0:
self.pv = self.vol = 0.0
@property
def value(self):
return self.pv / self.vol if self.vol > 0 else None
class State:
def __init__(self, short_ms=30_000, long_ms=120_000):
self.short = VWAP(short_ms)
self.long = VWAP(long_ms)
self.last_trade = None
self.last_sale = None
self.signal = None
self.n = 0
self.vol = 0
self.odd = 0
self.zero = 0
self.warming = False
def apply(self, m):
ts = m["timestamp_ms"]
px = m["price"]
sz = m["size"]
meta = m["metadata"]
self.short.add(ts, px, sz)
self.long.add(ts, px, sz)
self.last_trade = px
if meta["last_sale_eligible"]:
self.last_sale = px
self.n += 1
self.vol += sz
self.odd += meta["odd_lot"]
self.zero += meta["zero_size"]
s = self.short.value
l = self.long.value
if s is not None and l is not None:
self.signal = (
"SHORT_ABOVE"
if s > l
else "SHORT_BELOW"
)
def line(self):
f = lambda v: "--" if v is None else f"{v:.4f}"
return (
f"n={self.n:>7,} "
f"vol={self.vol:>9,} "
f"trade={f(self.last_trade):>9} "
f"sale={f(self.last_sale):>9} "
f"vwap30s={f(self.short.value):>9} "
f"vwap2m={f(self.long.value):>9} "
f"sig={self.signal or '--':<11} "
f"odd={100*self.odd/max(1,self.n):4.1f}% "
f"zero={100*self.zero/max(1,self.n):4.1f}%"
)
async def run(url, every=3000):
st = State()
async with websockets.connect(
url,
max_size=None
) as ws:
print("[consumer] connected", flush=True)
async for raw in ws:
m = json.loads(raw)
t = m["type"]
if t == "trade":
st.apply(m)
if not st.warming and st.n % every == 0:
print(
f"[consumer] {st.line()}",
flush=True
)
elif t == "replay_reset":
print(
f"[consumer] RESET -> "
f"{m['target_timestamp_ms']} "
f"({m['warmup_events']} warmup, "
f"{m['purged_stale_events']} purged)",
flush=True
)
st = State()
st.warming = True
elif t == "warmup_complete":
st.warming = False
print(
f"[consumer] WARM DONE {st.line()}",
flush=True
)
elif t in (
"session_completed",
"session_stopped"
):
print(
f"[consumer] {t.upper()} "
f"{st.line()}",
flush=True
)
break
else:
print(
f"[consumer] {t}",
flush=True
)
if __name__ == "__main__":
p = argparse.ArgumentParser()
p.add_argument(
"--url",
required=True
)
p.add_argument(
"--every",
type=int,
default=3000
)
a = p.parse_args()
asyncio.run(
run(a.url, a.every)
)
رولنگ VWAP ونڈوز لین دین کی تعداد کے بجائے مارکیٹ ٹائم اسٹیمپ پر مبنی ہیں۔ تمام آنے والے لین دین دونوں ونڈوز میں جاتے ہیں، اور پلے بیک وقت کے ساتھ 30 سیکنڈ یا 2 منٹ سے زیادہ پرانے مشاہدات کو ہٹا دیا جاتا ہے۔
لہذا، صارف ریاست آہستہ آہستہ ترقی کرتی ہے.

اہم بات یہ ہے کہ ان میں سے کوئی بھی ریاست اصل حالت سے نہیں آتی۔ TradeTape. صارف صرف ان واقعات کے بارے میں جانتا ہے جو WebSocket سے گزرے ہیں۔
یہ کھیل کے وقت کے دوران صاف طور پر کام کرتا ہے۔ نیویگیشن چیزوں کو مزید مشکل بنا دیتا ہے کیونکہ صارف کو ری سیٹ کیے بغیر پلے کرسر کو منتقل کرنے سے پلے کرسر ٹریڈنگ کے دن غلط مقام پر پھنس جائے گا۔
اپنی براؤزنگ کو محفوظ رکھیں
نیویگیشن صرف پلے کرسر کو حرکت دینے کا معاملہ نہیں ہے۔ اگر صارف پہلے ہی تجارتی دن میں ایک مقام پر رولنگ سٹیٹ قائم کر چکا ہے، تو اس حالت کو ری سیٹ کیے بغیر کسی دوسرے مقام پر جانے کے نتیجے میں دو مختلف مارکیٹ ریکارڈز کا مرکب ہو گا۔
فرض کریں کہ صارف 14:00 تک پہنچ جاتا ہے۔ 2 منٹ طویل VWAP میں تقریباً 13:58 کے بعد بھی لین دین شامل ہے۔ بس پلے بیک کرسر کو 13:30 پر واپس لے جائیں اور ٹرانزیکشنز کو ایکسپورٹ کرنا جاری رکھیں، اور مستقبل کے مشاہدات میموری میں رہیں گے۔

لہذا، پلے بیک کو نیچے کی دھارے کی حالت کو دوبارہ ترتیب دینا چاہیے اور عام پلے بیک کے جاری رہنے سے پہلے نئے ٹائم اسٹیمپ کی بنیاد پر دوبارہ تعمیر کرنا چاہیے۔
کنزیومر ری سیٹ اور وارم اپ
کہ seek() ہم نے جو طریقے شامل کیے ہیں۔ ReplaySession ہم پہلے ہی اس ترتیب پر کارروائی کر رہے ہیں۔ اہم حصہ موجودہ پروڈیوسر کو روکنے اور ٹیپ پر مطلوبہ مقام تلاش کرنے کے ساتھ شروع ہوتا ہے۔
was = self.state
await self._kill()
async with self._lock:
idx = max(
0,
min(
self.tape.index_at(target_ms),
len(self.tape)-1
)
)
self._epoch += 1
self.cursor = idx
self.state = State.PAUSED
self.clock.pause()
اگلا، اپنے مقصد سے 2 منٹ پہلے اپنے تیار پوائنٹ کا حساب لگائیں۔
warm = max(0, self.tape.index_at(int(self.tape.ts[idx]) - self.warmup_ms))
ہم 2 منٹ استعمال کرنے کی وجہ یہ ہے کہ یہ صارفین کے برقرار رکھنے والے طویل ترین رولنگ پیریڈ سے میل کھاتا ہے۔ اس وقفے کو دوبارہ چلانا نئی جگہ پر 30 سیکنڈ اور 2 منٹ کے VWAP دونوں کو دوبارہ تشکیل دینے کے لیے کافی ہے۔
تیار ٹرانزیکشن بھیجنے سے پہلے، سیشن کی قطار میں انتظار کرنے والے کسی بھی باقاعدہ لین دین کے پیغامات کو ہٹا دیا جاتا ہے۔
purged = sum(1 for m in self._q if m.get("type") == "trade")
self._q = collections.deque(m for m in self._q if m.get("type") != "trade")
سیشن پھر ایک واضح پیغام بھیجتا ہے۔ replay_reset واقعہ:
self._ctrl(
"replay_reset",
reason="seek",
target_timestamp_ms=int(self.tape.ts[idx]),
warmup_from_ms=int(self.tape.ts[warm]),
warmup_events=idx-warm,
purged_stale_events=purged,
epoch=self._epoch
)
صارف موجودہ حالت کو حذف کرکے جواب دیتا ہے۔
elif t == "replay_reset":
st = State()
st.warming = True
اب سیشن ہدف سے عین پہلے ماضی کے لین دین کو بھیج سکتا ہے۔
for k in range(warm, idx):
await self._put({
**self.tape[k].to_wire(),
"warmup": True
})
self._ctrl(
"warmup_complete",
market_ts_ms=int(self.tape.ts[idx])
)
یہ لین دین بالکل اسی طرح کام کرتا ہے۔ State.apply() منطق کو ایک عام ری پلے ایونٹ سمجھا جاتا ہے، لیکن صارف عام آؤٹ پٹ کو دبا دیتا ہے۔ warming ہے True.
لہذا مجموعی طور پر نیویگیشن بہاؤ مندرجہ ذیل ہے:

دوبارہ تعمیر کی حیثیت کو چیک کریں۔
مکمل سیشن چلانے کے لیے، میں نے پلے بیک موقوف کیا اور 13:30 کا تعاقب کیا۔ درخواست کردہ ٹائم اسٹیمپ پر یا اس کے بعد پہلا اصل واقعہ یہ ہے: 1784136600030.
صارفین کو موصول ہوا:

موجودہ صارف قوم غائب ہو چکی ہے، 3,456 تاریخی ٹریڈنگ میں سیشن میں نئے پوائنٹس کے ارد گرد دو رولنگ VWAP ونڈوز کو دوبارہ ترتیب دیا گیا ہے۔ نارمل ٹائم پلے بیک اب نیویگیشن باؤنڈری سے مارکیٹ کے حالات کو گزرے بغیر دوبارہ شروع کیا جا سکتا ہے۔
مکمل AAPL تجارتی دن دوبارہ دیکھیں
اب تمام ٹکڑے جڑے ہوئے ہیں۔ فاسٹ اے پی آئی کے ذریعے پورے دن کے ٹیپس کو کنٹرول کیا جا سکتا ہے، الگ الگ صارفین صرف ویب ساکٹ کے ذریعے لین دین اور کنٹرول ایونٹس کو دیکھتے ہیں۔
پہلے ٹرمینل میں پلے بیک سروس شروع کریں۔
python -m api.run --port 8765
ہماری آخری دوڑ کے لیے، آئیے اس کے ساتھ شروع کریں: 10xمارکیٹ کو روکیں اور اس پر سوئچ کریں: 50xدوبارہ شروع کریں، دوبارہ توقف کریں، 13:30 تک تلاش کریں، صارف کی حالت کو دوبارہ بنائیں، اور آخر کار ختم ہونے کی طرف دوڑیں۔ 400x.
مکمل کھیل چلائیں۔
عارضی طور پر درج ذیل کو محفوظ کریں۔ demo.py منصوبے کی جڑ میں. یہ اسکرپٹ صرف ایک ڈیمو ڈرائیور ہے۔ دوبارہ پیدا کرنے والے انجن اور صارفین اس پیکیج میں رہتے ہیں جو ہم نے پہلے ہی بنایا ہے۔
import asyncio, os, subprocess, sys
import httpx
BASE = "http://127.0.0.1:8765"
ROOT = os.getcwd()
SEEK_1330_MS = 1784136600000
async def demo():
async with httpx.AsyncClient(base_url=BASE, timeout=120) as c:
r = await c.post("/sessions", json={
"symbol": "AAPL",
"date": "2026-07-15",
"tag": "fullday",
"speed": 10.0
})
sid = r.json()["session_id"]
consumer = subprocess.Popen([
sys.executable,
"-u",
"-m",
"consumer.consumer",
"--url",
f"ws://127.0.0.1:8765/sessions/{sid}/stream",
"--every",
"25000"
], cwd=ROOT)
await asyncio.sleep(1.5)
controls = [
("START @10.0x", f"/sessions/{sid}/start", None, 4),
("PAUSE", f"/sessions/{sid}/pause", None, 1.5),
(
"SPEED 50x while paused",
f"/sessions/{sid}/speed",
{"speed": 50.0},
0.3
),
("RESUME", f"/sessions/{sid}/resume", None, 3),
("PAUSE", f"/sessions/{sid}/pause", None, 1),
(
"SEEK 13:30 while paused",
f"/sessions/{sid}/seek",
{"target_timestamp_ms": SEEK_1330_MS},
3
),
(
"RESUME after seek",
f"/sessions/{sid}/resume",
None,
3
),
(
"SPEED 400.0x to the close",
f"/sessions/{sid}/speed",
{"speed": 400.0},
2
)
]
for label, path, payload, wait in controls:
print(f"n--- {label} ---")
if payload is None:
await c.post(path)
else:
await c.post(path, json=payload)
await asyncio.sleep(wait)
for _ in range(600):
await asyncio.sleep(1)
state = (
await c.get(f"/sessions/{sid}")
).json()
if state["state"] in ("completed", "stopped"):
break
print(
f"nfinal: {state['state']} "
f"{state['cursor']:,}/{state['total_events']:,}"
)
print(
"timing:",
(
await c.get(f"/sessions/{sid}/timing")
).json()
)
if consumer.poll() is None:
consumer.terminate()
asyncio.run(demo())
اسے دوسرے ٹرمینل میں چلائیں۔
python demo.py
صارف اپنے عمل کے طور پر شروع ہوتا ہے اور پلے بیک شروع ہونے سے پہلے WebSocket سے جڑ جاتا ہے۔
اصل دوڑ اس طرح شروع ہوئی:

یہ آپ کو کسی سیشن کو روکنے، اسے مختلف رفتار سے دوبارہ لاک کرنے اور پلے بیک کو دوبارہ شروع کیے بغیر دوبارہ شروع کرنے کی اجازت دیتا ہے۔
درج ذیل کمانڈ براہ راست 13:30 پر جاتی ہے۔

یہ پچھلے حصے سے ریاست کے لیے محفوظ نیویگیشن ہے جو پورے سسٹم میں ہوتا ہے۔ صارفین پچھلی حالت کو چھوڑ دیتے ہیں۔ 3,456 یہ وارم اپ ایونٹ ہونے کے بعد ہی مارکیٹ کے نئے ٹائم اسٹیمپ سے جاری رہتا ہے۔
اس کے بعد آپ اپنے باقی سیشن کو تیز کر سکتے ہیں۔

سیشن ٹیپ کے اختتام تک پہنچنے تک صارف آنے والے واقعات سے اسٹیٹ کو اپ ڈیٹ کرتا رہتا ہے۔

دونوں شمار مختلف چیزوں کو بیان کرتے ہیں۔ سیشن کا کرسر اس پر ختم ہوتا ہے: 1,032,409/1,032,409اس کا مطلب ہے کہ ہم ایک ٹیپ کے اختتام پر پہنچ گئے ہیں جو سارا دن رہتا ہے۔ صارفین کی رپورٹ 306,343 اس کی وجہ یہ ہے کہ نیویگیشن کے دوران، اس حالت کو ایک نئے مقام سے صاف اور دوبارہ بنایا جاتا ہے۔ ایکسپلوریشن ریئل ٹائم میں چھوڑے گئے تمام لین دین کو اسٹریم کرنے کے بجائے تاریخی ٹیپ کے کچھ حصوں پر چھلانگ لگاتی ہے۔
realized_speed اسے اس رن میں جان بوجھ کر بغیر سیٹ کیا گیا تھا، کیونکہ پلے بیک کو کئی بار توقف، رفتار میں تبدیلی اور تلاش کرنے کی وجہ سے دوبارہ منجمد کیا گیا تھا۔ ایک واحد اختتام سے آخر تک رفتار کا تناسب معنی خیز طور پر ان سیشنز کا حساب نہیں رکھتا ہے جو جان بوجھ کر اپنی گھڑیوں کو درمیانی سلسلہ میں تبدیل کرتے ہیں۔
یہاں جو چیز اہم ہے وہ یہ ہے کہ وہی ریکارڈ شدہ ٹیپ پورے کنٹرول کے سلسلے میں زندہ رہتی ہے، صارف بازیافت کے بعد اپنی حالت کو دوبارہ تشکیل دیتا ہے، اور پلے بیک سیشن کے ختم ہونے تک جاری رہتا ہے۔
ری پلے انجن ٹیسٹ
اسے سارا دن چلانے سے یہ ظاہر ہوگا کہ سسٹم پورے کنٹرول کے سلسلے سے گزرنے کے قابل ہے، لیکن اکیلے ٹرمینل آؤٹ پٹ آپ کو یہ نہیں بتائے گا کہ آیا پلے بیک ترتیب میں رہا، توقف کی حدود کا مشاہدہ کیا، یا نیویگیشن کے بعد درست حالت کو دوبارہ قائم کیا۔
ہم اس رویے کی جانچ چھوٹے پیمانے پر کریں گے۔ quiet15m پہلے بنائے گئے ٹیپ:
market-time-machine/
└── tests/
├── __init__.py
├── conftest.py
└── test_replay.py
ٹیسٹ سوٹ چار شعبوں پر محیط ہے: ایونٹ کی ترتیب، پلے بیک ٹائمنگ، موقوف/دوبارہ شروع کرنے کا برتاؤ، اور بازیافت کے بعد ریاست کی تعمیر نو۔
بنانا tests/test_replay.py
بنانا tests/test_replay.py:
import asyncio
import numpy as np
import pytest
from replay import config
from replay.events import TradeTape
from replay.session import ReplaySession
from replay.clock import ReplayClock, replay_batches, new_stats, summarize
TAPE = sorted(config.PROCESSED.glob("*_quiet15m.npz"))[0]
@pytest.fixture
def tape():
return TradeTape.load(TAPE)
async def collect(sess, seconds):
out = []
async def drain():
async for m in sess.events():
out.append(m)
t = asyncio.create_task(drain())
await asyncio.sleep(seconds)
return out, t
@pytest.mark.asyncio
async def test_ordering(tape):
s = ReplaySession(tape, speed=500)
out, t = await collect(s, 0.1)
await s.start()
await asyncio.sleep(2)
await s.stop()
t.cancel()
trades = [
m for m in out
if m["type"] == "trade"
]
assert len(trades) > 1000
keys = [
(m["timestamp_ms"], m["sequence"])
for m in trades
]
assert keys == sorted(keys)
assert len(set(keys)) == len(keys)
@pytest.mark.asyncio
@pytest.mark.parametrize("speed", [10, 50, 100])
async def test_timing(tape, speed):
end = int(
np.searchsorted(
tape.ts,
tape.ts[0] + 60_000,
side="right"
)
)
clock = ReplayClock(tape.ts[0], speed).start()
st = new_stats(speed)
async for i, j in replay_batches(tape, clock, 0, st):
if j >= end:
break
r = summarize(st)
assert abs(r["speed_error_pct"]) < 5
assert r["lateness_p95_ms"] < 50
@pytest.mark.asyncio
async def test_pause_resume(tape):
s = ReplaySession(tape, speed=100)
out, t = await collect(s, 0.05)
await s.start()
await asyncio.sleep(1)
await s.pause()
n = len([
m for m in out
if m["type"] == "trade"
])
await asyncio.sleep(1)
assert len([
m for m in out
if m["type"] == "trade"
]) == n
await s.resume()
await asyncio.sleep(1)
await s.stop()
t.cancel()
seqs = [
m["sequence"]
for m in out
if m["type"] == "trade"
]
assert seqs == sorted(seqs)
assert len(set(seqs)) == len(seqs)
@pytest.mark.asyncio
async def test_pause_seek_resume(tape):
s = ReplaySession(
tape,
speed=200,
warmup_ms=120_000
)
out, t = await collect(s, 0.05)
await s.start()
await asyncio.sleep(0.5)
await s.pause()
target = int(tape.ts[0]) + 300_000
await s.seek(target)
assert s.info()["state"] == "paused"
def past():
return [
m for m in out
if m["type"] == "trade"
and not m.get("warmup")
and m["timestamp_ms"] >= target
]
await asyncio.sleep(0.4)
assert not past()
await s.resume()
await asyncio.sleep(1)
got = past()
await s.stop()
t.cancel()
assert got
seqs = [m["sequence"] for m in got]
assert seqs == sorted(seqs)
assert len(set(seqs)) == len(seqs)
@pytest.mark.asyncio
async def test_seek_state_equivalence(tape):
import sys
sys.path.insert(0, str(config.ROOT))
from consumer.consumer import State as ConsumerState
s = ReplaySession(
tape,
speed=200,
warmup_ms=120_000
)
live = ConsumerState()
reset = None
snap = None
out = []
async def drain():
nonlocal live, reset, snap
async for m in s.events():
out.append(m)
if m["type"] == "trade":
live.apply(m)
elif m["type"] == "replay_reset":
reset = m
live = ConsumerState()
elif m["type"] == "warmup_complete":
snap = (
live.n,
live.vol,
live.short.value,
live.long.value
)
t = asyncio.create_task(drain())
await s.start()
await asyncio.sleep(1)
await s.seek(
int(tape.ts[0]) + 600_000
)
for _ in range(100):
if snap:
break
await asyncio.sleep(0.05)
await s.stop()
t.cancel()
assert snap
fresh = ConsumerState()
lo = tape.index_at(
reset["warmup_from_ms"]
)
hi = tape.index_at(
reset["target_timestamp_ms"]
)
for k in range(lo, hi):
fresh.apply(tape[k].to_wire())
n, vol, short, long = snap
assert n == fresh.n == reset["warmup_events"]
assert vol == fresh.vol
assert short == pytest.approx(
fresh.short.value,
rel=1e-12
)
assert long == pytest.approx(
fresh.long.value,
rel=1e-12
)
kinds = [m["type"] for m in out]
seg = out[
kinds.index("replay_reset") + 1:
kinds.index("warmup_complete")
]
assert not [
m for m in seg
if m["type"] == "trade"
and not m.get("warmup")
]
test_ordering() اس بات کی توثیق کریں کہ برآمد شدہ لین دین اس طرح ترتیب دیے گئے ہیں: (timestamp, sequence) اور ایک ہی واقعہ دو بار نہیں ہوتا۔
ٹائمنگ ٹیسٹ تاریخی مارکیٹ ٹائم کے 60 سیکنڈ تک چلتا ہے۔ 10x, 50xاور 100x. ایونٹ لوپ سے سخت ریئل ٹائم شیڈولر کی طرح برتاؤ کرنے کی توقع کرنے کے بجائے، ہم تھوڑی سی رواداری کی اجازت دیتے ہیں۔ یعنی، وصولی کی رفتار ہدف کے 5% کے اندر رہنی چاہیے، اور 95ویں پرسنٹائل تاخیر 50 ms سے کم رہنی چاہیے۔
test_pause_resume() کچھ مختلف کے لیے چیک کریں۔ ایک بار pause() ایک بار واپس آنے کے بعد، پلے بیک دوبارہ شروع ہونے تک موصول ہونے والی لین دین کی تعداد میں کوئی تبدیلی نہیں ہونی چاہیے۔ دوبارہ شروع کرنے کے بعد، نتیجے میں ترتیب کو سیدھ میں رکھنا جاری رکھنا چاہیے اور اس میں کوئی ڈپلیکیٹ نہیں ہونا چاہیے۔
test_pause_seek_resume() مکمل پلے بیک کے لیے استعمال ہونے والے عین کنٹرول پیٹرن کا احاطہ کرتا ہے۔ سیشن موقوف ہو جائے گا، 5 منٹ کے لیے ٹیپ پر جائے گا، نئی جگہ پر موقوف رہے گا، اور تب ہی عام نیویگیشن کے بعد لین دین جاری کرنا شروع کر دے گا۔ resume().
آزادانہ طور پر ریاست کی تعمیر نو کی تصدیق کریں۔
سب سے طاقتور امتحان ہے۔ test_seek_state_equivalence().
جب پلے بیک کا پتہ چلتا ہے، صارف کو ایک ری سیٹ ملتا ہے اور پھر دو منٹ کا وارم اپ ایونٹ ملتا ہے۔ محض جانچ پڑتال کے بجائے warmup_complete اشارہ کرنے پر، یہ ٹیسٹ بالکل نیا ٹیسٹ ہے۔ ConsumerState آزادانہ طور پر ٹیپ سے براہ راست ایک جیسے ریکارڈنگ وقفے فراہم کرتا ہے۔
for k in range(lo, hi):
fresh.apply(tape[k].to_wire())
ری پلے تعمیرات اور آزادانہ طور پر دوبارہ تعمیر شدہ ریاستوں کو درج ذیل سے متفق ہونا چاہیے:
event count
cumulative volume
30-second VWAP
2-minute VWAP
VWAP قدر کا موازنہ نسبتا رواداری سے کیا جاتا ہے۔ 1e-12. ٹیسٹ اس بات کو بھی یقینی بناتا ہے کہ عام پلے بیک لین دین درمیان میں ندی میں نہ پھسل جائے۔ replay_reset اور warmup_complete.
pytest ترتیب
غیر مطابقت پذیر ٹیسٹوں میں pytest-asyncio. بنانا tests/conftest.py:
import pytest
def pytest_configure(config):
config.addinivalue_line(
"markers",
"asyncio"
)
پھر شامل کریں۔ pytest.ini منصوبے کی جڑ میں:
[pytest]
asyncio_mode = auto
مکمل سویٹ چلاتا ہے۔
pytest tests/ -v
ریکارڈ شدہ عملدرآمد کے نتائج درج ذیل ہیں:

ٹیسٹ اس سے زیادہ کا احاطہ کرتا ہے کہ آیا پلے بیک آخر میں ٹیپ کے آخر تک پہنچ جاتا ہے۔ اس بات کو یقینی بنائیں کہ پلے بیک کے بعد ریکارڈنگ کی ترتیب برقرار رہے، تیز رفتاری کا وقت متوقع رواداری کے اندر رہے، جو کنٹرول کرتا ہے واقعات کی ترتیب کو برقرار رکھتا ہے، اور یہ کہ نیویگیشن کے بعد کی تعمیر نو کی حالت بنیادی ریکارڈ شدہ ڈیٹا سے ایک آزاد تعمیر نو سے میل کھاتی ہے۔
نتیجہ
مجھے اس تعمیر کے بارے میں جو چیز سب سے زیادہ پسند آئی وہ یہ ہے کہ جب دوبارہ گھڑی دی جائے تو وہی تاریخی ڈیٹا سیٹ کتنا مختلف محسوس ہوتا ہے۔
ہم نے EODHD میں ایک مکمل AAPL سیشن کے ساتھ آغاز کیا اور اس کے ساتھ اختتام پذیر ہوا جس نے ہمیں آہستہ آہستہ آگے بڑھنے، آگے بڑھنے، درمیان میں توقف کرنے، دن میں مختلف پوائنٹس پر چھلانگ لگانے، اور جاری رکھنے کی اجازت دی جب کہ صارف صرف اس پر ردعمل ظاہر کرتا ہے جو انہیں اب تک حاصل ہوا ہے۔
اس منصوبے کو مزید ترقی دینے کے لیے ابھی بھی کافی گنجائش ہے۔ ری پلے متعدد علامتوں، بھرپور مارکیٹ اسٹیٹس، متعدد نیچے دھارے صارفین، جاری ری پلے سیشنز، اور یہاں تک کہ حکمت عملی اور عمل درآمد کے اجزاء کو براہ راست دھارے سے منسلک کرنے کی حمایت کر سکتے ہیں۔ ہم جان بوجھ کر ان حصوں کو موجودہ ورژن میں رکھتے ہیں، لیکن اب ہم بنیادی پلے بیک پرت پر تعمیر کر سکتے ہیں۔
میرے نزدیک یہ اس منصوبے کا ایک مفید نتیجہ ہے۔ EODHD تاریخی واقعات فراہم کرتا ہے، لیکن ری پلے لیئر دوسرے سافٹ ویئر کو ان واقعات کو ڈیٹا سیٹ کے بجائے تجارتی دن کے طور پر تجربہ کرنے کی اجازت دیتا ہے جہاں ہم پہلے ہی جانتے ہیں کہ دن کیسے ختم ہوتا ہے۔