Pure Market Making with Inventory Bounds and Order Management Controls
Summary
This document presents a configurable pure market making controller for a single trading pair. It defines buy and sell spreads and order amounts, portfolio allocation, and a target base-asset share bounded by minimum and maximum levels. The configuration also includes leverage and position mode, order types, executor refresh timing, and a cap on active executors per level. Additional controls cover side-specific cooldowns, a delay before positions count as effective, and price-distance and refresh tolerances for deciding when orders are too close or should be replaced.
The controller configuration exposes optional take-profit and global profit or loss thresholds, with settings for when those thresholds activate and whether profit is measured against the position or the portfolio. Validators parse and check settings such as spreads, amount lists, order types, and position side. The excerpt’s description points to hanging executors and breakeven awareness, but the source is truncated before its full trading logic is visible. It provides no backtest, market conditions, or performance evidence. These parameters describe operational controls, not proof that the controller is profitable; inventory exposure, leverage, fills, fees, and adverse selection remain material considerations.
Key ideas
- The controller configures quote spreads and order amounts separately for buying and selling.
- Target base-asset inventory is bounded by minimum and maximum allocation settings.
- Cooldowns, price-distance tolerances, refresh thresholds, and executor limits shape order management.
- Optional global profit and loss controls can use position-based or portfolio-based references.
- The source excerpt is incomplete and includes no backtest or performance evidence.
Tags
Full text
# PMMister
# PMMister
Advanced PMM (Pure Market Making) controller with sophisticated position management.
Features hanging executors, price distance requirements, and breakeven awareness.
## Source (Apache-2.0)
```python
from collections import defaultdict
from decimal import Decimal
from typing import Dict, List, Optional, Tuple, Union
from pydantic import Field, field_validator
from pydantic_core.core_schema import ValidationInfo
from hummingbot.core.data_type.common import (
MarketDict,
OrderType,
PositionAction,
PositionMode,
PositionSide,
PriceType,
TradeType,
)
from hummingbot.strategy_v2.controllers.controller_base import ControllerBase, ControllerConfigBase
from hummingbot.strategy_v2.executors.data_types import ConnectorPair
from hummingbot.strategy_v2.executors.order_executor.data_types import ExecutionStrategy, OrderExecutorConfig
from hummingbot.strategy_v2.executors.position_executor.data_types import PositionExecutorConfig, TripleBarrierConfig
from hummingbot.strategy_v2.models.executor_actions import CreateExecutorAction, ExecutorAction, StopExecutorAction
from hummingbot.strategy_v2.utils.common import parse_comma_separated_list, parse_enum_value
class PMMisterConfig(ControllerConfigBase):
"""
Advanced PMM (Pure Market Making) controller with sophisticated position management.
Features hanging executors, price distance requirements, and breakeven awareness.
"""
controller_type: str = "generic"
controller_name: str = "pmm_mister"
connector_name: str = Field(default="binance")
trading_pair: str = Field(default="BTC-USDT")
portfolio_allocation: Decimal = Field(default=Decimal("0.1"), json_schema_extra={"is_updatable": True})
target_base_pct: Decimal = Field(default=Decimal("0.5"), json_schema_extra={"is_updatable": True})
min_base_pct: Decimal = Field(default=Decimal("0.3"), json_schema_extra={"is_updatable": True})
max_base_pct: Decimal = Field(default=Decimal("0.7"), json_schema_extra={"is_updatable": True})
buy_spreads: List[float] = Field(default="0.0005", json_schema_extra={"is_updatable": True})
sell_spreads: List[float] = Field(default="0.0005", json_schema_extra={"is_updatable": True})
buy_amounts_pct: Union[List[Decimal], None] = Field(default="1", json_schema_extra={"is_updatable": True})
sell_amounts_pct: Union[List[Decimal], None] = Field(default="1", json_schema_extra={"is_updatable": True})
executor_refresh_time: int = Field(default=30, json_schema_extra={"is_updatable": True})
# Enhanced timing parameters
buy_cooldown_time: int = Field(default=60, json_schema_extra={"is_updatable": True})
sell_cooldown_time: int = Field(default=60, json_schema_extra={"is_updatable": True})
buy_position_effectivization_time: int = Field(default=120, json_schema_extra={"is_updatable": True})
sell_position_effectivization_time: int = Field(default=120, json_schema_extra={"is_updatable": True})
# Price distance tolerance - prevents placing new orders when existing ones are too close to current price
price_distance_tolerance: Decimal = Field(default=Decimal("0.0005"), json_schema_extra={"is_updatable": True})
# Refresh tolerance - triggers replacing open orders when price deviates from theoretical level
refresh_tolerance: Decimal = Field(default=Decimal("0.0005"), json_schema_extra={"is_updatable": True})
tolerance_scaling: Decimal = Field(default=Decimal("1.2"), json_schema_extra={"is_updatable": True})
leverage: int = Field(default=20, json_schema_extra={"is_updatable": True})
position_mode: PositionMode = Field(default=PositionMode.ONEWAY)
# LONG: buys accumulate, sells reduce. SHORT: sells accumulate, buys reduce.
position_side: TradeType = Field(default="BUY")
take_profit: Optional[Decimal] = Field(default=Decimal("0.0001"), gt=0, json_schema_extra={"is_updatable": True})
take_profit_order_type: Optional[OrderType] = Field(default=OrderType.LIMIT_MAKER, json_schema_extra={"is_updatable": True})
open_order_type: Optional[OrderType] = Field(default=OrderType.LIMIT_MAKER, json_schema_extra={"is_updatable": True})
max_active_executors_by_level: Optional[int] = Field(default=4, json_schema_extra={"is_updatable": True})
tick_mode: bool = Field(default=False, json_schema_extra={"is_updatable": True})
position_profit_protection: bool = Field(default=False, json_schema_extra={"is_updatable": True})
min_skew: Decimal = Field(default=Decimal("1.0"), json_schema_extra={"is_updatable": True})
global_take_profit: Decimal = Field(default=Decimal("0.03"), json_schema_extra={"is_updatable": True})
global_stop_loss: Decimal = Field(default=Decimal("0.05"), json_schema_extra={"is_updatable": True})
# Global TP/SL activation settings
global_tp_enabled: bool = Field(default=False, json_schema_extra={"is_updatable": True})
global_sl_enabled: bool = Field(default=False, json_schema_extra={"is_updatable": True})
# TP activates when position >= this threshold: "min_base" (earlier) or "target_base" (later)
global_tp_activation_from: str = Field(default="min_base", json_schema_extra={"is_updatable": True})
# SL activates when position >= this threshold: "target_base" (earlier) or "max_base" (later)
global_sl_activation_from: str = Field(default="target_base", json_schema_extra={"is_updatable": True})
# PnL reference: "position" = pnl/position_value, "portfolio" = pnl/total_amount_quote
global_pnl_reference: str = Field(default="position", json_schema_extra={"is_updatable": True})
@field_validator("take_profit", mode="before")
@classmethod
def validate_target(cls, v):
if isinstance(v, str):
if v == "":
return None
return Decimal(v)
return v
@field_validator('take_profit_order_type', mode="before")
@classmethod
def validate_order_type(cls, v) -> OrderType:
if v is None:
return OrderType.MARKET
return parse_enum_value(OrderType, v, "take_profit_order_type")
@field_validator('open_order_type', mode="before")
@classmethod
def validate_open_order_type(cls, v) -> OrderType:
if v is None:
return OrderType.MARKET
return parse_enum_value(OrderType, v, "open_order_type")
@field_validator('buy_spreads', 'sell_spreads', mode="before")
@classmethod
def parse_spreads(cls, v):
return parse_comma_separated_list(v)
@field_validator('buy_amounts_pct', 'sell_amounts_pct', mode="before")
@classmethod
def parse_and_validate_amounts(cls, v, validation_info: ValidationInfo):
field_name = validation_info.field_name
if v is None or v == "":
spread_field = field_name.replace('amounts_pct', 'spreads')
return [1 for _ in validation_info.data[spread_field]]
parsed = parse_comma_separated_list(v)
if isinstance(parsed, list) and len(parsed) != len(validation_info.data[field_name.replace('amounts_pct', 'spreads')]):
raise ValueError(
f"The number of {field_name} must match the number of {field_name.replace('amounts_pct', 'spreads')}.")
return parsed
@field_validator('position_mode', mode="before")
@classmethod
def validate_position_mode(cls, v) -> PositionMode:
return parse_enum_value(PositionMode, v, "position_mode")
@field_validator('position_side', mode="before")
@classmethod
def validate_position_side(cls, v) -> TradeType:
if isinstance(v, TradeType):
return v
# Accept the enum's integer value (e.g. from a serialized/reloaded config)
if isinstance(v, int) or (isinstance(v, str) and v.isdigit()):
try:
return TradeType(int(v))
except ValueError:
raise ValueError(f"position_side must be BUY/LONG or SELL/SHORT, got {v}")
mapping = {"BUY": TradeType.BUY, "SELL": TradeType.SELL, "LONG": TradeType.BUY, "SHORT": TradeType.SELL}
upper = str(v).upper()
if upper in mapping:
return mapping[upper]
raise ValueError(f"position_side must be BUY/LONG or SELL/SHORT, got {v}")
@field_validator('global_tp_activation_from', mode="before")
@classmethod
def validate_tp_activation_from(cls, v):
valid = {"always", "min_base", "target_base"}
if v not in valid:
raise ValueError(f"global_tp_activation_from must be one of {valid}")
return v
@field_validator('global_sl_activation_from', mode="before")
@classmethod
def validate_sl_activation_from(cls, v):
valid = {"target_base", "max_base"}
if v not in valid:
raise ValueError(f"global_sl_activation_from must be one of {valid}")
return v
@field_validator('global_pnl_reference', mode="before")
@classmethod
def validate_pnl_reference(cls, v):
valid = {"position", "portfolio"}
if v not in valid:
raise ValueError(f"global_pnl_reference must be one of {valid}")
return v
@field_validator('price_distance_tolerance', 'refresh_tolerance', 'tolerance_scaling', mode="before")
@classmethod
def validate_tolerance_fields(cls, v, validation_info: ValidationInfo):
field_name = validation_info.field_name
if isinstance(v, str):
return Decimal(v)
if field_name == 'tolerance_scaling' and Decimal(str(v)) <= 0:
raise ValueError(f"{field_name} must be greater than 0")
return v
@property
def is_short(self) -> bool:
return self.position_side == TradeType.SELL
@property
def triple_barrier_config(self) -> TripleBarrierConfig:
# Ensure we're passing OrderType enum values, not strings
open_order_type = self.open_order_type if isinstance(self.open_order_type, OrderType) else OrderType.LIMIT_MAKER
take_profit_order_type = self.take_profit_order_type if isinstance(self.take_profit_order_type, OrderType) else OrderType.LIMIT_MAKER
return TripleBarrierConfig(
take_profit=self.take_profit,
trailing_stop=None,
open_order_type=open_order_type,
take_profit_order_type=take_profit_order_type,
stop_loss_order_type=OrderType.MARKET,
time_limit_order_type=OrderType.MARKET
)
def get_cooldown_time(self, trade_type: TradeType) -> int:
"""Get cooldown time for specific trade type"""
return self.buy_cooldown_time if trade_type == TradeType.BUY else self.sell_cooldown_time
def get_position_effectivization_time(self, trade_type: TradeType) -> int:
"""Get position effectivization time for specific trade type"""
return self.buy_position_effectivization_time if trade_type == TradeType.BUY else self.sell_position_effectivization_time
def get_price_distance_level_tolerance(self, level: int) -> Decimal:
"""Get level-specific price distance tolerance (for new order placement).
Prevents placing new orders when existing ones are too close to current price.
"""
return self.price_distance_tolerance * (self.tolerance_scaling ** level)
def get_refresh_level_tolerance(self, level: int) -> Decimal:
"""Get level-specific refresh tolerance (for order replacement).
Triggers replacing open orders when price deviates from theoretical level.
"""
return self.refresh_tolerance * (self.tolerance_scaling ** level)
def update_parameters(self, trade_type: TradeType, new_spreads: Union[List[float], str],
new_amounts_pct: Optional[Union[List[int], str]] = None):
spreads_field = 'buy_spreads' if trade_type == TradeType.BUY else 'sell_spreads'
amounts_pct_field = 'buy_amounts_pct' if trade_type == TradeType.BUY else 'sell_amounts_pct'
setattr(self, spreads_field, self.parse_spreads(new_spreads))
if new_amounts_pct is not None:
setattr(self, amounts_pct_field,
self.parse_and_validate_amounts(new_amounts_pct, self.__dict__, self.__fields__[amounts_pct_field]))
else:
setattr(self, amounts_pct_field, [1 for _ in getattr(self, spreads_field)])
def get_spreads_and_amounts_in_quote(self, trade_type: TradeType) -> Tuple[List[float], List[float]]:
buy_amounts_pct = getattr(self, 'buy_amounts_pct')
sell_amounts_pct = getattr(self, 'sell_amounts_pct')
total_pct = sum(buy_amounts_pct) + sum(sell_amounts_pct)
if trade_type == TradeType.BUY:
normalized_amounts_pct = [amt_pct / total_pct for amt_pct in buy_amounts_pct]
else:
normalized_amounts_pct = [amt_pct / total_pct for amt_pct in sell_amounts_pct]
spreads = getattr(self, f'{trade_type.name.lower()}_spreads')
return spreads, [amt_pct * self.total_amount_quote * self.portfolio_allocation for amt_pct in normalized_amounts_pct]
def update_markets(self, markets: MarketDict) -> MarketDict:
return markets.add_or_update(self.connector_name, self.trading_pair)
class PMMister(ControllerBase):
"""
Advanced PMM (Pure Market Making) controller with sophisticated position management.
Features:
- Hanging executors system for better position control
- Price distance requirements to prevent over-accumulation
- Breakeven awareness for dynamic parameter adjustment
- Separate buy/sell cooldown and effectivization times
"""
def __init__(self, config: PMMisterConfig, *args, **kwargs):
super().__init__(config, *args, **kwargs)
self.config = config
self.market_data_provider.initialize_rate_sources(
[ConnectorPair(connector_name=config.connector_name, trading_pair=config.trading_pair)]
)
self.price_history = []
self.max_price_history = 60
self.order_history = []
self.max_order_history = 20
self.processed_data = {}
self._position_mode_verified = False
self._global_close_phase: Optional[str] = None # None | "stopping" | "closing"
self._global_close_side: Optional[TradeType] = None # Side of the position when TP/SL triggered
self._global_close_retries: int = 0 # Count how many times PHASE 2 has created a close executor
self._global_close_cooling_down: bool = False # True after a successful close until processed_data confirms 0
def _verify_position_mode(self) -> bool:
"""Check that the connector's position mode matches the config. Blocks trading until confirmed."""
if self._position_mode_verified:
return True
try:
connector = self.market_data_provider.get_connector(self.config.connector_name)
# Only perpetual connectors have position_mode; skip check for spot connectors
if not hasattr(connector, 'position_mode'):
self._position_mode_verified = True
return True
exchange_mode = connector.position_mode
config_mode = self.config.position_mode
if exchange_mode != config_mode:
self.logger().warning(
f"Position mode mismatch: exchange={exchange_mode}, config={config_mode}. "
f"Waiting for position mode to be set correctly before trading.")
return False
self._position_mode_verified = True
self.logger().info(
f"Position mode verified: {exchange_mode} matches config. Trading enabled.")
return True
except Exception as e:
self.logger().warning(f"Could not verify position mode: {e}. Blocking trading.")
return False
# ── Market data (called by framework) ─────────────────────────────────
async def update_processed_data(self):
"""Compute reference price and spread multiplier only. All executor analysis
is done in _compute_executor_analysis called from determine_executor_actions."""
try:
reference_price = self.market_data_provider.get_price_by_type(
self.config.connector_name, self.config.trading_pair, PriceType.MidPrice
)
if reference_price is None or reference_price <= 0:
self.logger().warning("Invalid reference price received, using previous price if available")
reference_price = self.processed_data.get("reference_price", Decimal("100"))
except Exception as e:
self.logger().warning(f"Error getting reference price: {e}, using previous price if available")
reference_price = self.processed_data.get("reference_price", Decimal("100"))
current_time = self.market_data_provider.time()
self.price_history.append({'timestamp': current_time, 'price': Decimal(reference_price)})
if len(self.price_history) > self.max_price_history:
self.price_history.pop(0)
if self.config.tick_mode:
spread_multiplier = (self.market_data_provider.get_trading_rules(
self.config.connector_name, self.config.trading_pair
).min_price_increment / reference_price)
else:
spread_multiplier = Decimal("1")
self.processed_data = {
"reference_price": Decimal(reference_price),
"spread_multiplier": spread_multiplier,
}
# ── Executor actions (called by framework) ────────────────────────────
def determine_executor_actions(self) -> List[ExecutorAction]:
# Guard: verify position mode matches config before operating
if not self._verify_position_mode():
return []
self._update_position_state()
self._compute_executor_analysis()
actions = []
# Check global TP/SL — two-phase: stop executors, then close position
tp_sl_actions = self._check_global_tp_sl()
if tp_sl_actions:
actions.extend(tp_sl_actions)
# Block normal trading while global close is in progress
if self._global_close_phase is not None:
return actions
actions.extend(self.create_actions_proposal())
actions.extend(self.stop_actions_proposal())
return actions
# ── Global TP/SL ──────────────────────────────────────────────────────
def _get_tp_activation_threshold(self) -> Decimal:
if self.config.global_tp_activation_from == "always":
return Decimal("0")
if self.config.global_tp_activation_from == "min_base":
return self.config.min_base_pct
return self.config.target_base_pct
def _get_sl_activation_threshold(self) -> Decimal:
if self.config.global_sl_activation_from == "target_base":
return self.config.target_base_pct
return self.config.max_base_pct
def _get_exchange_position(self) -> Tuple[Decimal, Optional[TradeType]]:
"""Read the REAL position from the exchange connector (WebSocket-updated, no orchestrator delay).
Returns (abs_amount, side) where side is BUY for long, SELL for short, None if no position."""
try:
connector = self.market_data_provider.get_connector(self.config.connector_name)
if not hasattr(connector, '_perpetual_trading'):
return Decimal("0"), None
perp = connector._perpetual_trading
pos = perp.get_position(self.config.trading_pair, PositionSide.BOTH)
if pos is None or pos.amount == Decimal("0"):
return Decimal("0"), None
amount = pos.amount
# Binance ONEWAY: positive amount = long, negative = short
if amount > 0:
return amount, TradeType.BUY
else:
return abs(amount), TradeType.SELL
except Exception as e:
self.logger().warning(f"Failed to read exchange position: {e}")
return Decimal("0"), None
def _check_global_tp_sl(self) -> List[ExecutorAction]:
"""Check global TP/SL using a two-phase approach:
Phase 1 (stopping): Stop all active executors with keep_position=True.
Phase 2 (closing): Once no active executors remain, close the actual position."""
# --- Phase: stopping --- wait for all executors to finish, then transition to closing
if self._global_close_phase == "stopping":
active_non_close = [
e for e in self.executors_info
if e.is_active and e.custom_info.get("level_id") != "global_close"
]
if active_non_close:
self.logger().debug(
f"Global close phase=stopping: waiting for {len(active_non_close)} executors to finish"
)
return []
# All executors stopped — transition to closing phase
self._global_close_phase = "closing"
self.logger().info("Global close phase=stopping complete. All executors stopped. Transitioning to closing.")
# --- Phase: closing --- create close executor for the real position
if self._global_close_phase == "closing":
# If a close executor is already active, wait for it
close_executors = [
e for e in self.executors_info
if e.is_active and e.custom_info.get("level_id") == "global_close"
]
if close_executors:
return []
# Guard: abort after too many failed close attempts (e.g. below min notional)
if self._global_close_retries >= 3:
self.logger().warning(
f"=== GLOBAL CLOSE ABORTED: {self._global_close_retries} close attempts failed. ===\n"
f" Position may be below minimum notional. Aborting to prevent infinite loop."
)
self._global_close_phase = None
self._global_close_side = None
self._global_close_retries = 0
return []
# Read position from the EXCHANGE CONNECTOR (WebSocket-updated, no orchestrator delay)
# This avoids the race condition where positions_held is stale
exchange_amount, exchange_side = self._get_exchange_position()
if exchange_amount == Decimal("0") or exchange_side is None:
self.logger().info("Global close phase=closing: exchange position is 0. Done.")
self._global_close_phase = None
self._global_close_side = None
self._global_close_retries = 0
self._global_close_cooling_down = True
return []
# SAFETY: Detect position side flip — if position flipped direction, abort close
if self._global_close_side is not None and exchange_side != self._global_close_side:
self.logger().warning(
f"=== GLOBAL CLOSE ABORTED: Position side flipped! ===\n"
f" Original side: {self._global_close_side.name} | "
f"Current side: {exchange_side.name} | Amount: {exchange_amount}\n"
f" This indicates over-selling. Aborting global close to prevent further damage."
)
self._global_close_phase = None
self._global_close_side = None
self._global_close_retries = 0
return []
quantized = self.market_data_provider.quantize_order_amount(
self.config.connector_name, self.config.trading_pair, exchange_amount
)
if quantized == Decimal("0"):
self._global_close_phase = None
self._global_close_side = None
self._global_close_retries = 0
return []
# Determine close side from the EXCHANGE position side
close_side = TradeType.SELL if exchange_side == TradeType.BUY else TradeType.BUY
self._global_close_retries += 1
self.logger().info(
f"=== GLOBAL CLOSE — PHASE 2: CLOSING POSITION (attempt {self._global_close_retries}/3) ===\n"
f" Exchange position: {exchange_side.name} {exchange_amount} | "
f"Close side: {close_side.name} | Creating close executor."
)
close_action = self._create_close_action_with_side(close_side, exchange_amount)
return [close_action] if close_action else []
# --- No phase active: check if TP/SL should trigger ---
current_base_pct = self.processed_data.get("current_base_pct", Decimal("0"))
unrealized_pnl_pct = self.processed_data.get("unrealized_pnl_pct", Decimal("0"))
position_amount = self.processed_data.get("position_amount", Decimal("0"))
if position_amount == Decimal("0"):
self._global_close_cooling_down = False
return []
# After a successful close, processed_data can lag behind the exchange by one tick.
# Suppress re-triggering until the exchange also confirms no position remains.
# If the exchange already shows a new non-zero position, a genuinely new position
# has opened and the cooldown no longer applies.
if self._global_close_cooling_down:
exchange_amount, _ = self._get_exchange_position()
if exchange_amount > Decimal("0"):
self._global_close_cooling_down = False
else:
return []
triggered = False
trigger_reason = ""
# Check take profit
tp_threshold = self._get_tp_activation_threshold()
if self.config.global_tp_enabled and current_base_pct >= tp_threshold:
if unrealized_pnl_pct >= self.config.global_take_profit:
triggered = True
trigger_reason = "take_profit"
# Check stop loss
sl_threshold = self._get_sl_activation_threshold()
if not triggered and self.config.global_sl_enabled and current_base_pct >= sl_threshold:
if unrealized_pnl_pct <= -self.config.global_stop_loss:
triggered = True
trigger_reason = "stop_loss"
if not triggered:
return []
# --- Trigger: enter stopping phase --- stop all active executors first
self._global_close_phase = "stopping"
self._global_close_retries = 0
# Remember the position side at trigger time so we always close in the right direction
position_held = next((p for p in self.positions_held if
p.trading_pair == self.config.trading_pair and
p.connector_name == self.config.connector_name), None)
self._global_close_side = position_held.side if position_held else None
active_executors = [
e for e in self.executors_info
if e.is_active and e.custom_info.get("level_id") != "global_close"
]
self.logger().info(
f"=== GLOBAL {trigger_reason.upper()} TRIGGERED — PHASE 1: STOPPING EXECUTORS ===\n"
f" PnL: {unrealized_pnl_pct:.4%} | Position: {position_amount} | Base%: {current_base_pct:.4%}\n"
f" Stopping {len(active_executors)} active executors before closing position."
)
stop_actions = []
for executor in active_executors:
stop_actions.append(StopExecutorAction(
controller_id=self.config.id,
keep_position=True,
executor_id=executor.id,
))
return stop_actions
def _create_close_action(self, position_amount: Decimal) -> Optional[CreateExecutorAction]:
"""Create a close action by inferring the side from position_held. Kept for backward compat."""
position_held = next((p for p in self.positions_held if
p.trading_pair == self.config.trading_pair and
p.connector_name == self.config.connector_name), None)
if position_held is None or position_amount == Decimal("0"):
return None
close_side = TradeType.SELL if position_held.side == TradeType.BUY else TradeType.BUY
return self._create_close_action_with_side(close_side, abs(position_amount))
def _create_close_action_with_side(self, side: TradeType, amount: Decimal) -> Optional[CreateExecutorAction]:
if amount == Decimal("0"):
return None
self.logger().info(
f"Creating close executor: side={side.name} amount={amount} "
f"action=CLOSE strategy=MARKET (reduceOnly on exchange)"
)
config = OrderExecutorConfig(
timestamp=self.market_data_provider.time(),
trading_pair=self.config.trading_pair,
connector_name=self.config.connector_name,
side=side,
amount=amount,
execution_strategy=ExecutionStrategy.MARKET,
position_action=PositionAction.CLOSE,
leverage=self.config.leverage,
level_id="global_close",
)
return CreateExecutorAction(
controller_id=self.config.id,
executor_config=config,
)
# ── Single-pass executor analysis ─────────────────────────────────────
def _compute_executor_analysis(self):
"""Analyse every executor and level once per tick. Results are stored
in self.processed_data and consumed by create/stop proposals and status display."""
current_time = self.market_data_provider.time()
reference_price = Decimal(str(self.processed_data.get("reference_price", 0)))
if reference_price <= 0:
return
# -- 1. Group executors by level_id in a single pass -----------------
executors_by_level: Dict[str, list] = defaultdict(list)
for e in self.executors_info:
level_id = e.custom_info.get("level_id")
if level_id:
executors_by_level[level_id].append(e)
# All configured levels (may not have executors yet)
all_level_ids = set()
for i in range(len(self.config.buy_spreads)):
all_level_ids.add(f"buy_{i}")
for i in range(len(self.config.sell_spreads)):
all_level_ids.add(f"sell_{i}")
all_level_ids.update(executors_by_level.keys())
# -- 2. Per-level analysis + blocking conditions ----------------------
levels_analysis: Dict[str, Dict] = {}
level_conditions: Dict[str, Dict] = {}
working_levels = set()
cooldown_status = {
"buy": {"active": False, "remaining_time": 0, "progress_pct": Decimal("0")},
"sell": {"active": False, "remaining_time": 0, "progress_pct": Decimal("0")},
}
current_pct = self.processed_data.get("current_base_pct", Decimal("0"))
breakeven_price = self.processed_data.get("breakeven_price")
for level_id in all_level_ids:
if not level_id.startswith(("buy_", "sell_")):
continue
executors = executors_by_level.get(level_id, [])
active = [e for e in executors if e.is_active]
active_not_trading = [e for e in active if not e.is_trading]
active_trading = [e for e in active if e.is_trading]
open_order_updates = [
e.custom_info.get("open_order_last_update") for e in executors
if e.custom_info.get("open_order_last_update") is not None
]
latest_update = max(open_order_updates) if open_order_updates else None
prices = [Decimal(str(e.config.entry_price)) for e in active if hasattr(e.config, 'entry_price')]
analysis = {
"active_not_trading": active_not_trading,
"active_trading": active_trading,
"total_active": len(active),
"open_order_last_update": latest_update,
"min_price": min(prices) if prices else None,
"max_price": max(prices) if prices else None,
}
levels_analysis[level_id] = analysis
trade_type = self.get_trade_type_from_level_id(level_id)
is_buy = level_id.startswith("buy")
level = self.get_level_from_level_id(level_id)
blocking: List[str] = []
# a) Has open (not yet filled) executors
if active_not_trading:
blocking.append("active_not_trading")
# b) Max executor cap reached
if analysis["total_active"] >= self.config.max_active_executors_by_level:
blocking.append("max_active_executors")
# c) Cooldown
if latest_update is not None:
cooldown_time = self.config.get_cooldown_time(trade_type)
time_since = current_time - latest_update
if time_since < cooldown_time:
blocking.append("cooldown")
# Track cooldown progress for display (keep the most recent)
side = "buy" if is_buy else "sell"
remaining = cooldown_time - time_since
progress = Decimal(str(time_since)) / Decimal(str(cooldown_time))
if not cooldown_status[side]["active"] or remaining > cooldown_status[side]["remaining_time"]:
cooldown_status[side].update(active=True, remaining_time=remaining, progress_pct=progress)
# d) Price distance violation
level_tolerance = self.config.get_price_distance_level_tolerance(level)
if is_buy and analysis["min_price"] is not None:
distance = (analysis["min_price"] - reference_price) / reference_price
if distance < level_tolerance:
blocking.append("price_distance")
elif not is_buy and analysis["max_price"] is not None:
distance = (reference_price - analysis["max_price"]) / reference_price
if distance < level_tolerance:
blocking.append("price_distance")
# e) Position constraints
is_accumulation = self._is_accumulation_side(trade_type)
if current_pct < self.config.min_base_pct and not is_accumulation:
blocking.append("below_min_position")
elif current_pct > self.config.max_base_pct and is_accumulation:
blocking.append("above_max_position")
# f) Position profit protection — block the reduction side when price is unfavorable
is_reduction = not is_accumulation
if self.config.position_profit_protection and is_reduction and breakeven_price and breakeven_price > 0:
if self.config.is_short:
# SHORT: buying to reduce — block if price > breakeven (would realize a loss)
if reference_price > breakeven_price:
blocking.append("position_profit_protection")
else:
# LONG: selling to reduce — block if price < breakeven (would realize a loss)
if reference_price < breakeven_price:
blocking.append("position_profit_protection")
# Execution-blocking conditions determine "working" levels
execution_blocking = {"active_not_trading", "max_active_executors", "cooldown", "price_distance"}
if any(b in execution_blocking for b in blocking):
working_levels.add(level_id)
level_conditions[level_id] = {
"trade_type": trade_type.name,
"can_execute": len(blocking) == 0,
"blocking_conditions": blocking,
"active_executors": len(active_not_trading),
"hanging_executors": len(active_trading),
}
# -- 3. Levels to execute (position-aware) ----------------------------
levels_to_execute = self._get_executable_levels(working_levels)
# -- 4. Executors to refresh + refresh tracking -----------------------
executors_to_refresh = []
refresh_tracking = {
"refresh_candidates": [], "near_refresh": 0,
"refresh_ready": 0, "distance_violations": 0,
}
for e in self.executors_info:
if not e.is_active or e.is_trading:
continue
age = current_time - e.timestamp
time_based = age > self.config.executor_refresh_time
distance_based = reference_price > 0 and self.should_refresh_executor_by_distance(e, reference_price)
if time_based or distance_based:
executors_to_refresh.append(e)
# Tracking data for display
time_to_refresh = max(0, self.config.executor_refresh_time - age)
progress = min(Decimal("1"), Decimal(str(age)) / Decimal(str(self.config.executor_refresh_time)))
ready = time_based or distance_based
near = time_to_refresh <= self.config.executor_refresh_time * 0.2
distance_deviation_pct = Decimal("0")
e_level_id = e.custom_info.get("level_id", "")
if e_level_id and hasattr(e.config, 'entry_price') and reference_price > 0:
theoretical = self.calculate_theoretical_price(e_level_id, reference_price)
if theoretical > 0:
distance_deviation_pct = abs(e.config.entry_price - theoretical) / theoretical
if ready:
refresh_tracking["refresh_ready"] += 1
elif near:
refresh_tracking["near_refresh"] += 1
if distance_based:
refresh_tracking["distance_violations"] += 1
e_level = self.get_level_from_level_id(e_level_id) if e_level_id else 0
refresh_tracking["refresh_candidates"].append({
"executor_id": e.id,
"level_id": e_level_id or "unknown",
"level": e_level,
"age": age,
"time_to_refresh": time_to_refresh,
"progress_pct": progress,
"ready": ready,
"ready_by_time": time_based,
"ready_by_distance": distance_based,
"distance_deviation_pct": distance_deviation_pct,
"distance_violation": distance_based,
"level_tolerance": self.config.get_refresh_level_tolerance(e_level),
"near_refresh": near,
})
# -- 5. Hanging executors to effectivize + tracking -------------------
executors_to_effectivize = []
effectivization_tracking = {
"hanging_executors": [], "total_hanging": 0, "ready_for_effectivization": 0,
}
for e in self.executors_info:
if not (e.is_active and e.is_trading):
continue
e_level_id = e.custom_info.get("level_id", "")
fill_time = e.custom_info.get("open_order_last_update")
if not e_level_id or fill_time is None:
continue
trade_type = self.get_trade_type_from_level_id(e_level_id)
eff_time = self.config.get_position_effectivization_time(trade_type)
elapsed = current_time - fill_time
remaining = max(0, eff_time - elapsed)
progress = min(Decimal("1"), Decimal(str(elapsed)) / Decimal(str(eff_time)))
ready = remaining == 0
if ready:
executors_to_effectivize.append(e)
effectivization_tracking["ready_for_effectivization"] += 1
effectivization_tracking["total_hanging"] += 1
effectivization_tracking["hanging_executors"].append({
"level_id": e_level_id,
"trade_type": trade_type.name,
"time_elapsed": elapsed,
"remaining_time": remaining,
"progress_pct": progress,
"ready": ready,
"executor_id": e.id,
})
# -- 6. Executor statistics -------------------------------------------
active_all = [e for e in self.executors_info if e.is_active]
total_trading = sum(1 for e in active_all if e.is_trading)
executor_stats = {
"total_active": len(active_all),
"total_trading": total_trading,
"total_not_trading": len(active_all) - total_trading,
}
# -- Store everything -------------------------------------------------
self.processed_data.update({
"levels_analysis": levels_analysis,
"level_conditions": level_conditions,
"levels_to_execute": levels_to_execute,
"executors_to_refresh": executors_to_refresh,
"executors_to_effectivize": executors_to_effectivize,
"cooldown_status": cooldown_status,
"effectivization_tracking": effectivization_tracking,
"refresh_tracking": refresh_tracking,
"executor_stats": executor_stats,
"current_time": current_time,
})
# ── Position state ────────────────────────────────────────────────────
def _update_position_state(self):
"""Recalculate position-derived fields (skews, deviation, breakeven) from positions_held."""
reference_price = self.processed_data.get("reference_price")
if reference_price is None:
return
position_held = next((p for p in self.positions_held if
p.trading_pair == self.config.trading_pair and
p.connector_name == self.config.connector_name), None)
target_position = self.config.total_amount_quote * self.config.target_base_pct
if position_held is not None:
# Use abs(amount_quote) so current_base_pct is always positive for both long and short
current_base_pct = abs(position_held.amount_quote) / self.config.total_amount_quote
deviation = (target_position - abs(position_held.amount_quote)) / target_position
breakeven_price = position_held.breakeven_price
position_amount = position_held.amount
position_cum_fees = position_held.cum_fees_quote
position_realized_pnl = position_held.realized_pnl_quote
position_unrealized_pnl = position_held.unrealized_pnl_quote
position_volume = position_held.volume_traded_quote
if self.config.global_pnl_reference == "portfolio":
pnl_denominator = self.config.total_amount_quote
else:
# Use entry value (breakeven * amount) for stable PnL % instead of mark-price based amount_quote
pnl_denominator = (abs(position_amount) * breakeven_price
if breakeven_price and breakeven_price > 0
else abs(position_held.amount_quote))
unrealized_pnl_pct = (position_held.unrealized_pnl_quote / pnl_denominator
if pnl_denominator != 0 else Decimal("0"))
else:
current_base_pct = Decimal("0")
deviation = Decimal("1")
unrealized_pnl_pct = Decimal("0")
breakeven_price = None
position_amount = Decimal("0")
position_cum_fees = Decimal("0")
position_realized_pnl = Decimal("0")
position_unrealized_pnl = Decimal("0")
position_volume = Decimal("0")
# Executor fees (from active executors)
executor_fees = sum(
(e.cum_fees_quote for e in self.executors_info if e.is_active),
Decimal("0")
)
min_pct = self.config.min_base_pct
max_pct = self.config.max_base_pct
if max_pct > min_pct:
if self.config.is_short:
# SHORT: sell accumulates → sell_skew high when position small, buy_skew high when position large
sell_skew = (max_pct - current_base_pct) / (max_pct - min_pct)
buy_skew = (current_base_pct - min_pct) / (max_pct - min_pct)
else:
# LONG: buy accumulates → buy_skew high when position small, sell_skew high when position large
buy_skew = (max_pct - current_base_pct) / (max_pct - min_pct)
sell_skew = (current_base_pct - min_pct) / (max_pct - min_pct)
buy_skew = max(min(buy_skew, Decimal("1.0")), self.config.min_skew)
sell_skew = max(min(sell_skew, Decimal("1.0")), self.config.min_skew)
else:
buy_skew = sell_skew = Decimal("1.0")
self.processed_data.update({
"deviation": deviation,
"current_base_pct": current_base_pct,
"unrealized_pnl_pct": unrealized_pnl_pct,
"breakeven_price": breakeven_price,
"position_amount": position_amount,
"buy_skew": buy_skew,
"sell_skew": sell_skew,
"position_cum_fees": position_cum_fees,
"position_realized_pnl": position_realized_pnl,
"position_unrealized_pnl": position_unrealized_pnl,
"position_volume": position_volume,
"executor_fees": executor_fees,
"total_fees": position_cum_fees + executor_fees,
})
# ── Create / stop proposals ───────────────────────────────────────────
def create_actions_proposal(self) -> List[ExecutorAction]:
create_actions = []
levels_to_execute = self.processed_data.get("levels_to_execute", [])
if not levels_to_execute:
return create_actions
buy_spreads, buy_amounts_quote = self.config.get_spreads_and_amounts_in_quote(TradeType.BUY)
sell_spreads, sell_amounts_quote = self.config.get_spreads_and_amounts_in_quote(TradeType.SELL)
reference_price = Decimal(self.processed_data["reference_price"])
buy_skew = self.processed_data["buy_skew"]
sell_skew = self.processed_data["sell_skew"]
for level_id in levels_to_execute:
trade_type = self.get_trade_type_from_level_id(level_id)
level = self.get_level_from_level_id(level_id)
if trade_type == TradeType.BUY:
spread_in_pct = Decimal(buy_spreads[level]) * Decimal(self.processed_data["spread_multiplier"])
amount_quote = Decimal(buy_amounts_quote[level])
else:
spread_in_pct = Decimal(sell_spreads[level]) * Decimal(self.processed_data["spread_multiplier"])
amount_quote = Decimal(sell_amounts_quote[level])
skew = buy_skew if trade_type == TradeType.BUY else sell_skew
side_multiplier = Decimal("-1") if trade_type == TradeType.BUY else Decimal("1")
price = reference_price * (Decimal("1") + side_multiplier * spread_in_pct)
amount = self.market_data_provider.quantize_order_amount(
self.config.connector_name,
self.config.trading_pair,
(amount_quote / price) * skew
)
if amount == Decimal("0"):
self.logger().warning(f"The amount of the level {level_id} is 0. Skipping.")
continue
# Position profit protection: block reduction-side orders at unfavorable prices
if self.config.position_profit_protection and not self._is_accumulation_side(trade_type):
breakeven_price = self.processed_data.get("breakeven_price")
if breakeven_price is not None and breakeven_price > 0:
# LONG reduces by selling → skip if price < breakeven
# SHORT reduces by buying → skip if price > breakeven
if self.config.is_short and price > breakeven_price:
continue
elif not self.config.is_short and price < breakeven_price:
continue
executor_config = self.get_executor_config(level_id, price, amount)
if executor_config is not None:
self.order_history.append({
'timestamp': self.market_data_provider.time(),
'price': price,
'side': trade_type.name,
'level_id': level_id,
'action': 'CREATE'
})
if len(self.order_history) > self.max_order_history:
self.order_history.pop(0)
create_actions.append(CreateExecutorAction(
controller_id=self.config.id,
executor_config=executor_config
))
return create_actions
def stop_actions_proposal(self) -> List[ExecutorAction]:
stop_actions = []
for executor in self.processed_data.get("executors_to_refresh", []):
stop_actions.append(StopExecutorAction(
controller_id=self.config.id,
keep_position=True,
executor_id=executor.id
))
for executor in self.processed_data.get("executors_to_effectivize", []):
stop_actions.append(StopExecutorAction(
controller_id=self.config.id,
keep_position=True,
executor_id=executor.id
))
return stop_actions
# ── Helpers ───────────────────────────────────────────────────────────
def _get_executable_levels(self, working_levels: set) -> List[str]:
"""Get levels that should be executed, applying position constraints."""
buy_missing = [
f"buy_{i}" for i in range(len(self.config.buy_spreads))
if f"buy_{i}" not in working_levels
]
sell_missing = [
f"sell_{i}" for i in range(len(self.config.sell_spreads))
if f"sell_{i}" not in working_levels
]
# Determine which side accumulates vs reduces based on position_side
if self.config.is_short:
accumulation_levels = sell_missing
reduction_levels = buy_missing
else:
accumulation_levels = buy_missing
reduction_levels = sell_missing
current_pct = self.processed_data.get("current_base_pct", Decimal("0"))
# Below min → only accumulate
if current_pct < self.config.min_base_pct:
return accumulation_levels
# Above max → only reduce
elif current_pct > self.config.max_base_pct:
return reduction_levels
if self.config.position_profit_protection:
breakeven_price = self.processed_data.get("breakeven_price")
reference_price = self.processed_data["reference_price"]
target_pct = self.config.target_base_pct
if breakeven_price is not None and breakeven_price > 0:
if self.config.is_short:
# SHORT: below target & price above breakeven → only accumulate (sell more)
if current_pct < target_pct and reference_price > breakeven_price:
return accumulation_levels
# SHORT: above target & price below breakeven → only reduce (buy back)
elif current_pct > target_pct and reference_price < breakeven_price:
return reduction_levels
else:
# LONG: below target & price below breakeven → only accumulate (buy more)
if current_pct < target_pct and reference_price < breakeven_price:
return accumulation_levels
# LONG: above target & price above breakeven → only reduce (sell)
elif current_pct > target_pct and reference_price > breakeven_price:
return reduction_levels
return buy_missing + sell_missing
def calculate_theoretical_price(self, level_id: str, reference_price: Decimal) -> Decimal:
"""Calculate the theoretical price for a given level"""
trade_type = self.get_trade_type_from_level_id(level_id)
level = self.get_level_from_level_id(level_id)
spreads = self.config.buy_spreads if trade_type == TradeType.BUY else self.config.sell_spreads
if level >= len(spreads):
return reference_price
spread_in_pct = Decimal(spreads[level]) * Decimal(self.processed_data.get("spread_multiplier", 1))
side_multiplier = Decimal("-1") if trade_type == TradeType.BUY else Decimal("1")
return reference_price * (Decimal("1") + side_multiplier * spread_in_pct)
def should_refresh_executor_by_distance(self, executor_info, reference_price: Decimal) -> bool:
"""Check if executor should be refreshed due to price distance deviation"""
level_id = executor_info.custom_info.get("level_id", "")
if not level_id or not hasattr(executor_info.config, 'entry_price'):
return False
theoretical_price = self.calculate_theoretical_price(level_id, reference_price)
if theoretical_price == 0:
return False
distance_deviation = abs(executor_info.config.entry_price - theoretical_price) / theoretical_price
level = self.get_level_from_level_id(level_id)
return distance_deviation > self.config.get_refresh_level_tolerance(level)
def get_executor_config(self, level_id: str, price: Decimal, amount: Decimal):
trade_type = self.get_trade_type_from_level_id(level_id)
return PositionExecutorConfig(
timestamp=self.market_data_provider.time(),
level_id=level_id,
connector_name=self.config.connector_name,
trading_pair=self.config.trading_pair,
entry_price=price,
amount=amount,
triple_barrier_config=self.config.triple_barrier_config,
leverage=self.config.leverage,
side=trade_type,
)
def _is_accumulation_side(self, trade_type: TradeType) -> bool:
"""Returns True if trade_type is the side that accumulates position.
LONG: BUY accumulates. SHORT: SELL accumulates."""
if self.config.is_short:
return trade_type == TradeType.SELL
return trade_type == TradeType.BUY
def get_level_id_from_side(self, trade_type: TradeType, level: int) -> str:
return f"{trade_type.name.lower()}_{level}"
def get_trade_type_from_level_id(self, level_id: str) -> TradeType:
return TradeType.BUY if level_id.startswith("buy") else TradeType.SELL
def get_level_from_level_id(self, level_id: str) -> int:
parts = level_id.split('_')
try:
return int(parts[1])
except (ValueError, IndexError):
return -1
# ── Custom info (MQTT / broker) ──────────────────────────────────────
def get_custom_info(self) -> dict:
if not self.processed_data:
return {}
reference_price = self.processed_data.get("reference_price", Decimal("0"))
position_amount = self.processed_data.get("position_amount", Decimal("0"))
current_base_pct = self.processed_data.get("current_base_pct", Decimal("0"))
unrealized_pnl_pct = self.processed_data.get("unrealized_pnl_pct", Decimal("0"))
breakeven_price = self.processed_data.get("breakeven_price")
buy_skew = self.processed_data.get("buy_skew", Decimal("1"))
sell_skew = self.processed_data.get("sell_skew", Decimal("1"))
executor_stats = self.processed_data.get("executor_stats", {})
level_conditions = self.processed_data.get("level_conditions", {})
# Distance to global TP/SL
distance_to_tp = float(self.config.global_take_profit - unrealized_pnl_pct)
distance_to_sl = float(unrealized_pnl_pct + self.config.global_stop_loss)
# Executable levels count
can_buy = sum(1 for lc in level_conditions.values() if lc.get("trade_type") == "BUY" and lc.get("can_execute"))
can_sell = sum(1 for lc in level_conditions.values() if lc.get("trade_type") == "SELL" and lc.get("can_execute"))
return {
"reference_price": float(reference_price),
"position_amount": float(position_amount),
"current_base_pct": float(current_base_pct),
"unrealized_pnl_pct": float(unrealized_pnl_pct),
"breakeven_price": float(breakeven_price) if breakeven_price is not None else None,
"buy_skew": float(buy_skew),
"sell_skew": float(sell_skew),
"distance_to_tp": distance_to_tp,
"distance_to_sl": distance_to_sl,
"global_tp_enabled": self.config.global_tp_enabled,
"global_sl_enabled": self.config.global_sl_enabled,
"global_close_phase": self._global_close_phase,
"closing_position": self._global_close_phase is not None,
"active_executors": executor_stats.get("total_active", 0),
"trading_executors": executor_stats.get("total_trading", 0),
"executable_buy_levels": can_buy,
"executable_sell_levels": can_sell,
}
# ── Status display ────────────────────────────────────────────────────
def to_format_status(self) -> List[str]:
from decimal import Decimal
from itertools import zip_longest
status = []
outer_width = 170
inner_width = outer_width - 4
if not hasattr(self, 'processed_data') or not self.processed_data:
status.append("╒" + "═" * inner_width + "╕")
status.append(f"│ {'Initializing controller... please wait':<{inner_width}} │")
status.append(f"╘{'═' * inner_width}╛")
return status
base_pct = self.processed_data.get('current_base_pct', Decimal("0"))
min_pct = self.config.min_base_pct
max_pct = self.config.max_base_pct
target_pct = self.config.target_base_pct
pnl = self.processed_data.get('unrealized_pnl_pct', Decimal('0'))
breakeven = self.processed_data.get('breakeven_price')
current_price = self.processed_data.get('reference_price', Decimal("0"))
buy_skew = self.processed_data.get('buy_skew', Decimal("1.0"))
sell_skew = self.processed_data.get('sell_skew', Decimal("1.0"))
cooldown_status = self.processed_data.get('cooldown_status', {})
effectivization = self.processed_data.get('effectivization_tracking', {})
level_conditions = self.processed_data.get('level_conditions', {})
executor_stats = self.processed_data.get('executor_stats', {})
refresh_tracking = self.processed_data.get('refresh_tracking', {})
levels_analysis = self.processed_data.get('levels_analysis', {})
col1_width = 28
col2_width = 35
col3_width = 28
col4_width = 25
col5_width = inner_width - col1_width - col2_width - col3_width - col4_width - 4
half_width = inner_width // 2 - 1
bar_width = inner_width - 25
# Header
status.append("╒" + "═" * inner_width + "╕")
header_line = (
f"{self.config.connector_name}:{self.config.trading_pair} @ {current_price:.2f} "
f"Alloc: {self.config.portfolio_allocation:.1%} "
f"Spread×{self.processed_data.get('spread_multiplier', Decimal('1')):.3f} "
f"Dist: {self.config.price_distance_tolerance:.4%} Ref: {self.config.refresh_tolerance:.4%} (×{self.config.tolerance_scaling}) "
f"Pos Protect: {'ON' if self.config.position_profit_protection else 'OFF'}"
)
status.append(f"│ {header_line:<{inner_width}} │")
# REAL-TIME CONDITIONS DASHBOARD
status.append(f"├{'─' * inner_width}┤")
status.append(f"│ {'🔄 REAL-TIME CONDITIONS DASHBOARD':<{inner_width}} │")
Shown in full with attribution under the source's licence. Licence: Apache-2.0
This summary was written by Stratmill's research agent from the original; it is not a copy of the source.