diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9c011f2..88ff23f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -82,6 +82,8 @@ jobs: run: python -m unittest discover -s open-core/tests -p "test_phase8_identity.py" -v && python -m unittest discover -s open-core/tests -p "test_phase8_guarded_live.py" -v - name: Phase 8B explicit authenticated read-only preflight, fail-closed permissions, redaction, and socket-free default run: python -m unittest discover -s open-core/tests -p "test_phase8b_authenticated_readonly_preflight.py" -v + - name: Phase 8B audited operator bootstrap planning, confirmation, migration identity, and offline attack coverage + run: python -m unittest discover -s open-core/tests -p "test_phase8b_operator_bootstrap.py" -v - name: Compile package and scripts run: python -m compileall -q open-core/src open-core/scripts - run: secure-eval-backtest @@ -99,6 +101,7 @@ jobs: - run: secure-eval-live-status --help - run: secure-eval-live-reconcile --help - run: secure-eval-live-kill --help + - run: secure-eval-live-bootstrap --help postgres-integration: name: PostgreSQL 16 integration @@ -185,6 +188,10 @@ jobs: env: RUN_POSTGRES_INTEGRATION: "true" run: python -m unittest discover -s open-core/tests -p "test_phase8b_postgres_integration.py" -v + - name: Phase 8B dedicated operator bootstrap, exact catalog, configuration replay, rollback, conflict, and concurrency + env: + RUN_POSTGRES_INTEGRATION: "true" + run: python -m unittest discover -s open-core/tests -p "test_phase8b_operator_bootstrap_postgres.py" -v - name: Seeded 0023 to 0026 upgrade with existing Phase 5, Phase 6, Phase 7, and Phase 8A rows run: | python open-core/scripts/apply_postgres_migrations.py --database secure_eval_upgrade --create-database --through 0009 --seed-phase5 diff --git a/.project/implementation_status.json b/.project/implementation_status.json index e754eb9..2bd6619 100644 --- a/.project/implementation_status.json +++ b/.project/implementation_status.json @@ -4,7 +4,7 @@ "repository": "Tcx086/secure-eval-wrapper", "status_source": "docs/IMPLEMENTATION_STATUS.md", "schema": ".project/implementation_status.schema.json", - "updated_at_utc": "2026-07-15T18:55:51Z", + "updated_at_utc": "2026-07-16T21:50:54Z", "current_phase": "phase_8_guarded_live_execution", "rules": { "future_functional_prs_must_update_markdown_status": true, @@ -392,13 +392,17 @@ "Repair the Phase 8B audit findings by removing the unsupported instType=SPOT positions parameter, requiring an empty completed positions response, normalizing documented non-Spot position types as disallowed exposure, removing literal kill-switch success from the standalone proof credential gate, and documenting the normal-role direct-SQL guarantee without claiming resistance to an unrestricted authoritative database writer", "Independently audit and accept the Phase 8B implementation delivered by PR #5, candidate head c3d5dc4b7f4d91d752f2dc2c06aba2078a896180, and merge commit 9b23c71a31a6e4183c23b7504618ebff158fee1e after final-main Actions run 29441539059 passed all six jobs: Ubuntu Python 3.11 job 87441507235, Ubuntu Python 3.12 job 87441507210, Ubuntu Python 3.13 job 87441507245, Windows Python 3.12 job 87441507224, PostgreSQL 16 integration job 87441507309, and public/private and runtime boundary job 87441507243", "Verify migration 0026 SHA-256 698772fb68c5c4981256682d064c3be641193ab10c8dbf55e1a5b390ca7c504a, keep migrations 0001 through 0026 immutable, and accept only the exact ordered six authenticated read-only GETs for account config, balance, one Spot instrument, pending orders, unparameterized positions, and venue time with an empty-positions and position_count=0 success rule", - "Confirm that Phase 8B implementation and audit used no real credentials and performed no authenticated OKX request, production order, or production cancellation; the real local authenticated proof remains unexecuted and production writes remain disabled and unreachable" + "Confirm that Phase 8B implementation and audit used no real credentials and performed no authenticated OKX request, production order, or production cancellation; the real local authenticated proof remains unexecuted and production writes remain disabled and unreachable", + "Implement the dedicated secure-eval-live-bootstrap inspect, plan, initialize, and verify workflow with pinned immutable migration hashes, exact repository, database, and plan identity, explicit read-only confirmation, the fixed conservative BTC-USDT Spot configuration factory, typed atomic PostgreSQL persistence, public-safe verification, offline attack coverage, and isolated PostgreSQL tests; keep this operator bootstrap pending independent audit without accessing credentials, OKX, or the existing operator database", + "Harden bootstrap failure provenance with the exact last completed stage, require the complete Phase 8 schema contract before configuration insertion, derive the bootstrap result hash independently, and add altered-catalog plus concurrent-conflict regressions without changing migrations 0001 through 0026", + "Repair the dedicated bootstrap audit boundary with literal loopback-only PostgreSQL targets, a fixed non-substitutable production database-name rule, server current-user/cluster/version/OID-bound plans, and a fail-closed exact 0001-0026 catalog requirement for every pre-existing target, including collation-only and externally created empty databases; preserve the target-wide full-operation advisory lock, global configuration singleton, independently derived verify hashes, production-named concurrency regressions, immutable migrations 0001-0026, and unreachable production writes" ], "todo": [ - "Optionally execute exactly one controlled local authenticated read-only proof with operator-owned environment credentials that are never persisted, only after the separate exact operator authorization", + "Independently audit and accept the dedicated Phase 8B local operator bootstrap before using it for operator authorization", + "Keep real authenticated proof authorization at NO; optionally execute exactly one controlled local authenticated read-only proof with operator-owned environment credentials that are never persisted only after separate exact operator authorization", "Independently review the resulting redacted proof before accepting the operational Phase 8B checkpoint", - "Do not design Phase 8C until the operational Phase 8B checkpoint is separately accepted", - "Keep production orders and cancellations disabled until a later explicitly approved phase" + "Phase 8C is not started; do not design it until the operational Phase 8B checkpoint is separately accepted", + "Keep production submit and cancel disabled and unreachable until a later explicitly approved phase" ] }, { diff --git a/docs/IMPLEMENTATION_STATUS.md b/docs/IMPLEMENTATION_STATUS.md index 94387c7..022c19b 100644 --- a/docs/IMPLEMENTATION_STATUS.md +++ b/docs/IMPLEMENTATION_STATUS.md @@ -9,7 +9,7 @@ Every future functional PR must update both files in the same change: Completed work must be listed under `Completed`. Everything not done must remain under `Todo`. -Current phase: `phase_8_guarded_live_execution` (`in_progress`). Phase 8A guarded-live dry-run/read-only runtime is accepted; the Phase 8B authenticated read-only proof implementation is independently audited and accepted, while the real local authenticated proof has not yet been executed. Production writes remain disabled and unreachable, Phase 8 is not complete, and Phase 9 remains todo. +Current phase: `phase_8_guarded_live_execution` (`in_progress`). Phase 8A guarded-live dry-run/read-only runtime and the Phase 8B authenticated read-only proof implementation are independently audited and accepted. The real authenticated proof is unexecuted and its authorization remains NO; the operator bootstrap is an implementation candidate pending independent audit. Phase 8 remains `in_progress`, Phase 8C is not started, Phase 9 remains todo, and production submit/cancel remain disabled and unreachable. ## Non-Negotiable Constraints - PostgreSQL is the only authoritative storage layer. @@ -142,8 +142,8 @@ Current phase: `phase_8_guarded_live_execution` (`in_progress`). Phase 8A guarde - [x] Add unambiguous provider instrument identities that separate Spot, perpetual swaps, and dated futures. - [x] Verify and document current official public endpoint, request, response, pagination, limit, and authentication contracts. - [x] Implement Binance and OKX public Spot trade collection through injectable transports. -- [x] Implement Binance USDⓈ-M and OKX SWAP public funding history. -- [x] Implement Binance Spot/USDⓈ-M and OKX SPOT/SWAP public instrument metadata. +- [x] Implement Binance USDⓈ-M and OKX SWAP public funding history. +- [x] Implement Binance Spot/USDⓈ-M and OKX SPOT/SWAP public instrument metadata. - [x] Add deterministic normalization, validation reports, accepted/rejected gates, and quarantine for all three data types. - [x] Add PostgreSQL trade/funding persistence, immutable instrument metadata versions, conflict hashes, reads, indexes, constraints, and foreign-key verification. - [x] Add provider-neutral typed trade, funding, and instrument pipelines with fail-fast, partial, warning, and one-transaction persistence semantics. @@ -339,7 +339,7 @@ Current phase: `phase_8_guarded_live_execution` (`in_progress`). Phase 8A guarde - [x] Verify immutable migrations `0001` through `0025` and confirm that Phase 8A used no real credentials and performed no production orders or production cancellations. - [x] Merge the Phase 8A acceptance status as PR `#4` at `3fe4736832954b01a75906917af41a9bc55745d6`, establishing the sole Phase 8B baseline after main run `29377882425` passed all six jobs: Ubuntu Python 3.11 job `87235016240`, Ubuntu Python 3.12 job `87235016311`, Ubuntu Python 3.13 job `87235016296`, Windows Python 3.12 job `87235016261`, PostgreSQL 16 integration job `87235016225`, and public/private and runtime boundary job `87235016288`. -### Phase 8B: explicit authenticated read-only proof (implementation independently audited and accepted) +### Phase 8B: explicit authenticated read-only proof (implementation independently audited and accepted; real proof unexecuted; operator bootstrap candidate pending independent audit) - [x] Implement an explicit CLI-only authenticated read-only OKX production Spot preflight with required configuration hash, environment credential source, expected account fingerprint, expected reviewed SHA, and instrument; the no-flag path is socket-free and CI is blocked before PostgreSQL, repository identity, credential, or transport access. - [x] Restrict the proof operation to the exact ordered six GETs for account configuration, balances, one Spot instrument, pending orders, unparameterized account positions, and venue time; require the positions data array to be empty, reject `trade` or `withdraw` immediately after the first account-config response, and keep all write/account-power methods unreachable. @@ -350,6 +350,9 @@ Current phase: `phase_8_guarded_live_execution` (`in_progress`). Phase 8A guarde - [x] Independently audit and accept the Phase 8B implementation delivered by PR `#5`, candidate head `c3d5dc4b7f4d91d752f2dc2c06aba2078a896180`, and merge commit `9b23c71a31a6e4183c23b7504618ebff158fee1e`; final-main Actions run `29441539059` passed all six jobs: Ubuntu Python 3.11 job `87441507235`, Ubuntu Python 3.12 job `87441507210`, Ubuntu Python 3.13 job `87441507245`, Windows Python 3.12 job `87441507224`, PostgreSQL 16 integration job `87441507309`, and public/private and runtime boundary job `87441507243`. - [x] Verify migration `0026` SHA-256 `698772fb68c5c4981256682d064c3be641193ab10c8dbf55e1a5b390ca7c504a`, keep migrations `0001` through `0026` immutable, and accept only the exact ordered authenticated read-only GETs: `/api/v5/account/config`, `/api/v5/account/balance`, `/api/v5/public/instruments?instId=&instType=SPOT`, `/api/v5/trade/orders-pending?instId=&instType=SPOT`, `/api/v5/account/positions`, and `/api/v5/public/time`; a successful proof requires an empty positions response and `position_count=0`. - [x] Confirm that implementation and audit used no real credentials and performed no authenticated OKX request, production order, or production cancellation; the real local authenticated proof remains unexecuted and production writes remain disabled and unreachable. +- [x] Implement the dedicated `secure-eval-live-bootstrap` inspect/plan/initialize/verify workflow with pinned immutable migration hashes, exact repository/database/plan identity, an explicit read-only confirmation, the fixed conservative BTC-USDT Spot configuration factory, typed atomic PostgreSQL persistence, public-safe verification, offline attack coverage, and isolated PostgreSQL tests; this operator bootstrap remains pending independent audit and has not accessed credentials, OKX, or the existing operator database. +- [x] Harden bootstrap failure provenance with the exact last completed stage, require the complete Phase 8 schema contract before configuration insertion, derive the bootstrap result hash independently, and add altered-catalog plus concurrent-conflict regressions without changing migrations `0001` through `0026`. +- [x] Repair the dedicated bootstrap audit boundary with literal loopback-only PostgreSQL targets, a fixed non-substitutable production database-name rule, server current-user/cluster/version/OID-bound plans, and a fail-closed exact `0001`–`0026` catalog requirement for every pre-existing target, including collation-only and externally created empty databases; preserve the target-wide full-operation advisory lock, global configuration singleton, independently derived verify hashes, production-named concurrency regressions, immutable migrations `0001`–`0026`, and unreachable production writes. ## Todo @@ -359,10 +362,11 @@ Current phase: `phase_8_guarded_live_execution` (`in_progress`). Phase 8A guarde ### Phase 8: guarded live execution (remaining) -- [ ] Optionally execute exactly one controlled local authenticated read-only proof with operator-owned environment credentials that are never persisted, only after the separate exact operator authorization. +- [ ] Independently audit and accept the dedicated Phase 8B local operator bootstrap before using it for operator authorization. +- [ ] Keep real authenticated proof authorization at NO; optionally execute exactly one controlled local authenticated read-only proof with operator-owned environment credentials that are never persisted only after separate exact operator authorization. - [ ] Independently review the resulting redacted proof before accepting the operational Phase 8B checkpoint. -- [ ] Do not design Phase 8C until the operational Phase 8B checkpoint is separately accepted. -- [ ] Keep production orders and cancellations disabled until a later explicitly approved phase. +- [ ] Phase 8C is not started; do not design it until the operational Phase 8B checkpoint is separately accepted. +- [ ] Keep production submit and cancel disabled and unreachable until a later explicitly approved phase. ### Phase 9: reporting + public delivery - [ ] Build public report templates. diff --git a/open-core/pyproject.toml b/open-core/pyproject.toml index f846cb7..4886222 100644 --- a/open-core/pyproject.toml +++ b/open-core/pyproject.toml @@ -30,6 +30,7 @@ secure-eval-live-dry-run = "secure_eval_wrapper.live.cli:dry_run_main" secure-eval-live-status = "secure_eval_wrapper.live.cli:status_main" secure-eval-live-reconcile = "secure_eval_wrapper.live.cli:reconcile_main" secure-eval-live-kill = "secure_eval_wrapper.live.cli:kill_main" +secure-eval-live-bootstrap = "secure_eval_wrapper.live.bootstrap_cli:main" [tool.setuptools] package-dir = {"" = "src"} diff --git a/open-core/src/secure_eval_wrapper/live/__init__.py b/open-core/src/secure_eval_wrapper/live/__init__.py index 71bb868..36c7c5c 100644 --- a/open-core/src/secure_eval_wrapper/live/__init__.py +++ b/open-core/src/secure_eval_wrapper/live/__init__.py @@ -5,7 +5,6 @@ """ from .approval import LiveApprovalController, confirmation_challenge_hash, manifest_preview_hash from .authorities import * -from .broker import DryRunResult, GuardedLiveBroker from .configuration import GuardedLiveConfiguration, phase8a_dry_run_configuration from .endpoints import EndpointClass, LiveOperation, endpoint_catalog_hash from .gates import evaluate_live_write_authority @@ -14,4 +13,15 @@ from .models import * from .preflight import LivePreflightEngine, LivePreflightEvidence, OperationalPreflightEvidence -__all__ = [name for name in globals() if not name.startswith("_")] +_LAZY_BROKER_EXPORTS = {"DryRunResult", "GuardedLiveBroker"} + + +def __getattr__(name): + if name in _LAZY_BROKER_EXPORTS: + from .broker import DryRunResult, GuardedLiveBroker + + return {"DryRunResult": DryRunResult, "GuardedLiveBroker": GuardedLiveBroker}[name] + raise AttributeError(name) + + +__all__ = [name for name in globals() if not name.startswith("_")] + sorted(_LAZY_BROKER_EXPORTS) diff --git a/open-core/src/secure_eval_wrapper/live/bootstrap.py b/open-core/src/secure_eval_wrapper/live/bootstrap.py new file mode 100644 index 0000000..73dee50 --- /dev/null +++ b/open-core/src/secure_eval_wrapper/live/bootstrap.py @@ -0,0 +1,1167 @@ +"""Audited local PostgreSQL bootstrap for Phase 8B operator authorization setup. + +This module deliberately has no credential-provider or OKX transport dependency and never +invokes provider HTTP or sockets. It creates only a dedicated PostgreSQL database, installs the accepted +immutable migrations, and persists one fixed read-only guarded-live configuration. +""" +from __future__ import annotations + +import hashlib +import importlib.util +import re +from contextlib import contextmanager +from dataclasses import dataclass +from datetime import datetime, timezone +from pathlib import Path +from types import MappingProxyType +from typing import Callable, Mapping + +from secure_eval_wrapper.data_collection.hashing import sha256_payload + +from .configuration import phase8b_authenticated_readonly_configuration +from .durable_repository import DurablePostgresLiveRepository +from .identity import resolve_runtime_repository_identity, validate_git_commit_sha +from .models import live_uuid + + +BOOTSTRAP_VERSION = "phase8b-operator-bootstrap-v1" +DEFAULT_DATABASE = "secure_eval_phase8b" +FORBIDDEN_OPERATOR_DATABASE = "secure_eval_wrapper" +FORBIDDEN_TARGET_DATABASES = frozenset({ + FORBIDDEN_OPERATOR_DATABASE, "postgres", "template0", "template1", +}) +LOCAL_POSTGRES_HOSTS = frozenset({"127.0.0.1", "::1"}) +MAINTENANCE_DATABASES = frozenset({"postgres"}) +DEFAULT_POSTGRES_EXTENSIONS = frozenset({"plpgsql"}) +LATEST_MIGRATION = "0026_phase8b_authenticated_readonly_preflight" +CONFIRMATION_FLAG = "--confirm-readonly-bootstrap" +_IDENTIFIER = re.compile(r"^[a-z][a-z0-9_]{0,62}$") + +EXPECTED_MIGRATION_CATALOG: Mapping[str, str] = MappingProxyType({ + "0001_initial_schema": "598486e6af2eed4559564593adc0b66deff9e21ea91dbda560980c208a2950c5", + "0002_schema_migrations": "36c91efa851e10fcc6039ebd8715af1c985237af6ff556e6943e10329458f76f", + "0003_data_quality_quarantine": "d0b32a72ad98a9d1361bfa57770a9b7d58ae2323816e8b3d77c3d05f66b35a9a", + "0004_reconciliation_persistence": "efe77fa89b25f90dea3f49a70b22b8cc376c434333abbff6fd17cc9eb75fd7ba", + "0005_trade_funding_instrument_hardening": "b18d66f37df55923a1e1cfba709784de55ab90d0c5ff250b8d683dc6029f9d48", + "0006_phase2_final_hardening": "af507329f29e63ab260317b879da5e82917aafd7368d692b343a09ccafdace5d", + "0007_alpha_signal_library": "0a355d3238afcf8691b5366e46332c3e1e6862a9ed574e740e3435479d8883a4", + "0008_phase3_phase4_audit_repairs": "a59dff645009c117a5146d2bd4102a9ed048126ca77b61566f8d31bf1fcba64b", + "0009_phase5_simulated_execution_backtesting": "9b49718ee48e45dda42916568f815723f94578eff814ffe0e0b236aa3523c0d5", + "0010_phase5_second_audit_repairs": "1387ccf65a7a7ac8c2c7b4d93de8443e47963740dcefbf30a0ae248ea5e978a0", + "0011_phase5_run_membership_repairs": "0c0a0ed26ec7419e773e69e8c1ab07d4e220377059e0bf2358b519055e6540a8", + "0012_phase5_run_scoped_projection_repairs": "2a55979b6419bc3eb464d2374d68a40d8cc559fac1a984d2ebcada5974d82d4d", + "0013_phase6_monitoring_simulated_fix": "5e7eb61540507ce4c0f7fb92b78fdebf2fb551a770c129b45d82d19cff592761", + "0014_phase6_first_audit_repairs": "30971466069b6dbcb29f7b08568ebd7791097d996ccbc742cb7f3aa8096ba4fe", + "0015_phase6_concurrency_and_audit_integrity": "5ae8bcfa8db52110978dddd4864700dac6a8e549000dde50065969910b24aec1", + "0016_phase7_safe_paper_trading": "866179dc6a95bf65a416c62d891cd06ce34cf28bceaeb8f29223ad70ef863b0f", + "0017_phase7_durable_paper_recovery": "c2a2e4ca347775898c11443da89552685b0d723a094335719ef07515e4639302", + "0018_phase7_recovery_state_machine_integrity": "c49fad9ed9b5cf3eeee6ae071f6a8b6e4d73c67c80571b09030c4b21b519d59d", + "0019_phase7_venue_event_and_accounting_integrity": "7a139eb65b7ed66fd16b2e7e20794e57f28dc59ab5992c468cb21bae22d68457", + "0020_phase7_price_terminal_and_expiry_integrity": "ce24b36b2ff6e276ce69edeef3044ab7f891e154fa389d55e53619deab990ad5", + "0021_phase7_cancel_terminal_accounting_integrity": "a9a088b497addb45353a3b906caafd5e3532bb389a8ddbfe18626d26597c7506", + "0022_phase8_guarded_live_foundation": "b01c0c0c7801247594ee75009055f899c8902b6cfa1b44ed91ad8451e478e434", + "0023_phase8a_authority_recovery_and_cli_integrity": "cd06abb25ef7a9c178b5aad8c6378f982c879b2e3b52eb4667e067554b987eef", + "0024_phase8a_evidence_reconciliation_metadata_integrity": "3f5671e34d312770dd05763116ce0102da1534061df887dc0c6754f0cc48b214", + "0025_phase8a_okx_credential_permission_authority": "773b2cc2cfb8fcdc9cd9ce022904c096e1f9520915ad8549c5a92d3067d7fc61", + LATEST_MIGRATION: "698772fb68c5c4981256682d064c3be641193ab10c8dbf55e1a5b390ca7c504a", +}) + +PHASE8_REQUIRED_TABLES = frozenset({ + "live_configuration_snapshots", "live_credential_references", "live_account_snapshots", + "live_okx_response_bundles", "live_okx_response_envelopes", "live_preflight_sources", + "live_preflight_reports", "live_preflight_checks", "live_preflight_check_sources", + "live_approvals", "live_run_manifests", "live_runs", "live_kill_switches", + "live_kill_events", "live_run_risk_state", "live_order_intents", + "live_runtime_risk_decisions", "live_reservations", "live_dispatch_outbox", + "live_dispatch_events", "live_cancel_outbox", "live_transport_attempts", + "live_order_observations", "live_order_projections", "live_fill_observations", + "live_reconciliations", "live_reconciliation_differences", "live_recovery_records", + "live_lifecycle_events", "live_pre_run_summaries", "live_post_run_summaries", + "live_market_source_bindings", "live_instrument_metadata_sources", + "live_reconciliation_input_bundles", "live_recovery_query_completions", + "live_authenticated_readonly_proofs", +}) +PHASE8_REQUIRED_INDEXES = frozenset({ + "idx_live_preflight_sources_run_kind", "idx_live_risk_state_day", + "idx_live_reservation_balance", "idx_live_dispatch_claimable", + "idx_live_recovery_claims", "idx_live_okx_bundle_run_purpose", + "idx_live_metadata_run_instrument", "idx_live_recovery_query_matrix", + "idx_live_authenticated_readonly_account_time", +}) +PHASE8_REQUIRED_TRIGGERS = frozenset({ + "trg_guard_live_preflight_authority", "trg_guard_live_manifest_chain", + "trg_live_approval_consumption", "trg_live_intent_mutation", + "trg_live_dispatch_request_immutable", "trg_live_cancel_request_immutable", + "trg_live_dispatch_monotonic", "trg_live_reservation_monotonic", + "trg_live_projection_monotonic", "trg_live_collector_source", + "trg_guard_live_okx_response_payload_hash", "trg_validate_live_okx_bundle_matrix", + "trg_validate_live_okx_envelope_matrix", "trg_validate_live_preflight_graph", + "trg_validate_live_0024_source_details", "trg_validate_live_0025_credential_permission_source", + "trg_validate_live_0025_permission_report", "trg_guard_live_0024_reconciliation", + "trg_validate_live_0024_reconciliation_exact", "trg_guard_live_0024_intent_metadata", + "trg_guard_live_0024_outbox_metadata", "trg_validate_live_0024_recovery_outcome", + "trg_guard_live_0024_kill_reset", "trg_validate_live_authenticated_readonly_proof", + "trg_live_authenticated_readonly_proofs_immutable", +}) + + +class BootstrapSafetyError(PermissionError): + """A public-safe, fail-closed bootstrap refusal.""" + + +class BootstrapOperationError(BootstrapSafetyError): + """A public-safe bootstrap refusal with exact completed-stage provenance.""" + + def __init__(self, message: str, *, last_completed_stage: str) -> None: + super().__init__(message) + self.last_completed_stage = last_completed_stage + + +@dataclass(frozen=True) +class PostgresAdminTarget: + database: str = DEFAULT_DATABASE + host: str = "127.0.0.1" + port: int = 5432 + admin_database: str = "postgres" + admin_user: str = "postgres" + sslmode: str = "disable" + def __post_init__(self) -> None: + for name in ("database", "admin_database"): + value = getattr(self, name) + if not isinstance(value, str) or _IDENTIFIER.fullmatch(value) is None: + raise ValueError(f"{name} must be a conservative PostgreSQL identifier") + if self.database in FORBIDDEN_TARGET_DATABASES: + raise BootstrapSafetyError("the selected database is never a Phase 8B bootstrap target") + if self.database == self.admin_database: + raise BootstrapSafetyError("target database must differ from the admin database") + if self.admin_database not in MAINTENANCE_DATABASES: + raise BootstrapSafetyError("admin database must be a recognized maintenance database") + if re.fullmatch( + r"secure_eval_phase8b(?:_[a-z0-9][a-z0-9_]{0,42})?", self.database + ) is None: + raise BootstrapSafetyError("target database name is not dedicated to Phase 8B") + if self.host not in LOCAL_POSTGRES_HOSTS: + raise BootstrapSafetyError("PostgreSQL host must be literal 127.0.0.1 or ::1") + if isinstance(self.port, bool) or not isinstance(self.port, int) or not 1 <= self.port <= 65535: + raise ValueError("port must be an integer from 1 through 65535") + if not isinstance(self.admin_user, str) or _IDENTIFIER.fullmatch(self.admin_user) is None: + raise ValueError("admin_user must be a conservative PostgreSQL identifier") + if self.sslmode not in {"disable", "require", "verify-ca", "verify-full"}: + raise ValueError("unsupported PostgreSQL sslmode") + + def connection_kwargs(self, database: str, *, read_only: bool) -> dict[str, object]: + result: dict[str, object] = { + "host": self.host, + "port": self.port, + "dbname": database, + "user": self.admin_user, + "sslmode": self.sslmode, + } + if read_only: + result["options"] = "-c default_transaction_read_only=on" + return result + + @property + def public_identity(self) -> str: + return sha256_payload({ + "host": self.host, + "port": self.port, + "database": self.database, + "admin_database": self.admin_database, + "admin_user": self.admin_user, + "sslmode": self.sslmode, + }) + + +@dataclass(frozen=True) +class DatabaseReference: + target_host: str + target_port: int + target_database: str + admin_database: str + admin_user: str + postgres_current_user: str + postgres_system_identifier: str + postgres_server_version: str + database_exists: bool + target_database_oid: int | None + database_identity_sha256: str + + def public_fields(self) -> dict[str, object]: + return { + "target_host": self.target_host, + "target_port": self.target_port, + "target_database": self.target_database, + "admin_database": self.admin_database, + "admin_user": self.admin_user, + "postgres_current_user": self.postgres_current_user, + "postgres_system_identifier": self.postgres_system_identifier, + "postgres_server_version": self.postgres_server_version, + "database_exists": self.database_exists, + "target_database_oid": self.target_database_oid, + "database_identity_sha256": self.database_identity_sha256, + } + + def same_cluster_and_connection(self, other: "DatabaseReference") -> bool: + keys = ( + "target_host", "target_port", "target_database", "admin_database", + "admin_user", "postgres_current_user", "postgres_system_identifier", + "postgres_server_version", + ) + return all(getattr(self, key) == getattr(other, key) for key in keys) + + +@dataclass(frozen=True) +class DatabaseInspection: + reference: DatabaseReference + database_exists: bool + database_oid: int | None + database_identity_sha256: str + catalog_state: str + catalog: tuple[tuple[str, str, str], ...] + latest_migration: str | None + application_row_count: int + configuration_row_count: int + production_write_count: int + non_system_object_count: int + non_system_object_kinds: tuple[str, ...] + blockers: tuple[str, ...] + + +@dataclass(frozen=True) +class ExpectedDatabaseObjects: + schemas: frozenset[str] + tables: frozenset[tuple[str, str]] + functions: frozenset[tuple[str, str]] + triggers: frozenset[tuple[str, str, str, str, str]] + + +def _public_safety_flags() -> dict[str, bool]: + return { + "credentials_accessed": False, + "network_reads_occurred": False, + "network_writes_occurred": False, + "real_proof_executed": False, + } + + +def _migration_root() -> Path: + return Path(__file__).resolve().parents[3] / "db" / "migrations" + + +def verify_local_migration_files(root: Path | None = None) -> tuple[Path, ...]: + migration_root = _migration_root() if root is None else Path(root) + paths = tuple(sorted(migration_root.glob("[0-9][0-9][0-9][0-9]_*.sql"))) + observed = { + path.stem: hashlib.sha256(path.read_bytes().replace(b"\r\n", b"\n")).hexdigest() + for path in paths + } + if observed != dict(EXPECTED_MIGRATION_CATALOG): + raise BootstrapSafetyError("local immutable migration files do not match accepted 0001-0026 hashes") + return paths + + +def _expected_database_objects(paths: tuple[Path, ...]) -> ExpectedDatabaseObjects: + source = "\n".join(path.read_text(encoding="utf-8") for path in paths) + schemas = frozenset(re.findall( + r"(?im)^\s*CREATE\s+SCHEMA\s+(?:IF\s+NOT\s+EXISTS\s+)?([a-z][a-z0-9_]*)", + source, + )) + tables = frozenset( + (match.group(1), match.group(2)) + for match in re.finditer( + r"(?im)^\s*CREATE\s+TABLE\s+(?:IF\s+NOT\s+EXISTS\s+)?" + r"([a-z][a-z0-9_]*)\.([a-z][a-z0-9_]*)", + source, + ) + ) + functions = frozenset( + (match.group(1), match.group(2)) + for match in re.finditer( + r"(?im)^\s*CREATE\s+(?:OR\s+REPLACE\s+)?(?:FUNCTION|PROCEDURE)\s+" + r"([a-z][a-z0-9_]*)\.([a-z][a-z0-9_]*)\s*\(", + source, + ) + ) + triggers: dict[ + tuple[str, str, str], tuple[str, str, str, str, str] + ] = {} + for match in re.finditer( + r"(?ims)^\s*CREATE\s+(?:CONSTRAINT\s+)?TRIGGER\s+([a-z][a-z0-9_]*)\b(.*?);", + source, + ): + body = match.group(2) + relation = re.search( + r"\bON\s+([a-z][a-z0-9_]*)\.([a-z][a-z0-9_]*)\b", body, re.I + ) + function = re.search( + r"\bEXECUTE\s+FUNCTION\s+([a-z][a-z0-9_]*)\.([a-z][a-z0-9_]*)\s*\(", + body, + re.I, + ) + if relation is None or function is None: + raise BootstrapSafetyError("accepted migration trigger catalog could not be derived") + signature = ( + relation.group(1).lower(), relation.group(2).lower(), match.group(1).lower(), + function.group(1).lower(), function.group(2).lower(), + ) + triggers[signature[:3]] = signature + for table in ( + "paper_internal_venue_commands", "paper_internal_venue_events", + "paper_venue_order_observations", "paper_recovery_observation_bundles", + "paper_fill_recovery_lineage", + ): + signature = ( + "execution", table, f"phase7_{table}_append_only", + "execution", "phase7_reject_immutable_change", + ) + triggers[signature[:3]] = signature + for table in ( + "live_configuration_snapshots", "live_credential_references", + "live_account_snapshots", "live_preflight_sources", "live_preflight_checks", + "live_preflight_check_sources", "live_preflight_reports", "live_run_manifests", + "live_runtime_risk_decisions", "live_transport_attempts", + "live_order_observations", "live_fill_observations", "live_reconciliations", + "live_reconciliation_differences", "live_pre_run_summaries", + "live_post_run_summaries", "live_lifecycle_events", "live_kill_events", + "live_dispatch_events", "live_okx_response_bundles", + "live_okx_response_envelopes", "live_market_source_bindings", + "live_instrument_metadata_sources", "live_reconciliation_input_bundles", + "live_recovery_query_completions", + ): + signature = ( + "execution", table, f"trg_{table}_immutable", + "execution", "prevent_live_authority_mutation", + ) + triggers[signature[:3]] = signature + return ExpectedDatabaseObjects( + schemas, tables, functions, frozenset(triggers.values()) + ) + +def derive_bootstrap_record_hash(payload: Mapping[str, object]) -> str: + core = dict(payload) + core.pop("bootstrap_record_hash", None) + return sha256_payload(core) + + +def _default_connector(**kwargs): + try: + import psycopg + except ImportError as exc: + raise RuntimeError("operator bootstrap requires the PostgreSQL package extra") from exc + return psycopg.connect(**kwargs) + + +def _default_identity(): + return resolve_runtime_repository_identity(environment={}) + + +def _quoted(identifier: str) -> str: + return '"' + identifier.replace('"', '""') + '"' + + +class Phase8BOperatorBootstrap: + """Plan, initialize, and verify one dedicated local PostgreSQL target.""" + + def __init__( + self, + target: PostgresAdminTarget, + *, + connector: Callable[..., object] = _default_connector, + identity_resolver: Callable[[], object] = _default_identity, + clock: Callable[[], datetime] | None = None, + ) -> None: + self.target = target + self._connector = connector + self._identity_resolver = identity_resolver + self._clock = clock or (lambda: datetime.now(timezone.utc)) + + def _connect(self, database: str, *, read_only: bool): + return self._connector(**self.target.connection_kwargs(database, read_only=read_only)) + + @staticmethod + def _fetchone(connection, statement: str, params=()): + with connection.cursor() as cursor: + cursor.execute(statement, params) + return cursor.fetchone() + + @staticmethod + def _fetchall(connection, statement: str, params=()): + with connection.cursor() as cursor: + cursor.execute(statement, params) + return cursor.fetchall() + + def _repository_sha(self, expected_reviewed_sha: str | None = None) -> str: + identity = self._identity_resolver() + observed = validate_git_commit_sha(identity.observed_commit_sha, field_name="observed_repository_sha") + if expected_reviewed_sha is not None: + expected = validate_git_commit_sha(expected_reviewed_sha, field_name="expected_reviewed_sha") + if observed != expected: + raise BootstrapSafetyError("observed repository SHA does not match the exact expected SHA") + return observed + + def _database_reference(self, connection=None) -> DatabaseReference: + owned = connection is None + if owned: + connection = self._connect(self.target.admin_database, read_only=True) + try: + row = self._fetchone( + connection, + "SELECT current_user::text,current_setting('server_version')::text," + "control.system_identifier::text,d.oid::bigint " + "FROM pg_control_system() AS control " + "LEFT JOIN pg_database d ON d.datname=%s", + (self.target.database,), + ) + except Exception as exc: + raise BootstrapSafetyError( + "PostgreSQL cluster identity could not be established" + ) from exc + finally: + if owned: + connection.close() + if row is None or len(row) != 4 or any(value is None for value in row[:3]): + raise BootstrapSafetyError("PostgreSQL cluster identity could not be established") + current_user, server_version, system_identifier, oid_value = row + if not str(current_user) or not str(server_version) or not str(system_identifier).isdigit(): + raise BootstrapSafetyError("PostgreSQL cluster identity is invalid") + core = { + "target_host": self.target.host, + "target_port": self.target.port, + "target_database": self.target.database, + "admin_database": self.target.admin_database, + "admin_user": self.target.admin_user, + "postgres_current_user": str(current_user), + "postgres_system_identifier": str(system_identifier), + "postgres_server_version": str(server_version), + "database_exists": oid_value is not None, + "target_database_oid": None if oid_value is None else int(oid_value), + } + return DatabaseReference( + **core, + database_identity_sha256=sha256_payload(core), + ) + + def _verify_target_connection_identity( + self, connection, reference: DatabaseReference + ) -> None: + try: + row = self._fetchone( + connection, + "SELECT current_database()::text,d.oid::bigint,current_user::text," + "current_setting('server_version')::text,control.system_identifier::text " + "FROM pg_database d CROSS JOIN pg_control_system() AS control " + "WHERE d.datname=current_database()", + ) + except Exception as exc: + raise BootstrapSafetyError("target database identity could not be verified") from exc + expected = ( + reference.target_database, + reference.target_database_oid, + reference.postgres_current_user, + reference.postgres_server_version, + reference.postgres_system_identifier, + ) + observed = None if row is None else ( + str(row[0]), int(row[1]), str(row[2]), str(row[3]), str(row[4]) + ) + if observed != expected: + raise BootstrapSafetyError("target database or cluster identity changed") + + @contextmanager + def _locked_admin_connection(self): + connection = self._connect(self.target.admin_database, read_only=False) + try: + connection.autocommit = True + lock_name = "phase8b-operator-bootstrap:" + self.target.database + self._fetchone( + connection, + "SELECT pg_advisory_lock(hashtextextended(%s,0))", + (lock_name,), + ) + yield connection + finally: + connection.close() + + def _inspect_database_objects( + self, connection, *, exact_catalog: bool + ) -> tuple[int, tuple[str, ...], tuple[str, ...], list[tuple[str, str]]]: + system_filter = ( + "n.nspname NOT IN ('pg_catalog','information_schema') " + "AND n.nspname NOT LIKE 'pg_toast%%' AND n.nspname NOT LIKE 'pg_temp_%%'" + ) + schemas = {str(row[0]) for row in self._fetchall( + connection, + "SELECT nspname FROM pg_namespace n WHERE " + system_filter + " ORDER BY nspname", + )} + relations = [ + (str(schema), str(name), str(kind)) + for schema, name, kind in self._fetchall( + connection, + "SELECT n.nspname,c.relname,c.relkind::text FROM pg_class c " + "JOIN pg_namespace n ON n.oid=c.relnamespace WHERE " + system_filter + + " AND c.relpersistence<>'t' AND c.relkind IN ('r','p','v','m','S','f','c') " + "ORDER BY n.nspname,c.relname,c.relkind", + ) + ] + functions = [ + (str(schema), str(name), int(count)) + for schema, name, count in self._fetchall( + connection, + "SELECT n.nspname,p.proname,count(*)::bigint FROM pg_proc p " + "JOIN pg_namespace n ON n.oid=p.pronamespace WHERE " + system_filter + + " GROUP BY n.nspname,p.proname ORDER BY n.nspname,p.proname", + ) + ] + triggers = frozenset( + (str(a), str(b), str(c), str(d), str(e)) + for a, b, c, d, e in self._fetchall( + connection, + "SELECT n.nspname,rel.relname,t.tgname,pn.nspname,p.proname " + "FROM pg_trigger t JOIN pg_class rel ON rel.oid=t.tgrelid " + "JOIN pg_namespace n ON n.oid=rel.relnamespace " + "JOIN pg_proc p ON p.oid=t.tgfoid JOIN pg_namespace pn ON pn.oid=p.pronamespace " + "WHERE NOT t.tgisinternal AND " + system_filter + + " ORDER BY n.nspname,rel.relname,t.tgname", + ) + ) + extensions = { + str(row[0]) for row in self._fetchall( + connection, "SELECT extname FROM pg_extension ORDER BY extname" + ) + } + type_count = int(self._fetchone( + connection, + "SELECT count(*)::bigint FROM pg_type t JOIN pg_namespace n ON n.oid=t.typnamespace " + "WHERE " + system_filter + " AND t.typtype IN ('d','e')", + )[0]) + other_counts = { + str(kind): int(count) + for kind, count in self._fetchall( + connection, + "SELECT 'event_trigger',count(*)::bigint FROM pg_event_trigger UNION ALL " + "SELECT 'publication',count(*)::bigint FROM pg_publication UNION ALL " + "SELECT 'subscription',count(*)::bigint FROM pg_subscription UNION ALL " + "SELECT 'foreign_server',count(*)::bigint FROM pg_foreign_server UNION ALL " + "SELECT 'foreign_data_wrapper',count(*)::bigint FROM pg_foreign_data_wrapper UNION ALL " + "SELECT 'large_object',count(*)::bigint FROM pg_largeobject_metadata UNION ALL " + "SELECT 'database_setting',count(*)::bigint FROM pg_db_role_setting s " + "WHERE s.setdatabase=(SELECT oid FROM pg_database WHERE datname=current_database()) " + "OR (s.setdatabase=0 AND s.setrole=(SELECT oid FROM pg_roles WHERE rolname=current_user))", + ) + } + user_tables = [(schema, name) for schema, name, kind in relations if kind in {"r", "p"}] + counts = { + "schema": len(schemas - {"public"}), + "relation": len(relations), + "function_or_procedure": sum(count for _, _, count in functions), + "trigger": len(triggers), + "extension": len(extensions - DEFAULT_POSTGRES_EXTENSIONS), + "type": type_count, + **other_counts, + } + blockers: list[str] = [] + if exact_catalog: + expected = _expected_database_objects(verify_local_migration_files()) + if schemas != set(expected.schemas) | {"public"}: + blockers.append("database_schema_catalog_mismatch") + if set(user_tables) != set(expected.tables) or any( + kind not in {"r", "p"} for _, _, kind in relations + ): + blockers.append("persistent_relation_catalog_mismatch") + if ( + {(schema, name) for schema, name, _ in functions} != set(expected.functions) + or any(count != 1 for _, _, count in functions) + ): + blockers.append("stored_function_catalog_mismatch") + if triggers != expected.triggers: + blockers.append("non_internal_trigger_catalog_mismatch") + if extensions - DEFAULT_POSTGRES_EXTENSIONS: + blockers.append("unexpected_extension") + if type_count or any(other_counts.values()): + blockers.append("unexpected_persistent_database_object") + object_kinds = tuple(sorted(kind for kind, count in counts.items() if count)) + return sum(counts.values()), object_kinds, blockers, user_tables + + def _inspect_existing_database( + self, + reference: DatabaseReference, + expected_configuration=None, + ) -> DatabaseInspection: + blockers: list[str] = [] + connection = self._connect(self.target.database, read_only=True) + try: + self._verify_target_connection_identity(connection, reference) + catalog_ready = bool(self._fetchone( + connection, + "SELECT to_regclass('audit.schema_migrations') IS NOT NULL", + )[0]) + catalog: tuple[tuple[str, str, str], ...] = () + if catalog_ready: + catalog = tuple((str(a), str(b), str(c)) for a, b, c in self._fetchall( + connection, + "SELECT migration_id,filename,sha256::text FROM audit.schema_migrations " + "ORDER BY migration_id", + )) + + expected_rows = tuple( + (migration_id, migration_id + ".sql", digest) + for migration_id, digest in EXPECTED_MIGRATION_CATALOG.items() + ) + exact_catalog = catalog_ready and catalog == expected_rows + ( + non_system_object_count, + non_system_object_kinds, + object_blockers, + user_tables, + ) = self._inspect_database_objects(connection, exact_catalog=exact_catalog) + if not catalog_ready: + catalog_state = "uncatalogued" + blockers.append("existing_uncatalogued_database_is_never_initialized") + elif exact_catalog: + catalog_state = "exact_0001_0026" + blockers.extend(object_blockers) + else: + observed_ids = {row[0] for row in catalog} + expected_ids = set(EXPECTED_MIGRATION_CATALOG) + if observed_ids - expected_ids: + catalog_state = "unknown_migrations" + blockers.append("migration_catalog_contains_unknown_entries") + elif observed_ids != expected_ids: + catalog_state = "partial_catalog" + blockers.append("partial_migration_catalog_is_never_auto_upgraded") + else: + catalog_state = "hash_or_filename_mismatch" + blockers.append("immutable_migration_catalog_mismatch") + + application_rows = 0 + configuration_rows: list[tuple] = [] + if catalog_state == "exact_0001_0026": + for schema, table in user_tables: + if (schema, table) == ("audit", "schema_migrations"): + continue + count = int(self._fetchone( + connection, + f"SELECT count(*)::bigint FROM {_quoted(schema)}.{_quoted(table)}", + )[0]) + if (schema, table) == ("execution", "live_configuration_snapshots"): + if count: + configuration_rows = list(self._fetchall( + connection, + "SELECT configuration_snapshot_id,configuration_sha256,record_sha256," + "account_fingerprint,dry_run,read_only_preflight,production_write_enabled " + "FROM execution.live_configuration_snapshots ORDER BY configuration_sha256", + )) + else: + application_rows += count + if application_rows: + blockers.append("existing_database_contains_unsafe_application_rows") + if configuration_rows: + if expected_configuration is None: + blockers.append("existing_configuration_requires_exact_plan_inputs") + else: + expected_id = live_uuid( + "configuration", {"hash": expected_configuration.configuration_hash} + ) + exact = len(configuration_rows) == 1 and ( + configuration_rows[0][0] == expected_id + and configuration_rows[0][1] == expected_configuration.configuration_hash + and configuration_rows[0][2] == expected_configuration.configuration_hash + and configuration_rows[0][3] == expected_configuration.account_fingerprint + and bool(configuration_rows[0][4]) + and bool(configuration_rows[0][5]) + and not bool(configuration_rows[0][6]) + ) + if not exact: + blockers.append("existing_guarded_live_configuration_conflicts") + + production_write_count = 0 + if catalog_state == "exact_0001_0026": + production_write_count = int(self._fetchone( + connection, + "SELECT count(*)::bigint FROM execution.live_transport_attempts " + "WHERE external_write_attempted OR successful_write", + )[0]) + if production_write_count: + blockers.append("production_write_history_is_not_zero") + finally: + connection.close() + + return DatabaseInspection( + reference=reference, + database_exists=True, + database_oid=reference.target_database_oid, + database_identity_sha256=reference.database_identity_sha256, + catalog_state=catalog_state, + catalog=catalog, + latest_migration=None if not catalog else catalog[-1][0], + application_row_count=application_rows, + configuration_row_count=len(configuration_rows), + production_write_count=production_write_count, + non_system_object_count=non_system_object_count, + non_system_object_kinds=non_system_object_kinds, + blockers=tuple(sorted(set(blockers))), + ) + + def inspect( + self, + *, + expected_configuration=None, + reference: DatabaseReference | None = None, + admin_connection=None, + ) -> DatabaseInspection: + verify_local_migration_files() + reference = reference or self._database_reference(admin_connection) + if not reference.database_exists: + return DatabaseInspection( + reference=reference, + database_exists=False, + database_oid=None, + database_identity_sha256=reference.database_identity_sha256, + catalog_state="absent", + catalog=(), + latest_migration=None, + application_row_count=0, + configuration_row_count=0, + production_write_count=0, + non_system_object_count=0, + non_system_object_kinds=(), + blockers=(), + ) + return self._inspect_existing_database(reference, expected_configuration) + def inspect_public(self) -> dict[str, object]: + observed_sha = self._repository_sha() + state = self.inspect() + return { + "command": "secure-eval-live-bootstrap inspect", + "version": BOOTSTRAP_VERSION, + "action": "inspect", + **state.reference.public_fields(), + "catalog_state": state.catalog_state, + "current_latest_migration": state.latest_migration, + "expected_latest_migration": LATEST_MIGRATION, + "expected_0026_sha256": EXPECTED_MIGRATION_CATALOG[LATEST_MIGRATION], + "immutable_catalog_verified": state.catalog_state == "exact_0001_0026", + "observed_repository_sha": observed_sha, + "application_row_count": state.application_row_count, + "configuration_row_count": state.configuration_row_count, + "production_write_count": state.production_write_count, + "non_system_object_count": state.non_system_object_count, + "non_system_object_kinds": list(state.non_system_object_kinds), + "blockers": list(state.blockers), + **_public_safety_flags(), + } + + def plan( + self, + *, + expected_reviewed_sha: str, + account_fingerprint: str, + instrument: str, + admin_connection=None, + reference: DatabaseReference | None = None, + ) -> dict[str, object]: + observed_sha = self._repository_sha(expected_reviewed_sha) + reference = reference or self._database_reference(admin_connection) + configuration = phase8b_authenticated_readonly_configuration( + account_fingerprint, instrument + ) + state = self.inspect( + expected_configuration=configuration, + reference=reference, + admin_connection=admin_connection, + ) + migrations_required = not state.database_exists + core = { + "command": "secure-eval-live-bootstrap plan", + "version": BOOTSTRAP_VERSION, + "action": "plan", + **reference.public_fields(), + "catalog_state": state.catalog_state, + "current_migration_count": len(state.catalog), + "current_latest_migration": state.latest_migration, + "expected_latest_migration": LATEST_MIGRATION, + "expected_0026_sha256": EXPECTED_MIGRATION_CATALOG[LATEST_MIGRATION], + "immutable_catalog_verified": state.catalog_state == "exact_0001_0026", + "observed_repository_sha": observed_sha, + "expected_reviewed_sha": expected_reviewed_sha, + "account_fingerprint": configuration.account_fingerprint, + "instrument": instrument, + "configuration_hash": configuration.configuration_hash, + "current_endpoint_catalog_hash": configuration.endpoint_catalog_hash, + "current_adapter_implementation_hash": configuration.provider_implementation_hash, + "intended_credential_policy": list(configuration.credential_source_policy), + "production_write_enabled": configuration.production_write_enabled, + "database_creation_required": not state.database_exists, + "migrations_required": migrations_required, + "migration_count_to_apply": len(EXPECTED_MIGRATION_CATALOG) if migrations_required else 0, + "configuration_insertion_required": state.configuration_row_count == 0, + "configuration_replay": state.configuration_row_count == 1, + "application_row_count": state.application_row_count, + "production_write_count": state.production_write_count, + "non_system_object_count": state.non_system_object_count, + "non_system_object_kinds": list(state.non_system_object_kinds), + "blockers": list(state.blockers), + **_public_safety_flags(), + } + return {**core, "plan_hash": sha256_payload(core)} + @staticmethod + def _load_migration_runner(): + path = Path(__file__).resolve().parents[3] / "scripts" / "apply_postgres_migrations.py" + spec = importlib.util.spec_from_file_location("secure_eval_phase8b_migration_runner", path) + if spec is None or spec.loader is None: + raise RuntimeError("accepted PostgreSQL migration runner is unavailable") + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + def _create_database(self, admin_connection) -> None: + if self._fetchone( + admin_connection, + "SELECT 1 FROM pg_database WHERE datname=%s", + (self.target.database,), + ) is not None: + raise BootstrapSafetyError("database identity changed after the confirmed plan") + try: + from psycopg import sql + except ImportError as exc: + raise RuntimeError("operator bootstrap requires the PostgreSQL package extra") from exc + with admin_connection.cursor() as cursor: + cursor.execute(sql.SQL("CREATE DATABASE {}").format(sql.Identifier(self.target.database))) + + def _apply_all_migrations(self, reference: DatabaseReference) -> None: + paths = verify_local_migration_files() + runner = self._load_migration_runner() + connection = self._connect(self.target.database, read_only=False) + try: + self._verify_target_connection_identity(connection, reference) + runner.bootstrap(connection) + connection.commit() + for path in paths: + digest, _ = runner._apply_migration(connection, path) + if digest != EXPECTED_MIGRATION_CATALOG[path.stem]: + raise BootstrapSafetyError("migration runner digest disagrees with accepted catalog") + finally: + connection.close() + + def _schema_contract(self, connection) -> dict[str, object]: + tables = {row[0] for row in self._fetchall( + connection, + "SELECT table_name FROM information_schema.tables " + "WHERE table_schema='execution' AND table_type='BASE TABLE'", + )} + indexes = {row[0] for row in self._fetchall( + connection, + "SELECT indexname FROM pg_indexes WHERE schemaname='execution'", + )} + triggers = {row[0] for row in self._fetchall( + connection, + "SELECT tgname FROM pg_trigger t JOIN pg_class c ON c.oid=t.tgrelid " + "JOIN pg_namespace n ON n.oid=c.relnamespace " + "WHERE n.nspname='execution' AND NOT t.tgisinternal", + )} + return { + "phase8_tables_verified": PHASE8_REQUIRED_TABLES <= tables, + "phase8_indexes_verified": PHASE8_REQUIRED_INDEXES <= indexes, + "phase8_triggers_verified": PHASE8_REQUIRED_TRIGGERS <= triggers, + } + + @staticmethod + def _require_same_reference( + expected: DatabaseReference, + observed: DatabaseReference, + *, + allow_creation: bool = False, + ) -> None: + if not expected.same_cluster_and_connection(observed): + raise BootstrapSafetyError("PostgreSQL connection or cluster identity changed") + if allow_creation and not expected.database_exists: + if not observed.database_exists or observed.target_database_oid is None: + raise BootstrapSafetyError("dedicated target database creation was not observable") + return + if observed.public_fields() != expected.public_fields(): + raise BootstrapSafetyError("target database identity changed") + + def verify( + self, + *, + expected_reviewed_sha: str, + account_fingerprint: str, + instrument: str, + expected_reference: DatabaseReference | None = None, + admin_connection=None, + ) -> dict[str, object]: + observed_sha = self._repository_sha(expected_reviewed_sha) + current_reference = self._database_reference(admin_connection) + if expected_reference is not None: + self._require_same_reference(expected_reference, current_reference) + configuration = phase8b_authenticated_readonly_configuration( + account_fingerprint, instrument + ) + state = self.inspect( + expected_configuration=configuration, + reference=current_reference, + admin_connection=admin_connection, + ) + contract = { + "phase8_tables_verified": False, + "phase8_indexes_verified": False, + "phase8_triggers_verified": False, + } + typed_reload_verified = False + snapshot_id = live_uuid("configuration", {"hash": configuration.configuration_hash}) + if ( + not state.blockers + and state.catalog_state == "exact_0001_0026" + and state.configuration_row_count == 1 + ): + connection = self._connect(self.target.database, read_only=True) + try: + self._verify_target_connection_identity(connection, current_reference) + contract = self._schema_contract(connection) + repository = DurablePostgresLiveRepository(connection) + typed_reload_verified = ( + repository.load_guarded_live_configuration(configuration.configuration_hash) + == configuration + ) + finally: + connection.close() + ready = all(contract.values()) and typed_reload_verified and ( + not state.blockers + and state.latest_migration == LATEST_MIGRATION + and state.production_write_count == 0 + and state.configuration_row_count == 1 + ) + blockers = list(state.blockers) + if state.catalog_state != "exact_0001_0026": + blockers.append("exact_migration_catalog_not_ready") + if not all(contract.values()): + blockers.append("phase8_schema_contract_not_ready") + if not typed_reload_verified: + blockers.append("typed_configuration_reload_not_ready") + if state.production_write_count: + blockers.append("production_write_history_is_not_zero") + core = { + "command": "secure-eval-live-bootstrap verify", + "version": BOOTSTRAP_VERSION, + "action": "verify", + **current_reference.public_fields(), + "ready_for_operator_authorization": ready, + "catalog_state": state.catalog_state, + "migration_count": len(state.catalog), + "migration_catalog": [ + {"migration_id": item[0], "filename": item[1], "sha256": item[2]} + for item in state.catalog + ], + "latest_migration": state.latest_migration, + "immutable_catalog_verified": state.catalog_state == "exact_0001_0026", + "migration_0026_installed": state.latest_migration == LATEST_MIGRATION, + "migration_hashes_verified": state.catalog == tuple( + (key, key + ".sql", value) for key, value in EXPECTED_MIGRATION_CATALOG.items() + ), + **contract, + "configuration_snapshot_id": str(snapshot_id), + "configuration_hash": configuration.configuration_hash, + "typed_configuration_reload_verified": typed_reload_verified, + "current_endpoint_catalog_hash": configuration.endpoint_catalog_hash, + "current_provider_implementation_hash": configuration.provider_implementation_hash, + "observed_repository_sha": observed_sha, + "expected_reviewed_sha": expected_reviewed_sha, + "account_fingerprint": configuration.account_fingerprint, + "instrument": instrument, + "allowed_instruments": list(configuration.allowed_instruments), + "allowed_instrument_types": list(configuration.allowed_instrument_types), + "allowed_settlement_assets": list(configuration.allowed_settlement_assets), + "base_currency": configuration.base_currency, + "allowed_order_types": list(configuration.allowed_order_types), + "credential_source_policy": list(configuration.credential_source_policy), + "dry_run": configuration.dry_run, + "read_only_preflight": configuration.read_only_preflight, + "production_write_enabled": configuration.production_write_enabled, + "automatic_flatten": configuration.automatic_flatten, + "allow_short": configuration.allow_short, + "allow_perpetual": configuration.allow_perpetual, + "production_write_count": state.production_write_count, + "blockers": sorted(set(blockers)), + **_public_safety_flags(), + } + return {**core, "bootstrap_record_hash": derive_bootstrap_record_hash(core)} + def initialize( + self, + *, + expected_reviewed_sha: str, + account_fingerprint: str, + instrument: str, + previous_plan_hash: str, + confirm_readonly_bootstrap: bool, + ) -> dict[str, object]: + progress = ["not_started"] + try: + return self._initialize_confirmed( + expected_reviewed_sha=expected_reviewed_sha, + account_fingerprint=account_fingerprint, + instrument=instrument, + previous_plan_hash=previous_plan_hash, + confirm_readonly_bootstrap=confirm_readonly_bootstrap, + progress=progress, + ) + except BootstrapOperationError: + raise + except (BootstrapSafetyError, ValueError) as exc: + raise BootstrapOperationError( + str(exc), last_completed_stage=progress[0] + ) from exc + except Exception as exc: + raise BootstrapOperationError( + "local_postgresql_operation_failed", + last_completed_stage=progress[0], + ) from exc + + def _initialize_confirmed( + self, + *, + expected_reviewed_sha: str, + account_fingerprint: str, + instrument: str, + previous_plan_hash: str, + confirm_readonly_bootstrap: bool, + progress: list[str], + ) -> dict[str, object]: + if not confirm_readonly_bootstrap: + raise BootstrapSafetyError( + f"initialization requires exact {CONFIRMATION_FLAG} confirmation" + ) + if not isinstance(previous_plan_hash, str) or not re.fullmatch( + r"[0-9a-f]{64}", previous_plan_hash + ): + raise BootstrapSafetyError("previous plan hash must be an exact lowercase SHA-256") + progress[0] = "confirmation_validated" + + with self._locked_admin_connection() as admin_connection: + progress[0] = "target_lock_acquired" + current = self.plan( + expected_reviewed_sha=expected_reviewed_sha, + account_fingerprint=account_fingerprint, + instrument=instrument, + admin_connection=admin_connection, + ) + if current["blockers"]: + raise BootstrapSafetyError("confirmed plan contains blockers") + if current["plan_hash"] != previous_plan_hash: + raise BootstrapSafetyError( + "provided plan hash does not match the locked current plan" + ) + planned_reference = self._database_reference(admin_connection) + planned_fields = planned_reference.public_fields() + if any(current.get(key) != value for key, value in planned_fields.items()): + raise BootstrapSafetyError("database identity changed during plan revalidation") + progress[0] = "plan_revalidated" + + if current["database_creation_required"]: + self._create_database(admin_connection) + active_reference = self._database_reference(admin_connection) + self._require_same_reference( + planned_reference, + active_reference, + allow_creation=current["database_creation_required"], + ) + progress[0] = "database_ready" + + before_migration = self._database_reference(admin_connection) + self._require_same_reference(active_reference, before_migration) + if current["migrations_required"]: + self._apply_all_migrations(active_reference) + progress[0] = "migrations_ready" + + configuration = phase8b_authenticated_readonly_configuration( + account_fingerprint, instrument + ) + before_schema = self._database_reference(admin_connection) + self._require_same_reference(active_reference, before_schema) + pre_persist = self.inspect( + expected_configuration=configuration, + reference=before_schema, + admin_connection=admin_connection, + ) + if pre_persist.catalog_state != "exact_0001_0026" or pre_persist.blockers: + raise BootstrapSafetyError( + "post-migration catalog is not safe for configuration persistence" + ) + schema_connection = self._connect(self.target.database, read_only=True) + try: + self._verify_target_connection_identity(schema_connection, active_reference) + schema_contract = self._schema_contract(schema_connection) + finally: + schema_connection.close() + if not all(schema_contract.values()): + raise BootstrapSafetyError( + "Phase 8 schema contract failed before configuration persistence" + ) + progress[0] = "schema_verified" + + before_persist = self._database_reference(admin_connection) + self._require_same_reference(active_reference, before_persist) + connection = self._connect(self.target.database, read_only=False) + try: + with connection.transaction(): + self._verify_target_connection_identity(connection, active_reference) + repository = DurablePostgresLiveRepository(connection) + snapshot_id = repository.persist_guarded_live_configuration_snapshot( + configuration=configuration, + created_at_utc=self._clock(), + ) + progress[0] = "configuration_persisted" + finally: + connection.close() + + before_verify = self._database_reference(admin_connection) + self._require_same_reference(active_reference, before_verify) + result = self.verify( + expected_reviewed_sha=expected_reviewed_sha, + account_fingerprint=account_fingerprint, + instrument=instrument, + expected_reference=active_reference, + admin_connection=admin_connection, + ) + if not result["ready_for_operator_authorization"]: + raise BootstrapSafetyError( + "configuration persisted but final verification failed closed" + ) + progress[0] = "verification_completed" + result_core = { + "command": "secure-eval-live-bootstrap initialize", + "version": BOOTSTRAP_VERSION, + **active_reference.public_fields(), + "observed_repository_sha": result["observed_repository_sha"], + "expected_reviewed_sha": result["expected_reviewed_sha"], + "migration_count": result["migration_count"], + "latest_migration": result["latest_migration"], + "immutable_catalog_verified": result["immutable_catalog_verified"], + "migration_0026_installed": result["migration_0026_installed"], + "configuration_snapshot_id": str(snapshot_id), + "configuration_hash": result["configuration_hash"], + "account_fingerprint": result["account_fingerprint"], + "instrument": instrument, + "credential_policy": result["credential_source_policy"], + "endpoint_catalog_hash": result["current_endpoint_catalog_hash"], + "adapter_implementation_hash": result["current_provider_implementation_hash"], + "dry_run": result["dry_run"], + "read_only_preflight": result["read_only_preflight"], + "production_write_enabled": result["production_write_enabled"], + "credentials_accessed": result["credentials_accessed"], + "network_reads_occurred": result["network_reads_occurred"], + "network_writes_occurred": result["network_writes_occurred"], + "real_proof_executed": result["real_proof_executed"], + } + return { + **result_core, + "bootstrap_record_hash": derive_bootstrap_record_hash(result_core), + } + +__all__ = [ + "BOOTSTRAP_VERSION", "BootstrapOperationError", "BootstrapSafetyError", + "CONFIRMATION_FLAG", "DEFAULT_DATABASE", + "EXPECTED_MIGRATION_CATALOG", "LATEST_MIGRATION", "DatabaseInspection", "DatabaseReference", + "Phase8BOperatorBootstrap", "PostgresAdminTarget", "derive_bootstrap_record_hash", + "verify_local_migration_files", +] diff --git a/open-core/src/secure_eval_wrapper/live/bootstrap_cli.py b/open-core/src/secure_eval_wrapper/live/bootstrap_cli.py new file mode 100644 index 0000000..8bae7e1 --- /dev/null +++ b/open-core/src/secure_eval_wrapper/live/bootstrap_cli.py @@ -0,0 +1,140 @@ +"""Command-line entry point for audited local Phase 8B PostgreSQL bootstrap.""" +from __future__ import annotations + +import argparse +import json +import sys + +from .bootstrap import ( + DEFAULT_DATABASE, + BootstrapSafetyError, + Phase8BOperatorBootstrap, + PostgresAdminTarget, +) + + +def _connection_arguments(parser: argparse.ArgumentParser) -> None: + parser.add_argument("--database", default=DEFAULT_DATABASE) + parser.add_argument("--host", default="127.0.0.1") + parser.add_argument("--port", type=int, default=5432) + parser.add_argument("--admin-database", default="postgres") + parser.add_argument("--admin-user", default="postgres") + parser.add_argument( + "--sslmode", + choices=("disable", "require", "verify-ca", "verify-full"), + default="disable", + ) + + +def _exact_plan_arguments(parser: argparse.ArgumentParser) -> None: + parser.add_argument("--expected-reviewed-sha", required=True) + parser.add_argument("--account-fingerprint", required=True) + parser.add_argument("--instrument", required=True) + + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser( + prog="secure-eval-live-bootstrap", + description=( + "Plan and initialize a dedicated local PostgreSQL Phase 8B read-only " + "operator database without loading credentials or making OKX requests." + ), + ) + commands = parser.add_subparsers(dest="command", required=True) + inspect_parser = commands.add_parser("inspect", help="read-only target inspection") + _connection_arguments(inspect_parser) + plan_parser = commands.add_parser("plan", help="produce a read-only hashed plan") + _connection_arguments(plan_parser) + _exact_plan_arguments(plan_parser) + initialize_parser = commands.add_parser( + "initialize", help="apply one exact previously reviewed plan" + ) + _connection_arguments(initialize_parser) + _exact_plan_arguments(initialize_parser) + initialize_parser.add_argument("--plan-hash", required=True) + initialize_parser.add_argument("--confirm-readonly-bootstrap", action="store_true") + verify_parser = commands.add_parser("verify", help="read-only readiness verification") + _connection_arguments(verify_parser) + _exact_plan_arguments(verify_parser) + return parser + + +def _target(args) -> PostgresAdminTarget: + return PostgresAdminTarget( + database=args.database, + host=args.host, + port=args.port, + admin_database=args.admin_database, + admin_user=args.admin_user, + sslmode=args.sslmode, + ) + + +def _failure_payload( + message: str, *, last_completed_stage: str | None = None +) -> dict[str, object]: + result: dict[str, object] = { + "status": "failed_closed", + "error": message, + "credentials_accessed": False, + "network_reads_occurred": False, + "network_writes_occurred": False, + "real_proof_executed": False, + } + if last_completed_stage is not None: + result["last_completed_stage"] = last_completed_stage + return result + + +def main(argv=None, *, bootstrap_factory=Phase8BOperatorBootstrap) -> int: + args = build_parser().parse_args(argv) + try: + bootstrap = bootstrap_factory(_target(args)) + if args.command == "inspect": + result = bootstrap.inspect_public() + elif args.command == "plan": + result = bootstrap.plan( + expected_reviewed_sha=args.expected_reviewed_sha, + account_fingerprint=args.account_fingerprint, + instrument=args.instrument, + ) + elif args.command == "initialize": + result = bootstrap.initialize( + expected_reviewed_sha=args.expected_reviewed_sha, + account_fingerprint=args.account_fingerprint, + instrument=args.instrument, + previous_plan_hash=args.plan_hash, + confirm_readonly_bootstrap=args.confirm_readonly_bootstrap, + ) + else: + result = bootstrap.verify( + expected_reviewed_sha=args.expected_reviewed_sha, + account_fingerprint=args.account_fingerprint, + instrument=args.instrument, + ) + except (BootstrapSafetyError, ValueError) as exc: + print(json.dumps( + _failure_payload( + str(exc), + last_completed_stage=getattr(exc, "last_completed_stage", None), + ), + sort_keys=True, + separators=(",", ":"), + )) + return 2 + except Exception: + print(json.dumps( + _failure_payload("local_postgresql_operation_failed"), + sort_keys=True, + separators=(",", ":"), + )) + return 3 + print(json.dumps(result, sort_keys=True, separators=(",", ":"), default=str)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) + + +__all__ = ["build_parser", "main"] diff --git a/open-core/src/secure_eval_wrapper/live/configuration.py b/open-core/src/secure_eval_wrapper/live/configuration.py index 6a5e1fc..a1bfbf5 100644 --- a/open-core/src/secure_eval_wrapper/live/configuration.py +++ b/open-core/src/secure_eval_wrapper/live/configuration.py @@ -8,7 +8,9 @@ from secure_eval_wrapper.data_collection.hashing import sha256_payload +from .endpoints import endpoint_catalog_hash from .identity import validate_okx_account_fingerprint +from .provider_identity import OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH def _text(value: str, name: str) -> str: @@ -136,6 +138,27 @@ def risk_limits(self) -> Mapping[str, object]: return MappingProxyType({name: getattr(self, name) for name in names}) +def guarded_configuration_from_json( + payload: Mapping[str, object], +) -> GuardedLiveConfiguration: + """Reconstruct a typed configuration without importing provider runtime code.""" + values = dict(payload) + for name in ( + "maximum_order_notional", "maximum_position_notional", "maximum_gross_exposure", + "maximum_net_exposure", "maximum_daily_submitted_notional", + "maximum_daily_realized_loss", "maximum_drawdown", "maximum_fee_bps", + "maximum_adverse_slippage_bps", "maximum_reference_price_deviation_bps", + ): + values[name] = Decimal(str(values[name])) + for name in ( + "allowed_instruments", "allowed_instrument_types", "allowed_settlement_assets", + "allowed_order_types", "credential_source_policy", + ): + values[name] = tuple(values[name]) + values.setdefault("maximum_transport_failures", 3) + return GuardedLiveConfiguration(**values) + + def phase8a_dry_run_configuration( *, account_fingerprint: str, @@ -162,4 +185,37 @@ def phase8a_dry_run_configuration( ) -__all__ = ["GuardedLiveConfiguration", "phase8a_dry_run_configuration"] +def phase8b_authenticated_readonly_configuration( + account_fingerprint: str, + instrument: str, +) -> GuardedLiveConfiguration: + """Build the one fixed, conservative Phase 8B operator-bootstrap profile.""" + if instrument != "BTC-USDT": + raise ValueError("Phase 8B operator bootstrap permits only the exact BTC-USDT spot instrument") + d = Decimal + return GuardedLiveConfiguration( + provider="okx", environment="production", account_fingerprint=account_fingerprint, + subaccount_fingerprint=None, allowed_instruments=(instrument,), + allowed_instrument_types=("spot",), allowed_settlement_assets=("USDT",), + base_currency="USDT", allowed_order_types=("limit",), maximum_order_notional=d("100"), + maximum_position_notional=d("250"), maximum_gross_exposure=d("500"), + maximum_net_exposure=d("500"), maximum_open_order_count=2, + maximum_daily_submitted_notional=d("1000"), maximum_daily_realized_loss=d("50"), + maximum_drawdown=d("50"), maximum_orders_per_minute=2, + maximum_cancellations_per_minute=3, market_data_freshness_seconds=30, + account_snapshot_freshness_seconds=30, reconciliation_freshness_seconds=30, + maximum_unknown_order_duration_seconds=30, + maximum_unacknowledged_order_duration_seconds=15, maximum_run_duration_seconds=300, + maximum_clock_skew_seconds=5, maximum_transport_failures=2, maximum_fee_bps=d("20"), + maximum_adverse_slippage_bps=d("50"), maximum_reference_price_deviation_bps=d("250"), + cancel_open_orders_on_kill=False, credential_source_policy=("environment",), + endpoint_catalog_hash=endpoint_catalog_hash(), + provider_implementation_hash=OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH, + ) + +__all__ = [ + "GuardedLiveConfiguration", + "guarded_configuration_from_json", + "phase8a_dry_run_configuration", + "phase8b_authenticated_readonly_configuration", +] diff --git a/open-core/src/secure_eval_wrapper/live/credentials.py b/open-core/src/secure_eval_wrapper/live/credentials.py index e2b0c74..a27c6f6 100644 --- a/open-core/src/secure_eval_wrapper/live/credentials.py +++ b/open-core/src/secure_eval_wrapper/live/credentials.py @@ -2,7 +2,6 @@ from __future__ import annotations import os -import re from abc import ABC, abstractmethod from collections.abc import Mapping from datetime import datetime @@ -10,48 +9,9 @@ from .gates import common_ci_indicators from .identity import validate_okx_account_fingerprint from .models import LiveCredentialReference - -_SECRET_KEY = re.compile(r"(api.?key|secret|passphrase|signature|authorization|cookie|token|ok-access-(?:key|sign|passphrase))", re.I) -_SECRET_VALUE = re.compile(r"(?i)(authorization:\s*|OK-ACCESS-(?:KEY|SIGN|PASSPHRASE)\s*[:=]\s*)[^\s,;]+") -_SECRET_QUERY = re.compile(r"([?&](?:api_?key|signature|token|passphrase|secret)=)[^&]+", re.I) -_EXPECTED_PERMISSION_NORMALIZATION = { - "read": "read", - "read_only": "read", - "trade": "trade", - "withdraw": "withdraw", -} - - -def redact(value): - if isinstance(value, Mapping): - return {str(key): ("[REDACTED]" if _SECRET_KEY.search(str(key)) else redact(item)) for key, item in value.items()} - if isinstance(value, (list, tuple)): - return type(value)(redact(item) for item in value) - if isinstance(value, str): - return _SECRET_QUERY.sub(r"\1[REDACTED]", _SECRET_VALUE.sub(r"\1[REDACTED]", value)) - return value - - -def normalize_expected_permission_summary(permissions: tuple[str, ...]) -> tuple[str, ...]: - values = tuple(permissions) - if not values: - return () - normalized = [] - for value in values: - if not isinstance(value, str) or not value or value != value.strip() or value != value.lower(): - raise PermissionError("expected credential permissions are malformed") - mapped = _EXPECTED_PERMISSION_NORMALIZATION.get(value) - if mapped is None: - raise PermissionError("expected credential permissions are unrecognized") - normalized.append(mapped) - return tuple(sorted(set(normalized))) - - -def validate_permission_summary(permissions: tuple[str, ...]) -> tuple[str, ...]: - normalized = normalize_expected_permission_summary(permissions) - if normalized != ("read",): - raise PermissionError("Phase 8A credential expectation must be exactly read-only") - return normalized +from .safety_policy import ( + normalize_expected_permission_summary, redact, validate_permission_summary, +) class LiveCredentialMaterial: diff --git a/open-core/src/secure_eval_wrapper/live/durable_repository.py b/open-core/src/secure_eval_wrapper/live/durable_repository.py index a9676ee..7c677a4 100644 --- a/open-core/src/secure_eval_wrapper/live/durable_repository.py +++ b/open-core/src/secure_eval_wrapper/live/durable_repository.py @@ -15,7 +15,7 @@ from .authorities import LiveRuntimeRiskState, OperationalPreflightEvidence, VerifiedOperationalSource from .collector_evidence import VerifiedOkxReadObservationBundle -from .credentials import normalize_expected_permission_summary, redact, validate_permission_summary +from .safety_policy import normalize_expected_permission_summary, redact, validate_permission_summary from .models import ( LiveKillState, LiveOrderState, @@ -29,7 +29,7 @@ from .reservations import calculate_live_reservation from .risk import evaluate_live_risk from .recovery import normalize_verified_recovery_observation -from .venues.okx_live import OkxProductionSpotAdapter +from .provider_identity import OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH class LiveConflictError(RuntimeError): @@ -189,7 +189,7 @@ def authenticated_readonly_storage_available(self) -> bool: return bool(row and row["table_ready"] and row["migration_ready"]) def load_guarded_live_configuration(self, configuration_hash: str): - from .readonly_preflight import guarded_configuration_from_json + from .configuration import guarded_configuration_from_json row = self._fetchone( "SELECT * FROM execution.live_configuration_snapshots WHERE configuration_sha256=%s", @@ -209,6 +209,93 @@ def load_guarded_live_configuration(self, configuration_hash: str): raise PermissionError("PostgreSQL guarded live configuration integrity check failed") return configuration + def persist_guarded_live_configuration_snapshot( + self, + *, + configuration, + created_at_utc, + fail_at=None, + ): + """Persist the exact Phase 8B bootstrap configuration as immutable authority.""" + from datetime import datetime + + from .configuration import ( + GuardedLiveConfiguration, + phase8b_authenticated_readonly_configuration, + ) + + if type(configuration) is not GuardedLiveConfiguration: + raise TypeError("configuration persistence requires exact GuardedLiveConfiguration") + if not isinstance(created_at_utc, datetime) or created_at_utc.tzinfo is None: + raise ValueError("created_at_utc must be a timezone-aware datetime") + if len(configuration.allowed_instruments) != 1: + raise ValueError("Phase 8B bootstrap requires one exact instrument") + expected = phase8b_authenticated_readonly_configuration( + configuration.account_fingerprint, + configuration.allowed_instruments[0], + ) + if configuration != expected: + raise PermissionError("configuration is not the exact current Phase 8B bootstrap profile") + + payload = { + name: getattr(configuration, name) + for name in configuration.__dataclass_fields__ + } + record_hash = sha256_payload(payload) + if record_hash != configuration.configuration_hash: + raise PermissionError("configuration hash derivation mismatch") + configuration_id = live_uuid( + "configuration", {"hash": configuration.configuration_hash} + ) + + with self.transaction(): + self._execute( + "SELECT pg_advisory_xact_lock(hashtextextended(%s,0))", + ("phase8b-bootstrap-global-configuration",), + ) + existing = self._fetchall( + "SELECT * FROM execution.live_configuration_snapshots " + "ORDER BY configuration_sha256 FOR UPDATE" + ) + if existing: + if len(existing) != 1 or ( + existing[0]["configuration_snapshot_id"] != configuration_id + or existing[0]["configuration_sha256"] != configuration.configuration_hash + or existing[0]["record_sha256"] != record_hash + or existing[0]["account_fingerprint"] != configuration.account_fingerprint + ): + raise LiveConflictError( + "bootstrap database already contains a different global guarded live configuration" + ) + if self.load_guarded_live_configuration(configuration.configuration_hash) != configuration: + raise PermissionError("exact configuration replay failed typed reload") + return configuration_id + if fail_at == "before_insert": + raise RuntimeError("injected configuration failure before insert") + self._strict_insert( + "execution.live_configuration_snapshots", + "configuration_snapshot_id", + configuration_id, + ( + "configuration_sha256", "provider", "environment", + "account_fingerprint", "dry_run", "read_only_preflight", + "production_write_enabled", "configuration_jsonb", "created_at_utc", + ), + ( + configuration.configuration_hash, configuration.provider, + configuration.environment, configuration.account_fingerprint, + configuration.dry_run, configuration.read_only_preflight, + configuration.production_write_enabled, _public_payload(payload), + created_at_utc, + ), + record_hash, + ) + if fail_at == "after_insert": + raise RuntimeError("injected configuration failure after insert") + if self.load_guarded_live_configuration(configuration.configuration_hash) != configuration: + raise PermissionError("persisted configuration failed exact typed reload") + return configuration_id + def persist_authenticated_readonly_proof( self, *, @@ -389,7 +476,7 @@ def _validate_start_bundle(*, configuration, credential_reference, account_snaps manifest.initial_account_snapshot_id == account_snapshot.snapshot_id, manifest.initial_account_snapshot_hash == account_snapshot.record_hash, manifest.credential_reference_hash == credential_reference.record_hash, - manifest.implementation_hash == report.implementation_hash == configuration.provider_implementation_hash, + manifest.implementation_hash == report.implementation_hash == configuration.provider_implementation_hash == OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH, kill_switch.state is LiveKillState.ARMED, manifest.dry_run and not manifest.production_write_enabled, ) @@ -865,6 +952,8 @@ def prepare_operational_dry_run( created_at_utc, fail_at=None, caller_risk_state=None, ): """Derive metadata normalization, risk, reservation, body, and hash under one lock.""" + from .venues.okx_live import OkxProductionSpotAdapter + del caller_risk_state with self.transaction(): self._lock_run(intent.live_run_id) @@ -1763,6 +1852,8 @@ def status(self, live_run_id): return self._fetchone("SELECT r.live_run_id,r.state,r.dry_run,r.production_write_enabled,r.started_at_utc,r.completed_at_utc,k.state AS kill_state,k.reason AS kill_reason,m.manifest_id,m.configuration_sha256,m.approval_id,m.preflight_report_id FROM execution.live_runs r JOIN execution.live_run_manifests m ON m.manifest_id=r.manifest_id JOIN execution.live_kill_switches k ON k.live_run_id=r.live_run_id WHERE r.live_run_id=%s", (live_run_id,)) def reconstruct(self, live_run_id): + from .venues.okx_live import OkxProductionSpotAdapter + run = self._fetchone("SELECT * FROM execution.live_runs WHERE live_run_id=%s", (live_run_id,)) manifest = self._fetchone("SELECT * FROM execution.live_run_manifests WHERE live_run_id=%s", (live_run_id,)) if run is None or manifest is None: raise LookupError("live run cannot be reconstructed") diff --git a/open-core/src/secure_eval_wrapper/live/preflight.py b/open-core/src/secure_eval_wrapper/live/preflight.py index 3f3fdc5..5f214fc 100644 --- a/open-core/src/secure_eval_wrapper/live/preflight.py +++ b/open-core/src/secure_eval_wrapper/live/preflight.py @@ -18,7 +18,7 @@ _issue_verified_source, ) from .collector_evidence import VerifiedOkxReadObservationBundle -from .credentials import normalize_expected_permission_summary +from .safety_policy import normalize_expected_permission_summary from .endpoints import endpoint_catalog_hash from .gates import common_ci_indicators from .identity import REPOSITORY_IDENTITY_RESOLVER_VERSION, RepositoryIdentityError, resolve_runtime_repository_identity, validate_git_commit_sha diff --git a/open-core/src/secure_eval_wrapper/live/provider_identity.py b/open-core/src/secure_eval_wrapper/live/provider_identity.py new file mode 100644 index 0000000..21b7904 --- /dev/null +++ b/open-core/src/secure_eval_wrapper/live/provider_identity.py @@ -0,0 +1,15 @@ +"""Transport-free provider implementation identity constants.""" +from secure_eval_wrapper.data_collection.hashing import sha256_payload + + +OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH = sha256_payload( + { + "adapter": "okx-production-spot", + "version": 5, + "writes": "phase8b-unreachable", + "authenticated_readonly_preflight": "exact-six-gets-unparameterized-positions", + } +) + + +__all__ = ["OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH"] diff --git a/open-core/src/secure_eval_wrapper/live/readonly_preflight.py b/open-core/src/secure_eval_wrapper/live/readonly_preflight.py index 180a8fd..fb190d5 100644 --- a/open-core/src/secure_eval_wrapper/live/readonly_preflight.py +++ b/open-core/src/secure_eval_wrapper/live/readonly_preflight.py @@ -12,7 +12,7 @@ from secure_eval_wrapper.data_collection.time_utils import require_utc_datetime from .collector_evidence import VerifiedOkxReadObservationBundle, expected_preflight_request_paths -from .configuration import GuardedLiveConfiguration +from .configuration import GuardedLiveConfiguration, guarded_configuration_from_json from .credentials import LiveCredentialProvider from .endpoints import endpoint_catalog_hash from .gates import common_ci_indicators @@ -61,33 +61,6 @@ def _at(value: object) -> datetime: ) -def guarded_configuration_from_json(payload: Mapping[str, object]) -> GuardedLiveConfiguration: - values = dict(payload) - for name in ( - "maximum_order_notional", - "maximum_position_notional", - "maximum_gross_exposure", - "maximum_net_exposure", - "maximum_daily_submitted_notional", - "maximum_daily_realized_loss", - "maximum_drawdown", - "maximum_fee_bps", - "maximum_adverse_slippage_bps", - "maximum_reference_price_deviation_bps", - ): - values[name] = Decimal(str(values[name])) - for name in ( - "allowed_instruments", - "allowed_instrument_types", - "allowed_settlement_assets", - "allowed_order_types", - "credential_source_policy", - ): - values[name] = tuple(values[name]) - values.setdefault("maximum_transport_failures", 3) - return GuardedLiveConfiguration(**values) - - @dataclass(frozen=True) class AuthenticatedReadOnlyProof: proof_id: UUID diff --git a/open-core/src/secure_eval_wrapper/live/safety_policy.py b/open-core/src/secure_eval_wrapper/live/safety_policy.py new file mode 100644 index 0000000..1fe0dd7 --- /dev/null +++ b/open-core/src/secure_eval_wrapper/live/safety_policy.py @@ -0,0 +1,71 @@ +"""Credential-free redaction and permission policy helpers.""" +from __future__ import annotations + +import re +from collections.abc import Mapping + + +_SECRET_KEY = re.compile( + r"(api.?key|secret|passphrase|signature|authorization|cookie|token|" + r"ok-access-(?:key|sign|passphrase))", + re.I, +) +_SECRET_VALUE = re.compile( + r"(?i)(authorization:\s*|OK-ACCESS-(?:KEY|SIGN|PASSPHRASE)\s*[:=]\s*)[^\s,;]+" +) +_SECRET_QUERY = re.compile( + r"([?&](?:api_?key|signature|token|passphrase|secret)=)[^&]+", re.I +) +_EXPECTED_PERMISSION_NORMALIZATION = { + "read": "read", + "read_only": "read", + "trade": "trade", + "withdraw": "withdraw", +} + + +def redact(value): + if isinstance(value, Mapping): + return { + str(key): ("[REDACTED]" if _SECRET_KEY.search(str(key)) else redact(item)) + for key, item in value.items() + } + if isinstance(value, (list, tuple)): + return type(value)(redact(item) for item in value) + if isinstance(value, str): + return _SECRET_QUERY.sub( + r"\1[REDACTED]", _SECRET_VALUE.sub(r"\1[REDACTED]", value) + ) + return value + + +def normalize_expected_permission_summary( + permissions: tuple[str, ...], +) -> tuple[str, ...]: + values = tuple(permissions) + if not values: + return () + normalized = [] + for value in values: + if ( + not isinstance(value, str) + or not value + or value != value.strip() + or value != value.lower() + ): + raise PermissionError("expected credential permissions are malformed") + mapped = _EXPECTED_PERMISSION_NORMALIZATION.get(value) + if mapped is None: + raise PermissionError("expected credential permissions are unrecognized") + normalized.append(mapped) + return tuple(sorted(set(normalized))) + + +def validate_permission_summary(permissions: tuple[str, ...]) -> tuple[str, ...]: + normalized = normalize_expected_permission_summary(permissions) + if normalized != ("read",): + raise PermissionError("Phase 8A credential expectation must be exactly read-only") + return normalized + + +__all__ = ["normalize_expected_permission_summary", "redact", "validate_permission_summary"] diff --git a/open-core/src/secure_eval_wrapper/live/venues/okx_live.py b/open-core/src/secure_eval_wrapper/live/venues/okx_live.py index c0fb645..ebe2bef 100644 --- a/open-core/src/secure_eval_wrapper/live/venues/okx_live.py +++ b/open-core/src/secure_eval_wrapper/live/venues/okx_live.py @@ -12,6 +12,8 @@ from secure_eval_wrapper.data_collection.hashing import sha256_payload +from ..provider_identity import OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH + from ..endpoints import EndpointClass, LiveOperation, OKX_PRODUCTION_ORIGIN, build_request_path, classify_exact, route_for from ..collector_evidence import QueryDisposition, _issue_okx_bundle, _issue_okx_envelope from ..gates import common_ci_indicators @@ -104,7 +106,7 @@ def _permissions(value: object) -> tuple[tuple[str, ...], tuple[str, ...]]: class OkxProductionSpotAdapter(GuardedLiveVenue): - provider_implementation_hash = sha256_payload({"adapter": "okx-production-spot", "version": 5, "writes": "phase8b-unreachable", "authenticated_readonly_preflight": "exact-six-gets-unparameterized-positions"}) + provider_implementation_hash = OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH def __init__(self, *, transport, credential_material=None, clock=None) -> None: self.transport = transport diff --git a/open-core/tests/test_phase8b_operator_bootstrap.py b/open-core/tests/test_phase8b_operator_bootstrap.py new file mode 100644 index 0000000..5d42286 --- /dev/null +++ b/open-core/tests/test_phase8b_operator_bootstrap.py @@ -0,0 +1,598 @@ +from __future__ import annotations + +import contextlib +import inspect +import io +import json +import socket +import subprocess +import sys +import tempfile +import unittest +from dataclasses import fields, replace +from pathlib import Path +from unittest.mock import patch + +from secure_eval_wrapper.data_collection.hashing import sha256_payload +from secure_eval_wrapper.live.bootstrap import ( + BootstrapOperationError, + BootstrapSafetyError, + DatabaseInspection, + DatabaseReference, + EXPECTED_MIGRATION_CATALOG, + Phase8BOperatorBootstrap, + PostgresAdminTarget, + _expected_database_objects, + derive_bootstrap_record_hash, + verify_local_migration_files, +) +from secure_eval_wrapper.live.bootstrap_cli import build_parser, main +from secure_eval_wrapper.live.configuration import ( + phase8b_authenticated_readonly_configuration, +) +from secure_eval_wrapper.live.identity import RuntimeRepositoryIdentity +from secure_eval_wrapper.live.provider_identity import ( + OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH, +) + + +SHA = "a" * 40 +FINGERPRINT = "1234567890abcdef" + + +def _reference( + target: PostgresAdminTarget, + *, + system_identifier: str = "7612345678901234567", + database_exists: bool = False, + database_oid: int | None = None, +) -> DatabaseReference: + core = { + "target_host": target.host, + "target_port": target.port, + "target_database": target.database, + "admin_database": target.admin_database, + "admin_user": target.admin_user, + "postgres_current_user": target.admin_user, + "postgres_system_identifier": system_identifier, + "postgres_server_version": "16.9", + "database_exists": database_exists, + "target_database_oid": database_oid, + } + return DatabaseReference(**core, database_identity_sha256=sha256_payload(core)) + + +class _PlanBootstrap(Phase8BOperatorBootstrap): + def __init__(self, target: PostgresAdminTarget, reference: DatabaseReference): + super().__init__( + target, + connector=lambda **kwargs: (_ for _ in ()).throw( + AssertionError("database connection attempted") + ), + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + ) + self.reference = reference + + def _database_reference(self, connection=None): + return self.reference + + def inspect(self, *, expected_configuration=None, reference=None, admin_connection=None): + reference = reference or self.reference + return DatabaseInspection( + reference=reference, + database_exists=reference.database_exists, + database_oid=reference.target_database_oid, + database_identity_sha256=reference.database_identity_sha256, + catalog_state="absent", + catalog=(), + latest_migration=None, + application_row_count=0, + configuration_row_count=0, + production_write_count=0, + non_system_object_count=0, + non_system_object_kinds=(), + blockers=(), + ) + + +class _FakeCliBootstrap: + calls = [] + + def __init__(self, target): + self.target = target + + def inspect_public(self): + self.calls.append(("inspect", self.target.database)) + return { + "action": "inspect", "credentials_accessed": False, + "network_reads_occurred": False, "network_writes_occurred": False, + "real_proof_executed": False, + } + + def plan(self, **kwargs): + self.calls.append(("plan", kwargs)) + return { + "action": "plan", "plan_hash": "b" * 64, "credentials_accessed": False, + "network_reads_occurred": False, "network_writes_occurred": False, + "real_proof_executed": False, + } + + def verify(self, **kwargs): + self.calls.append(("verify", kwargs)) + return { + "action": "verify", "ready_for_operator_authorization": False, + "credentials_accessed": False, "network_reads_occurred": False, + "network_writes_occurred": False, "real_proof_executed": False, + } + + def initialize(self, **kwargs): + self.calls.append(("initialize", kwargs)) + return { + "action": "initialize", "ready_for_operator_authorization": True, + "credentials_accessed": False, "network_reads_occurred": False, + "network_writes_occurred": False, "real_proof_executed": False, + } + + +class _ObjectOnlyBootstrap(Phase8BOperatorBootstrap): + class Connection: + def close(self): + pass + + def __init__(self, kind: str, reference: DatabaseReference): + super().__init__( + PostgresAdminTarget(), + connector=lambda **kwargs: self.Connection(), + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + ) + self.kind = kind + self.reference = reference + + def _verify_target_connection_identity(self, connection, reference): + return None + + def _fetchone(self, connection, statement, params=()): + if "to_regclass" in statement: + return (False,) + raise AssertionError(statement) + + def _inspect_database_objects(self, connection, *, exact_catalog): + self.assert_not_exact(exact_catalog) + return 1, (self.kind,), [], [] + + @staticmethod + def assert_not_exact(value): + if value: + raise AssertionError("uncatalogued database treated as exact") + + +class _InitializeHarness(Phase8BOperatorBootstrap): + def __init__(self, plan, reference, *, fail_create=False): + super().__init__( + PostgresAdminTarget(), + connector=lambda **kwargs: (_ for _ in ()).throw( + AssertionError("unexpected database connection") + ), + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + ) + self.plan_result = plan + self.reference = reference + self.fail_create = fail_create + self.created = False + + @contextlib.contextmanager + def _locked_admin_connection(self): + yield object() + + def plan(self, **kwargs): + return self.plan_result + + def _database_reference(self, connection=None): + return self.reference + + def _create_database(self, admin_connection): + if self.fail_create: + raise RuntimeError("injected create failure") + self.created = True + + +class Phase8BOperatorBootstrapOfflineTests(unittest.TestCase): + def test_factory_is_exact_readonly_spot_profile_with_current_hashes(self): + configuration = phase8b_authenticated_readonly_configuration(FINGERPRINT, "BTC-USDT") + from secure_eval_wrapper.live.venues.okx_live import OkxProductionSpotAdapter + + self.assertEqual(configuration.allowed_instruments, ("BTC-USDT",)) + self.assertEqual(configuration.allowed_instrument_types, ("spot",)) + self.assertEqual(configuration.allowed_settlement_assets, ("USDT",)) + self.assertEqual(configuration.base_currency, "USDT") + self.assertEqual(configuration.allowed_order_types, ("limit",)) + self.assertEqual(configuration.credential_source_policy, ("environment",)) + self.assertEqual( + configuration.provider_implementation_hash, + OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH, + ) + self.assertEqual( + configuration.provider_implementation_hash, + OkxProductionSpotAdapter.provider_implementation_hash, + ) + self.assertTrue(configuration.dry_run) + self.assertTrue(configuration.read_only_preflight) + self.assertFalse(configuration.production_write_enabled) + self.assertFalse(configuration.automatic_flatten) + self.assertFalse(configuration.allow_short) + self.assertFalse(configuration.allow_perpetual) + + def test_factory_rejects_bad_fingerprint_and_every_nonexact_instrument(self): + for fingerprint in ("", "0" * 16, "A" * 16, "short"): + with self.subTest(fingerprint=fingerprint), self.assertRaises(ValueError): + phase8b_authenticated_readonly_configuration(fingerprint, "BTC-USDT") + for instrument in ( + "btc-usdt", "BTC-USDT-SWAP", "BTC-USDT-PERP", "ETH-USDT", "BTC/USDT", "", + ): + with self.subTest(instrument=instrument), self.assertRaises(ValueError): + phase8b_authenticated_readonly_configuration(FINGERPRINT, instrument) + + def test_only_literal_local_hosts_are_accepted_before_connection(self): + for host in ("localhost", "db.internal", "192.168.1.20", "[::1]", " 127.0.0.1"): + with self.subTest(host=host), self.assertRaises(BootstrapSafetyError): + PostgresAdminTarget(host=host) + self.assertEqual(PostgresAdminTarget(host="127.0.0.1").host, "127.0.0.1") + self.assertEqual(PostgresAdminTarget(host="::1").host, "::1") + + def test_dedicated_database_and_maintenance_database_policy(self): + for database in ("secure_eval_wrapper", "postgres", "template0", "template1"): + with self.subTest(database=database), self.assertRaises(BootstrapSafetyError): + PostgresAdminTarget(database=database) + with self.assertRaises(BootstrapSafetyError): + PostgresAdminTarget(database="arbitrary_database") + with self.assertRaises(BootstrapSafetyError): + PostgresAdminTarget(admin_database="template1") + with self.assertRaises(BootstrapSafetyError): + PostgresAdminTarget( + database="secure_eval_phase8b", + admin_database="secure_eval_phase8b", + ) + self.assertEqual( + PostgresAdminTarget(database="secure_eval_phase8b_operator_01").database, + "secure_eval_phase8b_operator_01", + ) + + def test_production_database_policy_is_fixed_and_not_constructor_replaceable(self): + expected_fields = ( + "database", "host", "port", "admin_database", "admin_user", "sslmode", + ) + self.assertEqual( + tuple(inspect.signature(PostgresAdminTarget).parameters), expected_fields + ) + self.assertEqual( + tuple(item.name for item in fields(PostgresAdminTarget)), expected_fields + ) + post_init_source = inspect.getsource(PostgresAdminTarget.__post_init__) + self.assertIn( + r"secure_eval_phase8b(?:_[a-z0-9][a-z0-9_]{0,42})?", post_init_source + ) + self.assertNotIn("callable", post_init_source) + with self.assertRaises(TypeError): + PostgresAdminTarget(**{ + "database": "arbitrary_database", + "database_name_policy": lambda _: True, + }) + + def test_every_nonproduction_database_name_is_rejected_before_connection(self): + for database in ( + "secure_eval_wrapper", "postgres", "template0", "template1", + "arbitrary_database", "sew_phase8b_test", "secure_eval_phase8", + "secure_eval_phase8b-bad", "Secure_eval_phase8b", "secure_eval_phase8b test", + ): + with self.subTest(database=database), self.assertRaises( + (BootstrapSafetyError, ValueError) + ): + PostgresAdminTarget(database=database) + + def test_cli_rejects_nonproduction_database_before_bootstrap_construction(self): + output = io.StringIO() + + def fail_if_constructed(_target): + raise AssertionError("bootstrap constructed") + + with contextlib.redirect_stdout(output): + result = main( + ["inspect", "--database", "arbitrary_database"], + bootstrap_factory=fail_if_constructed, + ) + self.assertEqual(result, 2) + self.assertEqual(json.loads(output.getvalue())["status"], "failed_closed") + + def test_plan_hash_binds_all_public_connection_and_cluster_identity_fields(self): + target = PostgresAdminTarget() + reference = _reference(target) + service = _PlanBootstrap(target, reference) + first = service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + for key, value in reference.public_fields().items(): + self.assertEqual(first[key], value) + changed_cluster = _PlanBootstrap( + target, + _reference(target, system_identifier="7612345678901234568"), + ).plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + changed_user_target = replace(target, admin_user="phase8_admin") + changed_user = _PlanBootstrap( + changed_user_target, + _reference(changed_user_target), + ).plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertNotEqual(first["plan_hash"], changed_cluster["plan_hash"]) + self.assertNotEqual(first["plan_hash"], changed_user["plan_hash"]) + + def test_every_existing_uncatalogued_database_fails_closed(self): + reference = _reference( + PostgresAdminTarget(), database_exists=True, database_oid=16384 + ) + for kind in ("view", "function_or_procedure", "relation", "event_trigger"): + with self.subTest(kind=kind): + service = _ObjectOnlyBootstrap(kind, reference) + state = service._inspect_existing_database(reference) + self.assertEqual(state.catalog_state, "uncatalogued") + self.assertEqual(state.non_system_object_kinds, (kind,)) + self.assertIn( + "existing_uncatalogued_database_is_never_initialized", + state.blockers, + ) + + def test_verify_record_hash_is_independent_complete_and_deterministic(self): + configuration_hash = "c" * 64 + core = { + "database_identity_sha256": "d" * 64, + "postgres_system_identifier": "7612345678901234567", + "target_database_oid": 16384, + "observed_repository_sha": "a" * 40, + "expected_reviewed_sha": "a" * 40, + "migration_catalog": [{"migration_id": "0026", "sha256": "e" * 64}], + "phase8_tables_verified": True, + "phase8_indexes_verified": True, + "phase8_triggers_verified": True, + "configuration_hash": configuration_hash, + "account_fingerprint": FINGERPRINT, + "instrument": "BTC-USDT", + "ready_for_operator_authorization": True, + "production_write_count": 0, + "blockers": [], + "credentials_accessed": False, + "network_reads_occurred": False, + "network_writes_occurred": False, + "real_proof_executed": False, + } + first = derive_bootstrap_record_hash(core) + self.assertNotEqual(first, configuration_hash) + self.assertEqual(first, derive_bootstrap_record_hash(core)) + for key, value in ( + ("database_identity_sha256", "f" * 64), + ("observed_repository_sha", "b" * 40), + ("ready_for_operator_authorization", False), + ("blockers", ["not_ready"]), + ): + with self.subTest(key=key): + changed = dict(core) + changed[key] = value + self.assertNotEqual(first, derive_bootstrap_record_hash(changed)) + + def test_initialize_requires_confirmation_before_connection_or_mutation(self): + service = Phase8BOperatorBootstrap( + PostgresAdminTarget(), + connector=lambda **kwargs: (_ for _ in ()).throw( + AssertionError("connection attempted") + ), + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + ) + with self.assertRaises(BootstrapSafetyError): + service.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash="b" * 64, + confirm_readonly_bootstrap=False, + ) + + def test_initialize_failure_reports_exact_locked_stage(self): + reference = _reference(PostgresAdminTarget()) + plan = { + **reference.public_fields(), + "plan_hash": "b" * 64, + "blockers": [], + "database_creation_required": True, + "migrations_required": True, + } + service = _InitializeHarness(plan, reference, fail_create=True) + with self.assertRaises(BootstrapOperationError) as raised: + service.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash="b" * 64, + confirm_readonly_bootstrap=True, + ) + self.assertEqual(raised.exception.last_completed_stage, "plan_revalidated") + self.assertEqual(str(raised.exception), "local_postgresql_operation_failed") + + class _FailingCliBootstrap(_FakeCliBootstrap): + def initialize(self, **kwargs): + raise BootstrapOperationError( + "public_safe_failure", last_completed_stage="schema_verified" + ) + + output = io.StringIO() + argv = [ + "initialize", "--expected-reviewed-sha", SHA, + "--account-fingerprint", FINGERPRINT, "--instrument", "BTC-USDT", + "--plan-hash", "b" * 64, "--confirm-readonly-bootstrap", + ] + with contextlib.redirect_stdout(output): + self.assertEqual(main(argv, bootstrap_factory=_FailingCliBootstrap), 2) + payload = json.loads(output.getvalue()) + self.assertEqual(payload["last_completed_stage"], "schema_verified") + self.assertEqual(payload["error"], "public_safe_failure") + + def test_wrong_plan_hash_and_plan_blocker_fail_before_mutation(self): + reference = _reference(PostgresAdminTarget()) + base = { + **reference.public_fields(), + "plan_hash": "b" * 64, + "blockers": [], + "database_creation_required": True, + "migrations_required": True, + } + wrong = _InitializeHarness(base, reference) + with self.assertRaises(BootstrapSafetyError): + wrong.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash="c" * 64, + confirm_readonly_bootstrap=True, + ) + self.assertFalse(wrong.created) + blocked = _InitializeHarness( + {**base, "blockers": ["partial_catalog"]}, reference + ) + with self.assertRaises(BootstrapSafetyError): + blocked.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash="b" * 64, + confirm_readonly_bootstrap=True, + ) + self.assertFalse(blocked.created) + + def test_wrong_repository_sha_fails_before_database_connection(self): + service = Phase8BOperatorBootstrap( + PostgresAdminTarget(), + connector=lambda **kwargs: (_ for _ in ()).throw( + AssertionError("database connection occurred") + ), + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + ) + with self.assertRaises(BootstrapSafetyError): + service.plan( + expected_reviewed_sha="b" * 40, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + + def test_pinned_migration_and_expected_object_catalog(self): + paths = verify_local_migration_files() + self.assertEqual(len(paths), 26) + self.assertEqual( + EXPECTED_MIGRATION_CATALOG["0026_phase8b_authenticated_readonly_preflight"], + "698772fb68c5c4981256682d064c3be641193ab10c8dbf55e1a5b390ca7c504a", + ) + expected = _expected_database_objects(paths) + self.assertIn(("execution", "live_configuration_snapshots"), expected.tables) + self.assertIn( + ( + "execution", "live_configuration_snapshots", + "trg_live_configuration_snapshots_immutable", + "execution", "prevent_live_authority_mutation", + ), + expected.triggers, + ) + source = Path(__file__).resolve().parents[1] / "db" / "migrations" + with tempfile.TemporaryDirectory() as temporary: + root = Path(temporary) + for path in source.glob("*.sql"): + (root / path.name).write_bytes(path.read_bytes()) + target = root / "0026_phase8b_authenticated_readonly_preflight.sql" + target.write_bytes(target.read_bytes() + b"\n-- altered") + with self.assertRaises(BootstrapSafetyError): + verify_local_migration_files(root) + + def test_cli_has_no_destructive_options_and_read_commands_are_socket_free(self): + help_text = build_parser().format_help() + self.assertNotIn("--yes", help_text) + self.assertNotIn("drop", help_text.lower()) + self.assertNotIn("truncate", help_text.lower()) + _FakeCliBootstrap.calls.clear() + common = [ + "--expected-reviewed-sha", SHA, + "--account-fingerprint", FINGERPRINT, + "--instrument", "BTC-USDT", + ] + for command in ("inspect", "plan", "verify"): + argv = [command] if command == "inspect" else [command, *common] + output = io.StringIO() + with patch.object(socket, "socket", side_effect=AssertionError("socket used")): + with contextlib.redirect_stdout(output): + self.assertEqual(main(argv, bootstrap_factory=_FakeCliBootstrap), 0) + payload = json.loads(output.getvalue()) + rendered = json.dumps(payload).lower() + self.assertFalse(payload["credentials_accessed"]) + self.assertFalse(payload["network_reads_occurred"]) + for forbidden in ( + "api_secret", "passphrase", "ok-access-key", "balance", "account_uid", + ): + self.assertNotIn(forbidden, rendered) + + def test_initialize_cli_passes_only_exact_confirmation_and_plan_hash(self): + output = io.StringIO() + argv = [ + "initialize", "--expected-reviewed-sha", SHA, + "--account-fingerprint", FINGERPRINT, "--instrument", "BTC-USDT", + "--plan-hash", "b" * 64, "--confirm-readonly-bootstrap", + ] + with contextlib.redirect_stdout(output): + self.assertEqual(main(argv, bootstrap_factory=_FakeCliBootstrap), 0) + call = _FakeCliBootstrap.calls[-1] + self.assertEqual(call[0], "initialize") + self.assertTrue(call[1]["confirm_readonly_bootstrap"]) + self.assertEqual(call[1]["previous_plan_hash"], "b" * 64) + + def test_bootstrap_import_does_not_import_okx_transport_or_open_a_socket(self): + code = """ +import socket +import sys +socket.create_connection = lambda *a, **k: (_ for _ in ()).throw(AssertionError("socket")) +from secure_eval_wrapper.data_collection.hashing import sha256_payload +from secure_eval_wrapper.live.configuration import phase8b_authenticated_readonly_configuration +from secure_eval_wrapper.live.durable_repository import DurablePostgresLiveRepository +configuration = phase8b_authenticated_readonly_configuration("1234567890abcdef", "BTC-USDT") +payload = {name: getattr(configuration, name) for name in configuration.__dataclass_fields__} +repository = DurablePostgresLiveRepository(None) +repository._fetchone = lambda statement, params=(): { + "configuration_sha256": configuration.configuration_hash, + "configuration_jsonb": payload, + "record_sha256": sha256_payload(payload), + "dry_run": True, + "read_only_preflight": True, + "production_write_enabled": False, +} +assert repository.load_guarded_live_configuration(configuration.configuration_hash) == configuration +assert "secure_eval_wrapper.live.venues.okx_live" not in sys.modules +assert "secure_eval_wrapper.live.credentials" not in sys.modules +""" + completed = subprocess.run( + [sys.executable, "-c", code], text=True, capture_output=True, timeout=30 + ) + self.assertEqual(completed.returncode, 0, completed.stdout + completed.stderr) + + def test_stale_code_hash_overrides_are_rejected_by_typed_profile(self): + configuration = phase8b_authenticated_readonly_configuration(FINGERPRINT, "BTC-USDT") + self.assertNotEqual( + replace(configuration, endpoint_catalog_hash="f" * 64), configuration + ) + self.assertNotEqual( + replace(configuration, provider_implementation_hash="f" * 64), configuration + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/open-core/tests/test_phase8b_operator_bootstrap_postgres.py b/open-core/tests/test_phase8b_operator_bootstrap_postgres.py new file mode 100644 index 0000000..9426402 --- /dev/null +++ b/open-core/tests/test_phase8b_operator_bootstrap_postgres.py @@ -0,0 +1,698 @@ +from __future__ import annotations + +import os +import unittest +from concurrent.futures import ThreadPoolExecutor +from dataclasses import replace +from datetime import datetime, timezone +from unittest.mock import patch +from uuid import uuid4 + +from secure_eval_wrapper.data_collection.hashing import sha256_payload +from secure_eval_wrapper.live.bootstrap import ( + BootstrapOperationError, + BootstrapSafetyError, + EXPECTED_MIGRATION_CATALOG, + LATEST_MIGRATION, + Phase8BOperatorBootstrap, + PostgresAdminTarget, + derive_bootstrap_record_hash, +) +from secure_eval_wrapper.live.configuration import ( + phase8a_dry_run_configuration, + phase8b_authenticated_readonly_configuration, +) +from secure_eval_wrapper.live.durable_repository import ( + DurablePostgresLiveRepository, + LiveConflictError, + _public_payload, +) +from secure_eval_wrapper.live.endpoints import endpoint_catalog_hash +from secure_eval_wrapper.live.identity import RuntimeRepositoryIdentity +from secure_eval_wrapper.live.models import live_uuid +from secure_eval_wrapper.live.provider_identity import ( + OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH, +) + + +RUN = os.environ.get("RUN_POSTGRES_INTEGRATION", "").lower() == "true" +SHA = "a" * 40 +FINGERPRINT = "1234567890abcdef" + + +@unittest.skipUnless(RUN, "requires real PostgreSQL 16") +class Phase8BOperatorBootstrapPostgresTests(unittest.TestCase): + @classmethod + def setUpClass(cls): + import psycopg + + cls.psycopg = psycopg + cls.password = os.environ["POSTGRES_PASSWORD"] + suffix = uuid4().hex[:10] + cls.database = "secure_eval_phase8b_test_" + suffix + cls.partial_database = "secure_eval_phase8b_partial_" + suffix + cls.concurrent_database = "secure_eval_phase8b_concurrent_" + suffix + cls.oid_database = "secure_eval_phase8b_oid_" + suffix + cls.collation_database = "secure_eval_phase8b_collation_" + suffix + cls.external_empty_database = "secure_eval_phase8b_empty_" + suffix + cls.target = PostgresAdminTarget( + database=cls.database, + host=os.environ["POSTGRES_HOST"], + port=int(os.environ["POSTGRES_PORT"]), + admin_database="postgres", + admin_user=os.environ["POSTGRES_USER"], + sslmode=os.environ.get("POSTGRES_SSLMODE", "disable"), + ) + + def connector(**kwargs): + kwargs["password"] = cls.password + return psycopg.connect(**kwargs) + + cls.connector = staticmethod(connector) + cls.service = Phase8BOperatorBootstrap( + cls.target, + connector=connector, + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + clock=lambda: datetime(2026, 7, 15, tzinfo=timezone.utc), + ) + plan = cls.service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + cls.initial_plan = plan + cls.initialization_result = cls.service.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash=plan["plan_hash"], + confirm_readonly_bootstrap=True, + ) + cls.verification_result = cls.service.verify( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + + @classmethod + def tearDownClass(cls): + connection = cls.connector( + **cls.target.connection_kwargs("postgres", read_only=False) + ) + try: + connection.autocommit = True + with connection.cursor() as cursor: + for database in ( + cls.database, + cls.partial_database, + cls.concurrent_database, + cls.oid_database, + cls.collation_database, + cls.external_empty_database, + ): + if not database.startswith("secure_eval_phase8b_"): + raise AssertionError("refusing to clean non-test database") + cursor.execute( + "SELECT pg_terminate_backend(pid) FROM pg_stat_activity " + "WHERE datname=%s AND pid<>pg_backend_pid()", + (database,), + ) + from psycopg import sql + cursor.execute(sql.SQL("DROP DATABASE IF EXISTS {}").format(sql.Identifier(database))) + finally: + connection.close() + + def setUp(self): + self.connection = self.connector( + **self.target.connection_kwargs(self.database, read_only=False) + ) + with self.connection.cursor() as cursor: + cursor.execute("TRUNCATE execution.live_configuration_snapshots CASCADE") + self.connection.commit() + self.repository = DurablePostgresLiveRepository(self.connection) + self.configuration = phase8b_authenticated_readonly_configuration( + FINGERPRINT, "BTC-USDT" + ) + + def tearDown(self): + self.connection.close() + + def persist(self, configuration=None, **kwargs): + return self.repository.persist_guarded_live_configuration_snapshot( + configuration=self.configuration if configuration is None else configuration, + created_at_utc=datetime(2026, 7, 15, tzinfo=timezone.utc), + **kwargs, + ) + + def count(self): + with self.connection.cursor() as cursor: + cursor.execute("SELECT count(*) FROM execution.live_configuration_snapshots") + return cursor.fetchone()[0] + + def service_for_database(self, database): + target = replace(self.target, database=database) + return target, Phase8BOperatorBootstrap( + target, + connector=self.connector, + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + clock=lambda: datetime(2026, 7, 15, tzinfo=timezone.utc), + ) + + def create_external_database(self, target): + connection = self.connector( + **target.connection_kwargs("postgres", read_only=False) + ) + try: + connection.autocommit = True + from psycopg import sql + + with connection.cursor() as cursor: + cursor.execute( + sql.SQL("CREATE DATABASE {}").format(sql.Identifier(target.database)) + ) + finally: + connection.close() + + def assert_uncatalogued_database_was_not_migrated(self, target): + connection = self.connector( + **target.connection_kwargs(target.database, read_only=True) + ) + try: + with connection.cursor() as cursor: + cursor.execute( + "SELECT to_regclass('audit.schema_migrations') IS NOT NULL" + ) + self.assertFalse(cursor.fetchone()[0]) + finally: + connection.close() + + def assert_existing_uncatalogued_plan_fails_before_migrations(self, service, plan): + self.assertTrue(plan["database_exists"]) + self.assertFalse(plan["database_creation_required"]) + self.assertFalse(plan["migrations_required"]) + self.assertEqual(plan["migration_count_to_apply"], 0) + self.assertEqual(plan["catalog_state"], "uncatalogued") + self.assertIn( + "existing_uncatalogued_database_is_never_initialized", plan["blockers"] + ) + with self.assertRaises(BootstrapOperationError) as raised: + service.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash=plan["plan_hash"], + confirm_readonly_bootstrap=True, + ) + self.assertEqual(raised.exception.last_completed_stage, "target_lock_acquired") + + def test_existing_uncatalogued_collation_database_fails_closed_before_migrations(self): + target, service = self.service_for_database(self.collation_database) + self.create_external_database(target) + connection = self.connector( + **target.connection_kwargs(target.database, read_only=False) + ) + try: + with connection.cursor() as cursor: + cursor.execute('CREATE COLLATION public.phase8b_hidden FROM "C"') + connection.commit() + finally: + connection.close() + + inspection = service.inspect() + self.assertEqual(inspection.catalog_state, "uncatalogued") + self.assertIn( + "existing_uncatalogued_database_is_never_initialized", + inspection.blockers, + ) + plan = service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assert_existing_uncatalogued_plan_fails_before_migrations(service, plan) + self.assert_uncatalogued_database_was_not_migrated(target) + connection = self.connector( + **target.connection_kwargs(target.database, read_only=True) + ) + try: + with connection.cursor() as cursor: + cursor.execute( + "SELECT count(*) FROM pg_collation c JOIN pg_namespace n " + "ON n.oid=c.collnamespace WHERE n.nspname='public' " + "AND c.collname='phase8b_hidden'" + ) + self.assertEqual(cursor.fetchone()[0], 1) + finally: + connection.close() + + def test_externally_created_empty_uncatalogued_database_fails_closed_before_migrations(self): + target, service = self.service_for_database(self.external_empty_database) + self.create_external_database(target) + inspection = service.inspect() + self.assertEqual(inspection.catalog_state, "uncatalogued") + self.assertEqual(inspection.non_system_object_count, 0) + self.assertIn( + "existing_uncatalogued_database_is_never_initialized", + inspection.blockers, + ) + plan = service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assert_existing_uncatalogued_plan_fails_before_migrations(service, plan) + self.assert_uncatalogued_database_was_not_migrated(target) + + def test_clean_dedicated_initialization_reached_exact_0026_and_ready(self): + result = self.initialization_result + self.assertTrue(self.initial_plan["database_creation_required"]) + self.assertEqual(self.initial_plan["migration_count_to_apply"], 26) + self.assertEqual(result["migration_count"], 26) + self.assertEqual(result["latest_migration"], LATEST_MIGRATION) + self.assertTrue(result["immutable_catalog_verified"]) + self.assertTrue(result["migration_0026_installed"]) + self.assertFalse(result["credentials_accessed"]) + self.assertFalse(result["network_reads_occurred"]) + self.assertFalse(result["network_writes_occurred"]) + self.assertFalse(result["real_proof_executed"]) + self.assertTrue({ + "target_host", "target_port", "target_database", "admin_database", + "admin_user", "postgres_current_user", "postgres_system_identifier", + "postgres_server_version", "target_database_oid", "database_exists", + "database_identity_sha256", "observed_repository_sha", + "expected_reviewed_sha", "migration_count", "latest_migration", + "immutable_catalog_verified", "migration_0026_installed", + "configuration_snapshot_id", "configuration_hash", "account_fingerprint", + "instrument", "credential_policy", "endpoint_catalog_hash", + "adapter_implementation_hash", "dry_run", "read_only_preflight", + "production_write_enabled", "credentials_accessed", "network_reads_occurred", + "network_writes_occurred", "real_proof_executed", "bootstrap_record_hash", + } <= set(result)) + result_core = { + key: value for key, value in result.items() if key != "bootstrap_record_hash" + } + self.assertEqual( + result["bootstrap_record_hash"], derive_bootstrap_record_hash(result_core) + ) + self.assertNotEqual(result["bootstrap_record_hash"], result["configuration_hash"]) + verification = self.verification_result + self.assertTrue(verification["ready_for_operator_authorization"]) + self.assertTrue(verification["migration_hashes_verified"]) + self.assertTrue(verification["phase8_tables_verified"]) + self.assertTrue(verification["phase8_indexes_verified"]) + self.assertTrue(verification["phase8_triggers_verified"]) + self.assertEqual(verification["production_write_count"], 0) + verification_core = { + key: value for key, value in verification.items() + if key != "bootstrap_record_hash" + } + self.assertEqual( + verification["bootstrap_record_hash"], + derive_bootstrap_record_hash(verification_core), + ) + self.assertNotEqual( + verification["bootstrap_record_hash"], verification["configuration_hash"] + ) + self.persist() + repeated = self.service.verify( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertEqual( + verification["bootstrap_record_hash"], repeated["bootstrap_record_hash"] + ) + + def test_configuration_persistence_is_atomic_idempotent_and_exactly_reloadable(self): + expected_id = live_uuid( + "configuration", {"hash": self.configuration.configuration_hash} + ) + self.assertEqual(self.persist(), expected_id) + self.assertEqual(self.persist(), expected_id) + self.assertEqual(self.count(), 1) + self.assertEqual( + self.repository.load_guarded_live_configuration( + self.configuration.configuration_hash + ), + self.configuration, + ) + + def test_failures_before_and_after_insert_roll_back(self): + for fail_at in ("before_insert", "after_insert"): + with self.subTest(fail_at=fail_at): + with self.connection.cursor() as cursor: + cursor.execute("TRUNCATE execution.live_configuration_snapshots CASCADE") + self.connection.commit() + with self.assertRaises(RuntimeError): + self.persist(fail_at=fail_at) + self.connection.rollback() + self.assertEqual(self.count(), 0) + + def test_schema_contract_is_verified_before_configuration_insert(self): + plan_result = self.service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + incomplete_contract = { + "phase8_tables_verified": True, + "phase8_indexes_verified": False, + "phase8_triggers_verified": True, + } + with patch.object( + self.service, "_schema_contract", return_value=incomplete_contract + ): + with self.assertRaises(BootstrapOperationError) as raised: + self.service.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash=plan_result["plan_hash"], + confirm_readonly_bootstrap=True, + ) + self.assertEqual(raised.exception.last_completed_stage, "migrations_ready") + self.assertEqual(self.count(), 0) + + def test_stale_hashes_and_nonfactory_overrides_fail_before_insert(self): + for attacked in ( + replace(self.configuration, endpoint_catalog_hash="f" * 64), + replace(self.configuration, provider_implementation_hash="f" * 64), + ): + with self.subTest(hash=attacked.configuration_hash), self.assertRaises( + (PermissionError, ValueError) + ): + self.persist(attacked) + with self.assertRaises(ValueError): + replace(self.configuration, production_write_enabled=True) + self.assertEqual(self.count(), 0) + + def test_existing_conflicting_configuration_is_never_overwritten(self): + conflicting = phase8a_dry_run_configuration( + account_fingerprint=FINGERPRINT, + endpoint_catalog_hash=endpoint_catalog_hash(), + provider_implementation_hash=OKX_PRODUCTION_SPOT_ADAPTER_IMPLEMENTATION_HASH, + ) + payload = { + name: getattr(conflicting, name) + for name in conflicting.__dataclass_fields__ + } + with self.connection.cursor() as cursor: + cursor.execute( + "INSERT INTO execution.live_configuration_snapshots (" + "configuration_snapshot_id,configuration_sha256,provider,environment," + "account_fingerprint,dry_run,read_only_preflight,production_write_enabled," + "configuration_jsonb,record_sha256,created_at_utc) " + "VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)", + ( + live_uuid("configuration", {"hash": conflicting.configuration_hash}), + conflicting.configuration_hash, conflicting.provider, + conflicting.environment, conflicting.account_fingerprint, + conflicting.dry_run, conflicting.read_only_preflight, + conflicting.production_write_enabled, _public_payload(payload), + sha256_payload(payload), datetime(2026, 7, 15, tzinfo=timezone.utc), + ), + ) + self.connection.commit() + with self.assertRaises(LiveConflictError): + self.persist() + self.connection.rollback() + self.assertEqual(self.count(), 1) + + def test_concurrent_exact_replay_serializes_to_one_row(self): + def worker(): + connection = self.connector( + **self.target.connection_kwargs(self.database, read_only=False) + ) + try: + repository = DurablePostgresLiveRepository(connection) + return repository.persist_guarded_live_configuration_snapshot( + configuration=self.configuration, + created_at_utc=datetime(2026, 7, 15, tzinfo=timezone.utc), + ) + finally: + connection.close() + + with ThreadPoolExecutor(max_workers=2) as executor: + results = tuple(executor.map(lambda _: worker(), range(2))) + self.assertEqual(results[0], results[1]) + self.assertEqual(self.count(), 1) + + def test_concurrent_two_valid_fingerprints_enforce_global_singleton(self): + configuration_b = phase8b_authenticated_readonly_configuration( + "abcdef1234567890", "BTC-USDT" + ) + + def worker(configuration): + connection = self.connector( + **self.target.connection_kwargs(self.database, read_only=False) + ) + try: + repository = DurablePostgresLiveRepository(connection) + try: + repository.persist_guarded_live_configuration_snapshot( + configuration=configuration, + created_at_utc=datetime(2026, 7, 15, tzinfo=timezone.utc), + ) + except Exception as exc: + return "error", type(exc).__name__, configuration.account_fingerprint + return "ok", "persisted", configuration.account_fingerprint + finally: + connection.close() + + with ThreadPoolExecutor(max_workers=2) as executor: + outcomes = tuple(executor.map(worker, (self.configuration, configuration_b))) + self.assertEqual([outcome[0] for outcome in outcomes].count("ok"), 1) + self.assertEqual([outcome[0] for outcome in outcomes].count("error"), 1) + self.assertEqual( + [outcome[1] for outcome in outcomes if outcome[0] == "error"], + ["LiveConflictError"], + ) + self.assertEqual(self.count(), 1) + winner = next(outcome[2] for outcome in outcomes if outcome[0] == "ok") + verification = self.service.verify( + expected_reviewed_sha=SHA, + account_fingerprint=winner, + instrument="BTC-USDT", + ) + self.assertTrue(verification["ready_for_operator_authorization"]) + self.assertEqual(verification["production_write_count"], 0) + + def test_two_full_initialize_attempts_serialize_entire_operation(self): + target = replace(self.target, database=self.concurrent_database) + service = Phase8BOperatorBootstrap( + target, + connector=self.connector, + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + clock=lambda: datetime(2026, 7, 15, tzinfo=timezone.utc), + ) + plan = service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertFalse(plan["database_exists"]) + + def worker(): + try: + result = service.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash=plan["plan_hash"], + confirm_readonly_bootstrap=True, + ) + except Exception as exc: + return "error", type(exc).__name__ + return "ok", result["bootstrap_record_hash"] + + with ThreadPoolExecutor(max_workers=2) as executor: + outcomes = tuple(executor.map(lambda _: worker(), range(2))) + self.assertEqual([outcome[0] for outcome in outcomes].count("ok"), 1) + self.assertEqual([outcome[0] for outcome in outcomes].count("error"), 1) + self.assertEqual( + [outcome[1] for outcome in outcomes if outcome[0] == "error"], + ["BootstrapOperationError"], + ) + verification = service.verify( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertTrue(verification["ready_for_operator_authorization"]) + self.assertEqual(verification["production_write_count"], 0) + + def test_target_oid_change_after_plan_is_rejected(self): + target = replace(self.target, database=self.oid_database) + service = Phase8BOperatorBootstrap( + target, + connector=self.connector, + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + ) + admin = self.connector(**target.connection_kwargs("postgres", read_only=False)) + try: + admin.autocommit = True + from psycopg import sql + with admin.cursor() as cursor: + cursor.execute( + sql.SQL("CREATE DATABASE {}").format(sql.Identifier(self.oid_database)) + ) + finally: + admin.close() + plan = service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + original_oid = plan["target_database_oid"] + admin = self.connector(**target.connection_kwargs("postgres", read_only=False)) + try: + admin.autocommit = True + from psycopg import sql + with admin.cursor() as cursor: + cursor.execute( + sql.SQL("DROP DATABASE {}").format(sql.Identifier(self.oid_database)) + ) + cursor.execute( + sql.SQL("CREATE DATABASE {}").format(sql.Identifier(self.oid_database)) + ) + finally: + admin.close() + changed = service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertNotEqual(original_oid, changed["target_database_oid"]) + with self.assertRaises(BootstrapOperationError): + service.initialize( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + previous_plan_hash=plan["plan_hash"], + confirm_readonly_bootstrap=True, + ) + + def test_unexpected_non_internal_trigger_is_rejected(self): + with self.connection.cursor() as cursor: + cursor.execute( + "CREATE TRIGGER phase8b_unexpected_trigger BEFORE UPDATE ON " + "execution.live_configuration_snapshots FOR EACH ROW EXECUTE FUNCTION " + "execution.prevent_live_authority_mutation()" + ) + self.connection.commit() + try: + plan = self.service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertIn("non_internal_trigger_catalog_mismatch", plan["blockers"]) + finally: + with self.connection.cursor() as cursor: + cursor.execute( + "DROP TRIGGER phase8b_unexpected_trigger ON " + "execution.live_configuration_snapshots" + ) + self.connection.commit() + def test_altered_database_migration_hash_is_rejected_and_never_overwritten(self): + migration_id = next(iter(EXPECTED_MIGRATION_CATALOG)) + expected_hash = EXPECTED_MIGRATION_CATALOG[migration_id] + with self.connection.cursor() as cursor: + cursor.execute( + "UPDATE audit.schema_migrations SET sha256=%s WHERE migration_id=%s", + ("f" * 64, migration_id), + ) + self.connection.commit() + try: + plan_result = self.service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertEqual(plan_result["catalog_state"], "hash_or_filename_mismatch") + self.assertIn( + "immutable_migration_catalog_mismatch", plan_result["blockers"] + ) + self.assertEqual(self.count(), 0) + finally: + with self.connection.cursor() as cursor: + cursor.execute( + "UPDATE audit.schema_migrations SET sha256=%s WHERE migration_id=%s", + (expected_hash, migration_id), + ) + self.connection.commit() + + def test_unsafe_existing_application_rows_block_plan_without_overwrite(self): + with self.connection.cursor() as cursor: + cursor.execute("CREATE TABLE public.bootstrap_unsafe_test (value integer NOT NULL)") + cursor.execute("INSERT INTO public.bootstrap_unsafe_test VALUES (1)") + self.connection.commit() + try: + plan = self.service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertIn("existing_database_contains_unsafe_application_rows", plan["blockers"]) + self.assertTrue(plan["configuration_insertion_required"]) + self.assertEqual(self.count(), 0) + finally: + with self.connection.cursor() as cursor: + cursor.execute("DROP TABLE public.bootstrap_unsafe_test") + self.connection.commit() + def test_partial_catalog_and_existing_configuration_plan_fail_closed(self): + partial_target = replace(self.target, database=self.partial_database) + partial_service = Phase8BOperatorBootstrap( + partial_target, + connector=self.connector, + identity_resolver=lambda: RuntimeRepositoryIdentity(SHA, "git_checkout"), + ) + connection = self.connector( + **partial_target.connection_kwargs("postgres", read_only=False) + ) + try: + connection.autocommit = True + from psycopg import sql + with connection.cursor() as cursor: + cursor.execute(sql.SQL("CREATE DATABASE {}").format(sql.Identifier(self.partial_database))) + finally: + connection.close() + connection = self.connector( + **partial_target.connection_kwargs(self.partial_database, read_only=False) + ) + try: + with connection.cursor() as cursor: + cursor.execute("CREATE SCHEMA audit") + cursor.execute( + "CREATE TABLE audit.schema_migrations (migration_id text primary key," + "filename text not null unique,sha256 char(64) not null," + "applied_at_utc timestamptz not null default now(),description text not null)" + ) + cursor.execute( + "INSERT INTO audit.schema_migrations (migration_id,filename,sha256,description) " + "VALUES ('0001_initial_schema','0001_initial_schema.sql',%s,'partial test')", + ("598486e6af2eed4559564593adc0b66deff9e21ea91dbda560980c208a2950c5",), + ) + connection.commit() + finally: + connection.close() + partial = partial_service.plan( + expected_reviewed_sha=SHA, + account_fingerprint=FINGERPRINT, + instrument="BTC-USDT", + ) + self.assertEqual(partial["catalog_state"], "partial_catalog") + self.assertIn("partial_migration_catalog_is_never_auto_upgraded", partial["blockers"]) + + self.persist() + existing = self.service.plan( + expected_reviewed_sha=SHA, + account_fingerprint="abcdef1234567890", + instrument="BTC-USDT", + ) + self.assertIn("existing_guarded_live_configuration_conflicts", existing["blockers"]) + self.assertEqual(self.count(), 1) + + +if __name__ == "__main__": + unittest.main()