dita_v2(test): enforce post-ack atomicity and suite integrity

This commit is contained in:
Codex
2026-07-13 10:03:06 +02:00
parent e49b959b2e
commit b13b54e769
7 changed files with 2041 additions and 1 deletions

View File

@@ -0,0 +1,36 @@
"""Build and verify the DITAv2 Rust artifact outside the live runner.
Usage:
python -m prod.clean_arch.dita_v2.build_native_artifact
python -m prod.clean_arch.dita_v2.build_native_artifact --rollback
"""
from __future__ import annotations
import argparse
from pathlib import Path
from .native_artifact import build_verified_artifact, rollback_artifact
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument(
"--target-dir",
type=Path,
default=Path("/root/.cargo/dita_v2_target"),
)
parser.add_argument("--rollback", action="store_true")
args = parser.parse_args()
crate_dir = Path(__file__).resolve().with_name("_rust_kernel")
result = (
rollback_artifact(args.target_dir, crate_dir)
if args.rollback
else build_verified_artifact(crate_dir, args.target_dir)
)
print(result)
return 0
if __name__ == "__main__":
raise SystemExit(main())

View File

@@ -0,0 +1,10 @@
"""DITAv2 pytest collection policy.
The files below are source generators or frozen backups, not executable tests.
Both ``_gen_test.py`` paths match pytest's default ``*_test.py`` pattern and
perform file generation at import time, so collecting them corrupts the test
boundary and can fail before the real suite loads.
"""
collect_ignore = ["_gen_test.py"]
collect_ignore_glob = ["_backup_20260530/*"]

File diff suppressed because it is too large Load Diff

View File

@@ -0,0 +1,216 @@
"""Execution-atomicity regressions for BingX post-ack observability.
Once BingX has accepted an order, telemetry is no longer allowed to alter the
execution result. A 2026-07-13 signature skew raised after the HTTP 200, which
the kernel interpreted as a submit failure and rolled back while the exchange
kept the filled position. These tests protect the adapter and real kernel FSM
boundaries against that entire failure class.
"""
from __future__ import annotations
import ast
import asyncio
import inspect
import itertools
import threading
from datetime import datetime, timezone
from pathlib import Path
from types import SimpleNamespace
import pytest
from prod.clean_arch.dita_v2.bingx_venue import BingxVenueAdapter
from prod.clean_arch.dita_v2.contracts import (
KernelCommandType,
KernelEventKind,
KernelIntent,
TradeSide,
TradeStage,
)
from prod.clean_arch.dita_v2.rust_backend import ExecutionKernel
def _intent(*, trade_id: str = "post-ack-1") -> KernelIntent:
return KernelIntent(
timestamp=datetime.now(timezone.utc),
intent_id=f"intent-{trade_id}",
trade_id=trade_id,
slot_id=0,
asset="TRX-USDT",
action=KernelCommandType.ENTER,
side=TradeSide.SHORT,
reason="post-ack-regression",
target_size=10.0,
leverage=1.0,
reference_price=100.0,
exit_leg_ratios=(1.0,),
metadata={},
)
def _filled_receipt() -> SimpleNamespace:
return SimpleNamespace(
status="FILLED",
order_id="venue-order-1",
client_order_id="client-order-1",
price=100.0,
quantity=10.0,
timestamp=datetime.now(timezone.utc),
raw_ack={
"status": "FILLED",
"orderId": "venue-order-1",
"clientOrderId": "client-order-1",
"executedQty": "10.0",
"avgPrice": "100.0",
},
)
class _SyncBackend:
def __init__(self) -> None:
self.submit_count = 0
def submit_intent(self, _legacy_intent):
self.submit_count += 1
return _filled_receipt()
class _AsyncBackend:
def __init__(self) -> None:
self.submit_count = 0
async def submit_intent(self, _legacy_intent):
self.submit_count += 1
return _filled_receipt()
def _venue(backend) -> BingxVenueAdapter:
venue = BingxVenueAdapter.__new__(BingxVenueAdapter)
venue.backend = backend
venue._event_seq = itertools.count(1)
venue._snap_lock = threading.Lock()
venue._snapshot_ready = threading.Event()
venue._snapshot_ready.set()
venue._last_snapshot = None
return venue
class _RaiseAfterAck:
"""Allow pre-submit telemetry, then emulate any post-ack telemetry defect."""
def __init__(self) -> None:
self.phases: list[str] = []
def __call__(self, **fields) -> None:
phase = str(fields["phase"])
self.phases.append(phase)
if phase == "submit:done":
raise TypeError("simulated post-ack telemetry signature skew")
def _assert_fill_trace(events) -> None:
kinds = [event.kind for event in events]
assert kinds == [KernelEventKind.ORDER_ACK, KernelEventKind.FULL_FILL]
assert KernelEventKind.ORDER_REJECT not in kinds
fill = events[1]
assert fill.venue_order_id == "venue-order-1"
assert fill.venue_client_id == "client-order-1"
assert fill.filled_size == pytest.approx(10.0)
assert fill.remaining_size == pytest.approx(0.0)
def test_sync_submit_returns_fill_trace_when_post_ack_telemetry_raises(monkeypatch):
backend = _SyncBackend()
venue = _venue(backend)
telemetry = _RaiseAfterAck()
monkeypatch.setattr(venue, "_publish_telemetry", telemetry)
events = venue.submit(_intent(trade_id="sync-adapter"))
assert backend.submit_count == 1
assert telemetry.phases == ["submit:start", "submit:done"]
_assert_fill_trace(events)
def test_async_submit_returns_fill_trace_when_post_ack_telemetry_raises(monkeypatch):
backend = _AsyncBackend()
venue = _venue(backend)
telemetry = _RaiseAfterAck()
monkeypatch.setattr(venue, "_publish_telemetry", telemetry)
events = asyncio.run(venue.submit_async(_intent(trade_id="async-adapter")))
assert backend.submit_count == 1
assert telemetry.phases == ["submit:start", "submit:done"]
_assert_fill_trace(events)
def test_sync_kernel_does_not_roll_back_acknowledged_fill(monkeypatch):
backend = _SyncBackend()
venue = _venue(backend)
telemetry = _RaiseAfterAck()
monkeypatch.setattr(venue, "_publish_telemetry", telemetry)
with ExecutionKernel(max_slots=1, venue=venue) as kernel:
outcome = kernel.process_intent(_intent(trade_id="sync-kernel"))
slot = kernel._get_slot(0)
assert outcome.accepted is True
assert outcome.state is TradeStage.POSITION_OPEN
assert slot.fsm_state is TradeStage.POSITION_OPEN
assert slot.trade_id == "sync-kernel"
assert slot.size == pytest.approx(10.0)
_assert_fill_trace(outcome.emitted_events)
def test_async_kernel_does_not_roll_back_acknowledged_fill(monkeypatch):
backend = _AsyncBackend()
venue = _venue(backend)
telemetry = _RaiseAfterAck()
monkeypatch.setattr(venue, "_publish_telemetry", telemetry)
async def exercise():
with ExecutionKernel(max_slots=1, venue=venue) as kernel:
outcome = await kernel.process_intent_async(_intent(trade_id="async-kernel"))
slot = kernel._get_slot(0)
return outcome, slot
outcome, slot = asyncio.run(exercise())
assert outcome.accepted is True
assert outcome.state is TradeStage.POSITION_OPEN
assert slot.fsm_state is TradeStage.POSITION_OPEN
assert slot.trade_id == "async-kernel"
assert slot.size == pytest.approx(10.0)
_assert_fill_trace(outcome.emitted_events)
def test_every_publish_telemetry_call_site_binds_to_live_signature():
"""Catch call/signature skew before any order path can execute it."""
source_path = Path(inspect.getsourcefile(BingxVenueAdapter) or "")
tree = ast.parse(source_path.read_text(encoding="utf-8"), filename=str(source_path))
calls = [
node
for node in ast.walk(tree)
if isinstance(node, ast.Call)
and isinstance(node.func, ast.Attribute)
and node.func.attr == "_publish_telemetry"
]
# Protect against a broken discovery predicate making this test vacuous.
assert len(calls) >= 11, f"expected at least 11 telemetry call sites, found {len(calls)}"
signature = inspect.signature(BingxVenueAdapter._publish_telemetry)
failures: list[str] = []
for call in calls:
if any(keyword.arg is None for keyword in call.keywords):
failures.append(f"line {call.lineno}: **kwargs prevents static binding")
continue
positional = [object() for _ in call.args]
keywords = {str(keyword.arg): object() for keyword in call.keywords}
try:
signature.bind(object(), *positional, **keywords)
except TypeError as exc:
failures.append(f"line {call.lineno}: {exc}")
assert not failures, "telemetry call/signature skew:\n" + "\n".join(failures)

View File

@@ -0,0 +1,281 @@
"""Phase-3 tests: freshness-gated E-anchored capital provider.
Mutation litmus: each guard is paired with a mutation.
If removing the guard does NOT turn the test RED, the test is decoration.
"""
from __future__ import annotations
import sys
import time
sys.path.insert(0, "/mnt/dolphinng5_predict")
from prod.clean_arch.dita_v2.contracts import AccountStateSnapshot
from prod.clean_arch.dita_v2.zinc_plane import InMemoryZincPlane
from prod.clean_arch.dita_v2.e_capital_provider import (
ACCOUNT_STALENESS_NS,
make_e_capital_provider,
make_e_capital_provider_from_snap,
)
class TestFreshLiveReturnsWallet:
"""When e_live, wallet_balance > 0, and age < staleness → returns wallet_balance."""
def test_fresh_live_returns_wallet(self) -> None:
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, available_margin=20000.0, used_margin=5000.0,
event_seq=1, mono_ns=time.monotonic_ns(), e_live=True, reconcile_ok=True,
))
capital = provider()
assert capital == 25000.0
def test_fresh_live_different_amount(self) -> None:
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
plane.publish_account(AccountStateSnapshot(
wallet_balance=50000.0, mono_ns=time.monotonic_ns(), e_live=True,
))
assert provider() == 50000.0
def test_fresh_live_tiny_positive(self) -> None:
"""wallet_balance=0.01 — positive, so returns it."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
plane.publish_account(AccountStateSnapshot(
wallet_balance=0.01, mono_ns=time.monotonic_ns(), e_live=True,
))
assert provider() == 0.01
def test_zero_wallet_returns_none(self) -> None:
"""wallet_balance=0.0 is NOT > 0, so returns None."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
plane.publish_account(AccountStateSnapshot(
wallet_balance=0.0, mono_ns=time.monotonic_ns(), e_live=True,
))
assert provider() is None
class TestStaleReturnsNone:
"""When age >= staleness → returns None.
Mutation: change >= to > or remove age check → RED."""
def test_past_stale_returns_none(self) -> None:
"""MUTATION: age >= ACCOUNT_STALENESS_NS check removed → provider returns
old wallet_balance instead of None."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
# Publish with very old mono_ns
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=time.monotonic_ns() - ACCOUNT_STALENESS_NS - 1,
e_live=True,
))
capital = provider()
assert capital is None, f"MUTATION: stale snapshot returned {capital} instead of None"
def test_exactly_at_stale_returns_none(self) -> None:
"""age == ACCOUNT_STALENESS_NS is stale (strict <, not <=)."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
exact_stale_ns = time.monotonic_ns() - ACCOUNT_STALENESS_NS
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=exact_stale_ns, e_live=True,
))
assert provider() is None, "exactly at staleness should be None (strict <)"
def test_just_before_stale_returns_wallet(self, monkeypatch) -> None:
"""age == staleness - 1 is fresh."""
from prod.clean_arch.dita_v2 import e_capital_provider
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
now_ns = time.monotonic_ns()
monkeypatch.setattr(e_capital_provider.time, "monotonic_ns", lambda: now_ns)
fresh_ns = now_ns - ACCOUNT_STALENESS_NS + 1
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=fresh_ns, e_live=True,
))
assert provider() == 25000.0
class TestNotLiveReturnsNone:
"""e_live is False → returns None."""
def test_not_live_returns_none(self) -> None:
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=time.monotonic_ns(), e_live=False,
))
assert provider() is None
def test_default_snapshot_returns_none(self) -> None:
"""Fresh plane, never published — default snapshot has e_live=False."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
assert provider() is None
def test_live_then_not_live(self) -> None:
"""After going from live to not-live, returns None."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=time.monotonic_ns(), e_live=True,
))
assert provider() == 25000.0
# Overwrite with not-live
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=time.monotonic_ns(), e_live=False,
))
assert provider() is None
class TestProviderEdgeCases:
"""Edge cases for the capital provider."""
def test_never_stamped_returns_none(self) -> None:
"""No publish ever — provider returns None."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
assert provider() is None
def test_corrupt_snapshot_still_safe(self) -> None:
"""If plane returns a corrupt/zero snapshot, provider handles it."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
# Publish a snapshot with negative wallet_balance
plane.publish_account(AccountStateSnapshot(
wallet_balance=-100.0, mono_ns=time.monotonic_ns(), e_live=True,
))
# Negative wallet_balance means <= 0, so returns None
assert provider() is None
def test_negative_wallet_not_accepted(self) -> None:
"""Negative wallet_balance returns None (not >0)."""
snap = AccountStateSnapshot(
wallet_balance=-100.0, mono_ns=0, e_live=True,
)
result = make_e_capital_provider_from_snap(snap, now_ns=0)
assert result is None
def test_custom_staleness(self) -> None:
"""Different staleness threshold works."""
plane = InMemoryZincPlane()
# 1-second staleness
provider = make_e_capital_provider(plane, staleness_ns=1_000_000_000)
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=time.monotonic_ns(), e_live=True,
))
assert provider() == 25000.0
def test_very_old_snapshot(self) -> None:
"""Very old mono_ns — returns None."""
plane = InMemoryZincPlane()
provider = make_e_capital_provider(plane)
plane.publish_account(AccountStateSnapshot(
wallet_balance=25000.0,
mono_ns=time.monotonic_ns() - 100 * ACCOUNT_STALENESS_NS,
e_live=True,
))
assert provider() is None
class TestCapitalProviderFromSnap:
"""Pure-function variant tests."""
def test_fresh_live(self) -> None:
snap = AccountStateSnapshot(
wallet_balance=10000.0, mono_ns=1000, e_live=True,
)
assert make_e_capital_provider_from_snap(snap, now_ns=2000, staleness_ns=5000) == 10000.0
def test_stale(self) -> None:
snap = AccountStateSnapshot(
wallet_balance=10000.0, mono_ns=1000, e_live=True,
)
result = make_e_capital_provider_from_snap(snap, now_ns=7000, staleness_ns=5000)
assert result is None, f"stale: age={(7000-1000)} >= 5000"
def test_not_live(self) -> None:
snap = AccountStateSnapshot(
wallet_balance=10000.0, mono_ns=1000, e_live=False,
)
assert make_e_capital_provider_from_snap(snap, now_ns=2000) is None
def test_zero_wallet(self) -> None:
snap = AccountStateSnapshot(
wallet_balance=0.0, mono_ns=1000, e_live=True,
)
assert make_e_capital_provider_from_snap(snap, now_ns=2000) is None
def test_negative_wallet(self) -> None:
snap = AccountStateSnapshot(
wallet_balance=-50.0, mono_ns=1000, e_live=True,
)
assert make_e_capital_provider_from_snap(snap, now_ns=2000) is None
def test_exactly_at_boundary_stale(self) -> None:
"""age == staleness → stale (strict <)."""
snap = AccountStateSnapshot(
wallet_balance=10000.0, mono_ns=1000, e_live=True,
)
result = make_e_capital_provider_from_snap(snap, now_ns=6000, staleness_ns=5000)
assert result is None, "age == staleness should be stale (strict <)"
def test_one_ns_before_stale(self) -> None:
"""age = staleness - 1 → fresh."""
snap = AccountStateSnapshot(
wallet_balance=10000.0, mono_ns=1000, e_live=True,
)
result = make_e_capital_provider_from_snap(snap, now_ns=5999, staleness_ns=5000)
assert result == 10000.0, f"age={(5999-1000)} < 5000 should be fresh, got {result}"
def test_none_snap(self) -> None:
assert make_e_capital_provider_from_snap(None, 0) is None
class TestMutationLitmus:
"""Verification that mutations to the age check cause test RED."""
def test_mutation_remove_age_check(self) -> None:
"""If the age check (>= staleness) is removed, stale snapshots return
wallet_balance instead of None.
This test proves the guard exists by showing what happens without it."""
stale_snap = AccountStateSnapshot(
wallet_balance=25000.0,
mono_ns=time.monotonic_ns() - ACCOUNT_STALENESS_NS - 1000,
e_live=True,
)
# With the correct check, this should be None
result = make_e_capital_provider_from_snap(
stale_snap, now_ns=time.monotonic_ns()
)
assert result is None, (
f"MUTATION: stale snapshot returned {result} instead of None. "
"Age check guard is missing or wrong."
)
def test_mutation_remove_e_live_check(self) -> None:
"""If the e_live check is removed, not-live snapshots return wallet_balance."""
not_live_snap = AccountStateSnapshot(
wallet_balance=25000.0, mono_ns=0, e_live=False,
)
result = make_e_capital_provider_from_snap(not_live_snap, now_ns=1000)
assert result is None, (
f"MUTATION: not-live snapshot returned {result}. e_live guard missing."
)
def test_mutation_remove_wallet_positive_check(self) -> None:
"""If the wallet_balance > 0 check is removed, zero wallet returns 0.0."""
zero_snap = AccountStateSnapshot(
wallet_balance=0.0, mono_ns=0, e_live=True,
)
result = make_e_capital_provider_from_snap(zero_snap, now_ns=1000)
assert result is None, (
f"MUTATION: zero-wallet snapshot returned {result}. >0 guard missing."
)

View File

@@ -55,7 +55,10 @@ def _mk_intent(
) -> KernelIntent:
return KernelIntent(
timestamp=datetime.now(timezone.utc),
intent_id=kw.pop("intent_id", trade_id),
# Venue client IDs must distinguish entry and exit orders. Reusing the
# trade ID for both makes one fill match both active orders, which the
# identity-first Rust router correctly refuses to misclassify.
intent_id=kw.pop("intent_id", f"{trade_id}:{action.value.lower()}"),
trade_id=trade_id,
slot_id=slot_id,
asset=asset,

View File

@@ -0,0 +1,113 @@
"""Mutation-sensitive tests for DITAv2 native artifact provenance."""
from __future__ import annotations
import json
import shutil
from pathlib import Path
import pytest
from prod.clean_arch.dita_v2 import native_artifact
def _crate(tmp_path: Path) -> Path:
crate = tmp_path / "crate"
(crate / "src").mkdir(parents=True)
(crate / "Cargo.toml").write_text(
'[package]\nname = "x"\nversion = "0.1.0"\n', encoding="utf-8"
)
(crate / "Cargo.lock").write_text("version = 3\n", encoding="utf-8")
(crate / "src" / "lib.rs").write_text(
"pub fn x() -> u8 { 1 }\n", encoding="utf-8"
)
return crate
def _artifact(tmp_path: Path, crate: Path) -> Path:
library = tmp_path / "release" / native_artifact.library_name()
library.parent.mkdir(parents=True)
library.write_bytes(b"verified-native-artifact")
native_artifact.write_manifest(library, crate)
return library
def test_valid_pair_verifies(tmp_path: Path):
crate = _crate(tmp_path)
library = _artifact(tmp_path, crate)
result = native_artifact.verify_artifact(library, crate)
assert result.library_path == library
assert result.ffi_schema_version == native_artifact.FFI_SCHEMA_VERSION
def test_source_mutation_refuses_startup(tmp_path: Path):
crate = _crate(tmp_path)
library = _artifact(tmp_path, crate)
(crate / "src" / "lib.rs").write_text(
"pub fn x() -> u8 { 2 }\n", encoding="utf-8"
)
with pytest.raises(
native_artifact.ArtifactProvenanceError, match="source_tree_sha256"
):
native_artifact.verify_artifact(library, crate)
def test_artifact_mutation_refuses_startup(tmp_path: Path):
crate = _crate(tmp_path)
library = _artifact(tmp_path, crate)
library.write_bytes(b"tampered-native-artifact")
with pytest.raises(
native_artifact.ArtifactProvenanceError, match="library_sha256"
):
native_artifact.verify_artifact(library, crate)
def test_missing_sidecar_refuses_startup(tmp_path: Path):
crate = _crate(tmp_path)
library = tmp_path / "release" / native_artifact.library_name()
library.parent.mkdir(parents=True)
library.write_bytes(b"no-sidecar")
with pytest.raises(
native_artifact.ArtifactProvenanceError, match="manifest missing"
):
native_artifact.verify_artifact(library, crate)
def test_manifest_self_fingerprint_mutation_refuses_startup(tmp_path: Path):
crate = _crate(tmp_path)
library = _artifact(tmp_path, crate)
sidecar = native_artifact.manifest_path(library)
payload = json.loads(sidecar.read_text(encoding="utf-8"))
payload["ffi_schema_version"] += 1
sidecar.write_text(json.dumps(payload), encoding="utf-8")
with pytest.raises(native_artifact.ArtifactProvenanceError):
native_artifact.verify_artifact(library, crate)
def test_rust_loader_refuses_stale_existing_artifact(tmp_path: Path, monkeypatch):
crate = _crate(tmp_path)
library = _artifact(tmp_path, crate)
library.write_bytes(b"stale")
from prod.clean_arch.dita_v2 import rust_backend
monkeypatch.setattr(rust_backend, "_library_path", lambda: library)
monkeypatch.setattr(rust_backend, "_crate_dir", lambda: crate)
with pytest.raises(native_artifact.ArtifactProvenanceError):
rust_backend._ensure_library()
def test_rollback_restores_previous_verified_pair(tmp_path: Path):
crate = _crate(tmp_path)
target = tmp_path / "target"
library = target / "release" / native_artifact.library_name()
library.parent.mkdir(parents=True)
library.write_bytes(b"known-good")
sidecar = native_artifact.write_manifest(library, crate)
rollback = native_artifact.rollback_directory(target)
rollback.mkdir(parents=True)
shutil.copy2(library, rollback / library.name)
shutil.copy2(sidecar, rollback / sidecar.name)
library.write_bytes(b"bad-current")
result = native_artifact.rollback_artifact(target, crate)
assert library.read_bytes() == b"known-good"
assert result.library_path == library