Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,20 @@ Risk policies are read-only plugins: built-in asset/cash/position/margin checks
first, policies cannot transmit orders, and missing mark/FX context or policy errors
produce explicit rejection codes instead of bypassing risk.

`resolve_a_share_replay_status` converts a complete point-in-time listing/tradability/price-limit
snapshot into the existing `open`, `suspended`, `limit_up`, `limit_down`, or `closed` QExec status.
It rejects contradictory or missing-typed flags and deliberately has no universe-membership input:
leaving an index or strategy universe is not evidence of delisting or a disposal transaction.

Deferred cash dividends use two explicit ledger events. `cash_dividend_entitlement` recognizes an
`assets:dividend_receivable` on the ex-date so NAV includes the evidenced entitlement without
making it spendable cash. `cash_dividend_payment` moves that exact receivable into cash on the true
payment date. The payment remains valid after the position is sold because entitlement is fixed on
the ex-date. The ledger rejects a payment without the matching entitlement or with a different
cash-per-share amount or receivable total, while a registered zero-holding entitlement remains a
valid zero-effect lifecycle event. Missing payment dates must be rejected by the upstream
corporate-action bridge.

## Deterministic replay

`DeterministicRunEngine.replay` sorts by availability/event time and stable stream
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ requires-python = ">=3.10"
dependencies = [
"pyarrow>=14.0",
"jsonschema>=4.20",
"quant-data-kit @ git+https://github.com/PureSaber/quant-data-kit.git@fbbad913c9e592fdcbf8b473c4230e2da8594cef",
"quant-data-kit @ git+https://github.com/PureSaber/quant-data-kit.git@104f1ef8a3b1278c0ea5420fadaa9d1d863ce726",
]

[project.optional-dependencies]
Expand Down
2 changes: 1 addition & 1 deletion requirements.lock
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ pytz==2026.3.post1
# via pandas
pyyaml==6.0.3
# via quant-data-kit
quant-data-kit @ git+https://github.com/PureSaber/quant-data-kit.git@fbbad913c9e592fdcbf8b473c4230e2da8594cef
quant-data-kit @ git+https://github.com/PureSaber/quant-data-kit.git@104f1ef8a3b1278c0ea5420fadaa9d1d863ce726
# via quant-execution (pyproject.toml)
referencing==0.37.0
# via
Expand Down
2 changes: 2 additions & 0 deletions src/quant_execution/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@
LinearPerpetualRule,
MarketState,
RuleBookRiskGate,
resolve_a_share_replay_status,
)
from quant_execution.schemas import (
ACCOUNT_SNAPSHOT_SCHEMA_ID,
Expand Down Expand Up @@ -143,6 +144,7 @@
"get_json_schema",
"load_stored_artifacts",
"remaining_quantity",
"resolve_a_share_replay_status",
"transition_order",
"validate_arrow_table",
"validate_json_record",
Expand Down
177 changes: 163 additions & 14 deletions src/quant_execution/ledger.py
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,7 @@ def reset(self, *, opened_at: datetime | None = None) -> None:
self._fill_trading_days: dict[str, date] = {}
self._position_lots: dict[str, list[tuple[date, Decimal]]] = {}
self._fill_close_allocations: dict[str, tuple[Decimal, Decimal]] = {}
self._dividend_entitlements: dict[tuple[str, str, date], tuple[Decimal, Decimal]] = {}
self._posting_cache: dict[
tuple[str, str, Decimal, str | None, Decimal | None, int], Posting
] = {}
Expand Down Expand Up @@ -331,6 +332,7 @@ def capture_state(self) -> dict[str, object]:
"fill_trading_days": self._fill_trading_days,
"position_lots": self._position_lots,
"fill_close_allocations": self._fill_close_allocations,
"dividend_entitlements": self._dividend_entitlements,
"posting_cache": self._posting_cache,
"fx": self._fx,
"fx_history": self._fx_history,
Expand Down Expand Up @@ -360,6 +362,7 @@ def _restore_captured_state(self, state: dict[str, object]) -> None:
self._fill_trading_days = restored["fill_trading_days"]
self._position_lots = restored["position_lots"]
self._fill_close_allocations = restored["fill_close_allocations"]
self._dividend_entitlements = restored["dividend_entitlements"]
self._posting_cache = restored["posting_cache"]
self._fx = restored["fx"]
self._fx_history = restored["fx_history"]
Expand Down Expand Up @@ -504,6 +507,42 @@ def convert_to_base(
def cash_balance(self, currency: str) -> Decimal:
return self._accounts.get(("assets:cash", currency, None), Decimal(0))

def dividend_receivable_balance(
self,
currency: str,
*,
instrument_id: str | None = None,
) -> Decimal:
"""Return declared cash dividends that have not reached their payment date."""

currency = _currency(currency)
return sum(
(
amount
for (account, entry_currency, receivable_key), amount in self._accounts.items()
if account == "assets:dividend_receivable"
and entry_currency == currency
and (
instrument_id is None
or (
receivable_key is not None
and receivable_key.startswith(f"{instrument_id}@")
)
)
),
Decimal(0),
)

def _dividend_receivable_value(self, event_time: datetime) -> Decimal:
return sum(
(
self._to_base(amount, currency, event_time)
for (account, currency, _), amount in self._accounts.items()
if account == "assets:dividend_receivable"
),
Decimal(0),
)

@property
def has_open_derivative_position(self) -> bool:
return any(
Expand All @@ -525,6 +564,7 @@ def risk_balances(
(self._to_base(amount, currency, event_time) for currency, amount in cash.items()),
Decimal(0),
)
nav += self._dividend_receivable_value(event_time)
initial_margin = Decimal(0)
for instrument_id, quantity in self._positions.items():
spec = self._spec(instrument_id)
Expand Down Expand Up @@ -733,6 +773,7 @@ def liquidation_required(self, event_time: datetime | None = None) -> bool:
),
Decimal(0),
)
nav += self._dividend_receivable_value(at)
maintenance_margin = Decimal(0)
for instrument_id, quantity in self._positions.items():
spec = self._spec(instrument_id)
Expand Down Expand Up @@ -839,8 +880,10 @@ def _apply(
event.event_time,
event.settlement_id,
)
elif isinstance(event, CorporateActionEvent) and event.ratio is not None:
self._apply_split_state(event)
elif isinstance(event, CorporateActionEvent):
self._apply_dividend_state(event)
if event.ratio is not None:
self._apply_split_state(event)
if not trusted_unique:
self._event_fingerprints[reference_id] = (
self._event_fingerprint(event) if self._artifact_sink is not None else event
Expand Down Expand Up @@ -878,6 +921,12 @@ def _capture_apply_undo(
and posting.quantity_delta is not None
}
instrument_id = getattr(event, "instrument_id", None)
entitlement_key = (
self._dividend_entitlement_key(event)
if isinstance(event, CorporateActionEvent)
and event.action_type in {"cash_dividend_entitlement", "cash_dividend_payment"}
else None
)
return {
"transaction_count": len(self._transactions),
"transaction_key": transaction.idempotency_key in self._transaction_keys,
Expand Down Expand Up @@ -907,6 +956,12 @@ def _capture_apply_undo(
"mark": (
self._marks.get(instrument_id, _MISSING) if instrument_id is not None else _MISSING
),
"dividend_entitlement_key": entitlement_key,
"dividend_entitlement": (
self._dividend_entitlements.get(entitlement_key, _MISSING)
if entitlement_key is not None
else _MISSING
),
"event_time": self._event_time,
}

Expand Down Expand Up @@ -948,6 +1003,13 @@ def _rollback_apply(
undo["position_lots"],
)
self._restore_value(self._marks, instrument_id, undo["mark"])
entitlement_key = undo["dividend_entitlement_key"]
if entitlement_key is not None:
self._restore_value(
self._dividend_entitlements,
entitlement_key,
undo["dividend_entitlement"],
)
self._event_time = undo["event_time"]

@staticmethod
Expand Down Expand Up @@ -1086,6 +1148,7 @@ def snapshot(self, event_time: datetime | None = None) -> AccountSnapshot:
nav = sum(
(self._to_base(value, currency, at) for currency, value in cash.items()), Decimal(0)
)
nav += self._dividend_receivable_value(at)
initial_margin = Decimal(0)
maintenance_margin = Decimal(0)
for instrument_id, quantity in sorted(self._positions.items()):
Expand Down Expand Up @@ -1142,6 +1205,7 @@ def assert_nav_residual(self, snapshot: AccountSnapshot) -> None:
at = snapshot.event_time
for currency, balance in snapshot.cash_balances.items():
expected += self._to_base(decimal(balance), currency, at)
expected += self._dividend_receivable_value(at)
for instrument_id, quantity_fp in snapshot.positions.items():
spec = self._spec(instrument_id)
quantity = decimal(quantity_fp)
Expand Down Expand Up @@ -1222,8 +1286,60 @@ def _apply_split_state(self, event: CorporateActionEvent) -> None:
event.event_id,
)

@staticmethod
def _dividend_entitlement_key(event: CorporateActionEvent) -> tuple[str, str, date]:
assert event.currency is not None
return (event.instrument_id, str(event.currency), event.effective_date)

@staticmethod
def _dividend_receivable_instrument(key: tuple[str, str, date]) -> str:
return f"{key[0]}@{key[2].isoformat()}"

def _apply_dividend_state(self, event: CorporateActionEvent) -> None:
if event.action_type not in {"cash_dividend_entitlement", "cash_dividend_payment"}:
return
key = self._dividend_entitlement_key(event)
if event.action_type == "cash_dividend_payment":
del self._dividend_entitlements[key]
return
assert event.cash_amount is not None
receivable_key = self._dividend_receivable_instrument(key)
total = self._accounts.get(
("assets:dividend_receivable", str(event.currency), receivable_key),
Decimal(0),
)
self._dividend_entitlements[key] = (decimal(event.cash_amount), total)

def _validate_corporate_action(self, event: CorporateActionEvent) -> None:
spec = self._spec(event.instrument_id)
if event.action_type in {"cash_dividend_entitlement", "cash_dividend_payment"}:
if event.cash_amount is None or event.currency is None:
raise ValidationError(f"{event.action_type} requires cash_amount and currency")
if decimal(event.cash_amount) < 0:
raise ValidationError("cash dividend amount must be non-negative")
key = self._dividend_entitlement_key(event)
declaration = self._dividend_entitlements.get(key)
if event.action_type == "cash_dividend_entitlement" and declaration is not None:
raise ValidationError("cash dividend entitlement is already registered")
if event.action_type == "cash_dividend_payment" and event.ratio is not None:
raise ValidationError("cash dividend payment cannot carry a share ratio")
if event.action_type == "cash_dividend_payment":
if declaration is None:
raise ValidationError("cash dividend payment has no registered entitlement")
declared_amount, declared_total = declaration
if decimal(event.cash_amount) != declared_amount:
raise ValidationError(
"cash dividend payment amount does not match the registered entitlement"
)
receivable_key = self._dividend_receivable_instrument(key)
ledger_total = self._accounts.get(
("assets:dividend_receivable", str(event.currency), receivable_key),
Decimal(0),
)
if ledger_total != declared_total:
raise ValidationError(
"cash dividend receivable does not match the registered entitlement"
)
if event.ratio is None:
return
ratio = decimal(event.ratio)
Expand Down Expand Up @@ -1492,18 +1608,51 @@ def _corporate_action_transaction(self, event: CorporateActionEvent) -> LedgerTr
quantity = self._positions.get(event.instrument_id, Decimal(0))
postings: list[Posting] = []
if event.cash_amount is not None:
cash_delta = quantity * decimal(event.cash_amount) * decimal(spec.contract_multiplier)
postings.extend(
[
self._posting("assets:cash", str(event.currency), cash_delta),
self._posting(
"income:corporate_action",
str(event.currency),
-cash_delta,
instrument_id=event.instrument_id,
),
]
)
currency = str(event.currency)
receivable_key = f"{event.instrument_id}@{event.effective_date.isoformat()}"
if event.action_type == "cash_dividend_payment":
entitlement_key = self._dividend_entitlement_key(event)
cash_delta = self._dividend_entitlements[entitlement_key][1]
postings.extend(
[
self._posting("assets:cash", currency, cash_delta),
self._posting(
"assets:dividend_receivable",
currency,
-cash_delta,
instrument_id=receivable_key,
),
]
)
else:
cash_delta = (
quantity * decimal(event.cash_amount) * decimal(spec.contract_multiplier)
)
asset_account = (
"assets:dividend_receivable"
if event.action_type == "cash_dividend_entitlement"
else "assets:cash"
)
postings.extend(
[
self._posting(
asset_account,
currency,
cash_delta,
instrument_id=(
receivable_key
if event.action_type == "cash_dividend_entitlement"
else None
),
),
self._posting(
"income:corporate_action",
currency,
-cash_delta,
instrument_id=event.instrument_id,
),
]
)
if event.ratio is not None:
quantity_delta = quantity * (decimal(event.ratio) - Decimal(1))
postings.extend(
Expand Down
42 changes: 42 additions & 0 deletions src/quant_execution/rules.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,48 @@ class MarketState:
status: str


def resolve_a_share_replay_status(
*,
listed: bool,
delisted: bool,
tradable: bool,
limit_up: bool,
limit_down: bool,
) -> str:
"""Map complete point-in-time A-share flags to one QExec market status.

Listing lifecycle and index/universe membership are intentionally separate.
Callers must never turn a universe exit into ``delisted=True``.
"""

flags = {
"listed": listed,
"delisted": delisted,
"tradable": tradable,
"limit_up": limit_up,
"limit_down": limit_down,
}
if any(type(value) is not bool for value in flags.values()):
raise ValidationError("A-share replay status flags must be booleans")
if delisted and listed:
raise ValidationError("a delisted instrument cannot remain listed")
if not listed and (tradable or limit_up or limit_down):
raise ValidationError("an unlisted instrument cannot be tradable or price-limited")
if not tradable and (limit_up or limit_down):
raise ValidationError("a non-tradable instrument cannot be at a tradable price limit")
if limit_up and limit_down:
raise ValidationError("an instrument cannot be both limit-up and limit-down")
if not listed:
return "closed"
if not tradable:
return "suspended"
if limit_up:
return "limit_up"
if limit_down:
return "limit_down"
return "open"


@dataclass(frozen=True, slots=True)
class _RiskAccountView:
account_id: str
Expand Down
Loading
Loading