build 059.7: add websocket OHLC adapter

This commit is contained in:
2026-07-17 11:44:03 +03:00
parent 85e9662f5b
commit 20a70df1f0
3 changed files with 904 additions and 0 deletions

View File

@@ -5,16 +5,21 @@ from __future__ import annotations
from datetime import datetime, timezone
from src.market_data.acquisition.adapters.dzengi.mapper import (
map_dzengi_websocket_ohlc_to_candle_close_event,
map_dzengi_websocket_quote_to_quote,
)
from src.market_data.acquisition.adapters.dzengi.parser import (
parse_dzengi_websocket_ohlc,
parse_dzengi_websocket_quote,
)
from src.market_data.acquisition.models.candle_close import CandleCloseEvent
from src.market_data.acquisition.models.quote import Quote
from src.market_data.acquisition.validation.schema import (
validate_dzengi_websocket_ohlc_schema,
validate_dzengi_websocket_quote_schema,
)
from src.market_data.acquisition.validation.values import (
validate_dzengi_websocket_ohlc_values,
validate_dzengi_websocket_quote_values,
)
@@ -35,3 +40,22 @@ class DzengiWebSocketQuoteAdapter:
response,
received_at=received_at or datetime.now(timezone.utc),
)
# Преобразует одно декодированное сообщение Dzengi WebSocket
# в CandleCloseEvent.
class DzengiWebSocketOhlcAdapter:
def map_message(
self,
document: object,
*,
received_at: datetime | None = None,
) -> CandleCloseEvent:
validated = validate_dzengi_websocket_ohlc_schema(document)
event = parse_dzengi_websocket_ohlc(validated)
validate_dzengi_websocket_ohlc_values(event)
return map_dzengi_websocket_ohlc_to_candle_close_event(
event,
received_at=received_at or datetime.now(timezone.utc),
)

View File

@@ -0,0 +1,112 @@
# app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py
from __future__ import annotations
from datetime import datetime, timezone
import pytest
from src.market_data.acquisition.adapters.dzengi.websocket import (
DzengiWebSocketOhlcAdapter,
)
from src.market_data.acquisition.exceptions import (
CandleWebSocketValueError,
)
def test_adapter_maps_document_to_candle_close_event() -> None:
received_at = datetime(2026, 7, 13, tzinfo=timezone.utc)
result = DzengiWebSocketOhlcAdapter().map_message(
{
"status": "OK",
"destination": "ohlc.event",
"payload": {
"symbol": "BTC/USD",
"interval": "1m",
"type": "classic",
"t": 1752364800000,
"o": "100",
"h": "105",
"l": "99",
"c": "103",
},
},
received_at=received_at,
)
assert result.symbol == "BTC/USD"
assert result.interval == "1m"
assert result.candle_type == "classic"
assert str(result.open_price) == "100"
assert str(result.high_price) == "105"
assert str(result.low_price) == "99"
assert str(result.close_price) == "103"
assert result.received_at is received_at
assert result.source == "dzengi_websocket_ohlc"
def test_adapter_sets_received_at_when_missing() -> None:
result = DzengiWebSocketOhlcAdapter().map_message(
{
"status": "OK",
"destination": "ohlc.event",
"payload": {
"symbol": "BTC/USD",
"interval": "1m",
"type": "classic",
"t": 1752364800000,
"o": "100",
"h": "105",
"l": "99",
"c": "103",
},
}
)
assert result.received_at.tzinfo is timezone.utc
def test_adapter_preserves_value_validation_error() -> None:
with pytest.raises(CandleWebSocketValueError):
DzengiWebSocketOhlcAdapter().map_message(
{
"status": "OK",
"destination": "ohlc.event",
"payload": {
"symbol": "BTC/USD",
"interval": "1m",
"type": "classic",
"t": 1752364800000,
"o": "100",
"h": "90",
"l": "99",
"c": "95",
},
}
)
def test_adapter_runs_complete_pipeline() -> None:
result = DzengiWebSocketOhlcAdapter().map_message(
{
"status": "OK",
"destination": "ohlc.event",
"payload": {
"symbol": "BTC/USD",
"interval": "5m",
"type": "heikin-ashi",
"t": 1752364800000,
"o": "250",
"h": "260",
"l": "245",
"c": "255",
},
}
)
assert result.interval == "5m"
assert result.candle_type == "heikin-ashi"
assert str(result.close_price) == "255"