57 lines
1.6 KiB
Python
57 lines
1.6 KiB
Python
import asyncio
|
|
import json
|
|
from collections.abc import AsyncIterator
|
|
from pathlib import Path
|
|
|
|
from app.bars.models import Bar, Timeframe
|
|
|
|
|
|
class ReplaySource:
|
|
name = "replay"
|
|
delay_minutes = 0
|
|
|
|
def __init__(self, path: str | Path, speed: float = 0):
|
|
self.path = Path(path)
|
|
self.speed = speed
|
|
|
|
def supports_history(self) -> bool:
|
|
return True
|
|
|
|
def supports_stream(self) -> bool:
|
|
return True
|
|
|
|
def _rows(self) -> list[dict]:
|
|
if not self.path.exists():
|
|
return []
|
|
with self.path.open(encoding="utf-8") as tape:
|
|
return [json.loads(line) for line in tape if line.strip()]
|
|
|
|
async def history(
|
|
self,
|
|
symbol: str,
|
|
tf: Timeframe,
|
|
start: int | None = None,
|
|
end: int | None = None,
|
|
*,
|
|
range_: str | None = None,
|
|
) -> list[Bar]:
|
|
return [
|
|
bar
|
|
for row in self._rows()
|
|
if (bar := Bar.from_dict(row["bar"])).tf is tf
|
|
and (not symbol or bar.symbol == symbol)
|
|
and (start is None or bar.t >= start)
|
|
and (end is None or bar.t < end)
|
|
]
|
|
|
|
async def stream(self, symbol: str) -> AsyncIterator[Bar]:
|
|
previous_arrival: float | None = None
|
|
for row in self._rows():
|
|
bar = Bar.from_dict(row["bar"])
|
|
if symbol and bar.symbol != symbol:
|
|
continue
|
|
arrival = float(row["arrival_t"])
|
|
if self.speed > 0 and previous_arrival is not None:
|
|
await asyncio.sleep(max(0, arrival - previous_arrival) / self.speed)
|
|
previous_arrival = arrival
|
|
yield bar
|