SimpleChatbot / hermes_overlay /tests /test_phase4_locking_durability.py
Amin
Deploy: HermesFace finalized project to HF Space
2e658e7
Raw
History Blame Contribute Delete
8.76 kB
from __future__ import annotations
import asyncio
import pytest
import trading.domain.execution_service as execution_module
from trading.adapters.paper_adapter import PaperExchangeAdapter
from trading.domain.approvals import create_approval
from trading.domain.execution_contracts import ExecuteApprovedPlanCommand
from trading.domain.execution_contracts import ExecutionErrorCode
from trading.domain.execution_service import ExecutionService, ExecutionServiceError
from trading.domain.identity_repo import ensure_paper_exchange_account
from trading.domain.plan_revalidation import PlanRevalidationResult
from trading.domain.plans_repo import create_plan
from trading.domain.risk_limits import MarketRiskContext
from trading.domain.schema import apply_migrations, connect
@pytest.fixture(autouse=True)
def _db(tmp_path, monkeypatch):
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
monkeypatch.setenv("TRADING_MODE", "paper")
monkeypatch.setenv("HERMES_EXECUTION_ENABLED", "true")
ok, messages = apply_migrations()
assert ok, messages
def _plan():
return {
"symbol": "BTCUSDT",
"decision": "LONG",
"entry": 100.0,
"stop_loss": 95.0,
"take_profit": 110.0,
"leverage": 3,
"quantity": 1.0,
"risk_profile": "moderate",
"risk_percent": 0.01,
"futuresVerified": True,
"noTradeGuard": False,
"trading_readiness": "ready",
"executable": True,
"expires_at": "2099-01-01T00:00:00Z",
}
def _accepted(plan_id: str) -> PlanRevalidationResult:
return PlanRevalidationResult(
revalidation_id=f"rv-{plan_id}",
accepted=True,
reasons=(),
context_hash="context-hash",
observed_price=100.0,
observed_spread_bps=1.0,
observed_slippage_percent=0.0,
market_risk=MarketRiskContext(
provider_fresh=True,
reconciliation_clear=True,
spread_bps=1.0,
funding_rate=0.0,
liquidity_ok=True,
),
orderbook={"bids": [[99.9, 100]], "asks": [[100.0, 100]]},
)
def _setup(owner_id: str = "owner-lock"):
account = ensure_paper_exchange_account(owner_id)
plan = create_plan(_plan(), owner_id=owner_id, account_id=account.account_id)
approval_id = create_approval(
plan["plan_id"], plan["plan_hash"], owner_id, account_id=account.account_id,
)
command = ExecuteApprovedPlanCommand(
plan_id=plan["plan_id"], approval_id=approval_id, owner_id=owner_id,
account_id=account.account_id, idempotency_key=f"lock:{plan['plan_id']}:{approval_id}",
)
return account, plan, command
class ObservingAdapter(PaperExchangeAdapter):
def __init__(self, *, idem_key: str) -> None:
super().__init__(default_price=100.0)
self.idem_key = idem_key
self.configure_observation = None
self.submit_observation = None
self.configure_calls = 0
self.create_order_calls = 0
async def configure_account(self, symbol, leverage, margin_mode, position_mode):
self.configure_calls += 1
conn = connect()
try:
idem = conn.execute(
"SELECT status FROM idempotency_keys WHERE idem_key=?", (self.idem_key,)
).fetchone()
execution = conn.execute(
"SELECT execution_id,state FROM execution_history WHERE idempotency_key=?",
(self.idem_key,),
).fetchone()
events = conn.execute(
"SELECT to_state FROM execution_events WHERE execution_id=? ORDER BY id",
(execution[0],),
).fetchall()
self.configure_observation = {
"idem_status": idem[0] if idem else None,
"execution_state": execution[1] if execution else None,
"events": [row[0] for row in events],
}
finally:
conn.close()
return await super().configure_account(symbol, leverage, margin_mode, position_mode)
async def create_order(self, req):
self.create_order_calls += 1
conn = connect()
try:
order = conn.execute(
"SELECT order_id,state FROM orders WHERE client_order_id=?", (req.client_order_id,)
).fetchone()
execution = conn.execute(
"SELECT state FROM execution_history WHERE idempotency_key=?", (self.idem_key,)
).fetchone()
self.submit_observation = {
"order_exists": order is not None,
"order_state": order[1] if order else None,
"execution_state": execution[0] if execution else None,
}
finally:
conn.close()
return await super().create_order(req)
def test_durable_idempotency_execution_transitions_and_order_precede_adapter_calls(monkeypatch):
_account, plan, command = _setup()
async def accepted(plan_id, **_kwargs):
return _accepted(plan_id)
monkeypatch.setattr(execution_module, "revalidate_plan", accepted)
adapter = ObservingAdapter(idem_key=command.idempotency_key or "")
result = asyncio.run(ExecutionService(adapter).execute_approved_plan(command))
assert result["status"] == "protected"
assert adapter.configure_observation == {
"idem_status": "in_progress",
"execution_state": "CONFIGURING_ACCOUNT",
"events": ["CREATED", "VALIDATING", "LOCKED", "CONFIGURING_ACCOUNT"],
}
assert adapter.submit_observation == {
"order_exists": True,
"order_state": "submitted",
"execution_state": "SUBMITTING_ENTRY",
}
conn = connect()
try:
lock = conn.execute(
"SELECT holder,expires_at,fencing_token FROM execution_locks WHERE lock_name=?",
(f"exec:{command.account_id}:BTC/USDT:USDT:one_way",),
).fetchone()
idem = conn.execute(
"SELECT status,result_ref FROM idempotency_keys WHERE idem_key=?",
(command.idempotency_key,),
).fetchone()
finally:
conn.close()
assert lock is not None and lock[0] == "" and float(lock[1]) == 0.0 and int(lock[2]) >= 1
assert idem is not None and idem[0] == "completed" and idem[1]
def test_concurrent_same_command_has_one_mutation_winner(monkeypatch):
_account, _plan_row, command = _setup("owner-concurrent")
async def accepted(plan_id, **_kwargs):
await asyncio.sleep(0)
return _accepted(plan_id)
monkeypatch.setattr(execution_module, "revalidate_plan", accepted)
adapter = ObservingAdapter(idem_key=command.idempotency_key or "")
service = ExecutionService(adapter)
async def run():
first, second = await asyncio.gather(
service.execute_approved_plan(command),
service.execute_approved_plan(command),
)
return first, second
first, second = asyncio.run(run())
assert sorted([first["status"], second["status"]]) == ["duplicate", "protected"]
assert adapter.configure_calls == 1
assert adapter.create_order_calls == 1
conn = connect()
try:
order_count = conn.execute("SELECT COUNT(*) FROM orders").fetchone()[0]
position_count = conn.execute(
"SELECT COUNT(*) FROM positions WHERE account_id=? AND status IN ('open','partially_open')",
(command.account_id,),
).fetchone()[0]
finally:
conn.close()
assert order_count == 1
assert position_count == 1
def test_execution_record_failure_marks_idempotency_failed_and_prevents_adapter_call(monkeypatch):
_account, _plan_row, command = _setup("owner-db-failure")
async def accepted(plan_id, **_kwargs):
return _accepted(plan_id)
monkeypatch.setattr(execution_module, "revalidate_plan", accepted)
def explode_begin(**_kwargs):
raise RuntimeError("database unavailable before execution record")
monkeypatch.setattr(execution_module, "begin_execution", explode_begin)
adapter = ObservingAdapter(idem_key=command.idempotency_key or "")
with pytest.raises(ExecutionServiceError, match="approved plan execution failed") as caught:
asyncio.run(ExecutionService(adapter).execute_approved_plan(command))
assert caught.value.code is ExecutionErrorCode.EXECUTION_FAILED
assert adapter.configure_calls == 0
assert adapter.create_order_calls == 0
conn = connect()
try:
row = conn.execute(
"SELECT status,error_class FROM idempotency_keys WHERE idem_key=?",
(command.idempotency_key,),
).fetchone()
finally:
conn.close()
assert row is not None and row[0] == "failed" and row[1] == "RuntimeError"