diff --git a/.gitignore b/.gitignore index 740f69cf..a2fbcb1a 100644 --- a/.gitignore +++ b/.gitignore @@ -50,6 +50,8 @@ temp/ # Market and analytical data # ------------------------------------------------------------ data/ +!/src/ai4binance/data/ +!/src/ai4binance/data/** datasets/ market_data/ historical_data/ diff --git a/config/operations/services.json b/config/operations/services.json index 4b3ee239..0e1f42a5 100644 --- a/config/operations/services.json +++ b/config/operations/services.json @@ -103,7 +103,7 @@ "command": "python -m ai4binance.cli.futures_multitf daemon", "mode": "RunFuturesMultiTf", "required": false, - "enabled": false, + "enabled": true, "health_file": "futures-multitf-health.json", "lock_file": "futures-multitf.lock", "max_health_age_seconds": 900, diff --git a/scripts/primary_local_reasoning_agent.ps1 b/scripts/primary_local_reasoning_agent.ps1 index c03903a3..0a906caa 100644 --- a/scripts/primary_local_reasoning_agent.ps1 +++ b/scripts/primary_local_reasoning_agent.ps1 @@ -9,6 +9,7 @@ $ErrorActionPreference = "Stop" $root = (Resolve-Path -LiteralPath (Join-Path $PSScriptRoot "..")).Path $taskName = "AI4BINANCE-Primary-Local-Reasoning" $healthPath = Join-Path $root "runtime\state\primary-local-reasoning-health.json" +$lockPath = Join-Path $root "runtime\state\primary-local-reasoning.lock" $serverScript = Join-Path $PSScriptRoot "start_llama_server.ps1" function Test-PrimaryListener { @@ -32,6 +33,7 @@ function Write-PrimaryHealth { endpoint = "http://127.0.0.1:8080" status = $Status updated_at = [DateTimeOffset]::UtcNow.ToString("o") + pid = $PID listener_pids = @($listeners | ForEach-Object { $_.OwningProcess } | Sort-Object -Unique) blockers = @($Blockers) execution_allowed = $false @@ -41,6 +43,31 @@ function Write-PrimaryHealth { } | ConvertTo-Json -Depth 4 | Set-Content -LiteralPath $healthPath -Encoding UTF8 } +function Write-PrimaryLock { + New-Item -ItemType Directory -Path (Split-Path $lockPath -Parent) -Force | Out-Null + $temporary = "$lockPath.$PID.tmp" + [System.IO.File]::WriteAllText( + $temporary, + [string]$PID, + [System.Text.Encoding]::ASCII + ) + Move-Item -LiteralPath $temporary -Destination $lockPath -Force +} + +function Remove-PrimaryLock { + if (-not (Test-Path -LiteralPath $lockPath -PathType Leaf)) { + return + } + try { + $ownerPid = [int](Get-Content -LiteralPath $lockPath -Raw -Encoding ASCII) + if ($ownerPid -eq $PID) { + Remove-Item -LiteralPath $lockPath -Force + } + } + catch { + } +} + function Start-PrimaryProvider { try { & $serverScript | Out-Null @@ -84,12 +111,14 @@ $hasHandle = $false try { $hasHandle = $mutex.WaitOne(0) if (-not $hasHandle) { exit 0 } + Write-PrimaryLock while ($true) { [void](Start-PrimaryProvider) Start-Sleep -Seconds $PollSeconds } } finally { + Remove-PrimaryLock if ($hasHandle) { $mutex.ReleaseMutex() } $mutex.Dispose() } diff --git a/scripts/ykb_daily_report.ps1 b/scripts/ykb_daily_report.ps1 index a6635e47..aea144f7 100644 --- a/scripts/ykb_daily_report.ps1 +++ b/scripts/ykb_daily_report.ps1 @@ -102,6 +102,7 @@ function Sync-LatestHumanReport { Reason = "ykb_report_MISSING" ReportId = "" Status = "" + Blockers = @() ObservedAt = "" MarkdownPath = "" LatestMarkdown = "" @@ -124,6 +125,7 @@ function Sync-LatestHumanReport { Reason = "YKB_HUMAN_REPORT_READY" ReportId = [string]$payload.report_id Status = [string]$payload.status + Blockers = @($payload.blockers) ObservedAt = [string]$payload.observed_at MarkdownPath = $markdownPath LatestMarkdown = $markdownPath @@ -135,6 +137,7 @@ function Sync-LatestHumanReport { Reason = "YKB_HUMAN_REPORT_SYNC_FAILED" ReportId = "" Status = "" + Blockers = @() ObservedAt = "" MarkdownPath = "" LatestMarkdown = "" @@ -184,10 +187,18 @@ function Write-Health { [AllowEmptyCollection()][string[]]$Blockers = @(), [int]$RefreshExitCode = 0, [int]$ReportExitCode = 0, + [string]$ReportStatus = "NOT_EVALUATED", + [AllowEmptyCollection()][string[]]$ReportBlockers = @(), [string]$LastReportObservedAt = "", [string]$LatestMarkdownPath = "" ) New-Item -ItemType Directory -Path $stateDirectory -Force | Out-Null + [string[]]$normalizedBlockers = @( + $Blockers | Where-Object { -not [string]::IsNullOrWhiteSpace($_) } + ) + [string[]]$normalizedReportBlockers = @( + $ReportBlockers | Where-Object { -not [string]::IsNullOrWhiteSpace($_) } + ) $payload = [ordered]@{ service = "ykb-report" status = $Status @@ -199,7 +210,9 @@ function Write-Health { latest_report_observed_at = $LastReportObservedAt refresh_exit_code = $RefreshExitCode report_exit_code = $ReportExitCode - blockers = @($Blockers) + report_status = $ReportStatus + report_blockers = $normalizedReportBlockers + blockers = $normalizedBlockers execution_allowed = $false live_eligibility_status = "LIVE_ORDER_BLOCKED" } @@ -437,20 +450,26 @@ function Invoke-YkbRefreshAndReport { } $state = Read-LatestReportState $humanReport = Sync-LatestHumanReport - $blockers = [System.Collections.Generic.List[string]]::new() + $operationalBlockers = [System.Collections.Generic.List[string]]::new() if ($refreshExitCode -ne 0) { - $blockers.Add("RUNTIME_RESEARCH_REFRESH_WITH_BLOCKERS") + $operationalBlockers.Add("RUNTIME_RESEARCH_REFRESH_WITH_BLOCKERS") } - if ($reportExitCode -ne 0) { - $blockers.Add("ykb_report_WITH_BLOCKERS") + if ($reportExitCode -notin @(0, 2)) { + $operationalBlockers.Add("YKB_REPORT_COMMAND_FAILED") + } + if ( + $reportExitCode -eq 2 -and + $humanReport.Status -ne "RUNNING_WITH_BLOCKERS" + ) { + $operationalBlockers.Add("YKB_REPORT_EXIT_STATUS_MISMATCH") } if (-not $state.Fresh) { - $blockers.Add([string]$state.Reason) + $operationalBlockers.Add([string]$state.Reason) } if (-not $humanReport.Synced) { - $blockers.Add([string]$humanReport.Reason) + $operationalBlockers.Add([string]$humanReport.Reason) } - $status = if ($state.Fresh -and $blockers.Count -eq 0) { + $status = if ($state.Fresh -and $operationalBlockers.Count -eq 0) { "READY" } elseif ($state.Fresh) { @@ -461,9 +480,11 @@ function Invoke-YkbRefreshAndReport { } Write-Health ` -Status $status ` - -Blockers $blockers.ToArray() ` + -Blockers $operationalBlockers.ToArray() ` -RefreshExitCode $refreshExitCode ` -ReportExitCode $reportExitCode ` + -ReportStatus ([string]$humanReport.Status) ` + -ReportBlockers @($humanReport.Blockers) ` -LastReportObservedAt ([string]$state.ObservedAt) ` -LatestMarkdownPath ([string]$humanReport.LatestMarkdown) if ($humanReport.Synced) { @@ -474,7 +495,7 @@ function Invoke-YkbRefreshAndReport { Write-Output "markdown_path=$($humanReport.MarkdownPath)" Write-Output "latest_markdown_path=$($humanReport.LatestMarkdown)" } - return $(if ($state.Fresh) { 0 } else { 1 }) + return $(if ($status -eq "READY") { 0 } else { 1 }) } function Start-YkbLoop { @@ -616,13 +637,17 @@ if ($Mode -eq "RunLoop") { $state = Read-LatestReportState if ($state.Fresh -and -not $Force) { $humanReport = Sync-LatestHumanReport + $status = if ($humanReport.Synced) { "READY" } else { "DEGRADED" } + $blockers = if ($humanReport.Synced) { @() } else { @($humanReport.Reason) } Write-Health ` - -Status "READY" ` - -Blockers @() ` + -Status $status ` + -Blockers $blockers ` + -ReportStatus ([string]$humanReport.Status) ` + -ReportBlockers @($humanReport.Blockers) ` -LastReportObservedAt ([string]$state.ObservedAt) ` -LatestMarkdownPath ([string]$humanReport.LatestMarkdown) Write-Output "ykb_report_FRESH" - exit 0 + exit $(if ($humanReport.Synced) { 0 } else { 1 }) } exit (Invoke-YkbRefreshAndReport) diff --git a/src/ai4binance/accounting/user_stream.py b/src/ai4binance/accounting/user_stream.py index b1f1b358..ebe44c58 100644 --- a/src/ai4binance/accounting/user_stream.py +++ b/src/ai4binance/accounting/user_stream.py @@ -16,6 +16,8 @@ from urllib.parse import urlencode from urllib.request import Request, urlopen +from websockets.exceptions import WebSocketException + from ai4binance.accounting.collectors import ( AccountingWebSocketCollector, ) @@ -306,67 +308,28 @@ def __post_init__(self) -> None: raise ValueError("WebSocket receive timeout is outside the safe range") def collect_once(self, *, sync_run_id: str) -> WebSocketCollectionSummary: - accepted = 0 - duplicates = 0 - rejected = 0 blockers: list[str] = [] collector = AccountingWebSocketCollector(self.ledger) deadline = time.monotonic() + self.collect_seconds - opened: list[UserDataStreamSession] = [] - now = datetime.now(UTC) + opened, accepted, duplicates, rejected = self._open_sessions( + sync_run_id=sync_run_id, + received_at=datetime.now(UTC), + blockers=blockers, + ) try: - for session in self.sessions: - try: - session.open() - opened.append(session) - if self._append_subscription_event(session, sync_run_id, now): - accepted += 1 - else: - duplicates += 1 - except (OSError, RuntimeError, ValueError) as error: - rejected += 1 - blockers.append( - f"{session.product_type.value}_WEBSOCKET_OPEN_FAILED:{type(error).__name__}" - ) - while ( - accepted + duplicates < self.event_limit and time.monotonic() < deadline - ): - made_progress = False - for session in opened: - if accepted + duplicates >= self.event_limit: - break - try: - event = session.receive_event( - timeout_seconds=self.receive_timeout_seconds - ) - except TimeoutError: - continue - except (OSError, RuntimeError, ValueError) as error: - blockers.append( - f"{session.product_type.value}_WEBSOCKET_EVENT_FAILED:{type(error).__name__}" - ) - continue - result = collector.ingest_event( - event, - snapshot_id=sync_run_id, - sync_run_id=sync_run_id, - received_at=datetime.now(UTC), - ) - accepted += result.accepted_count - duplicates += result.duplicate_count - rejected += result.rejected_count - blockers.extend(result.blockers) - made_progress = True - if not made_progress and not opened: - break + event_counts = self._collect_events( + collector=collector, + opened=opened, + sync_run_id=sync_run_id, + deadline=deadline, + accepted=accepted, + duplicates=duplicates, + rejected=rejected, + blockers=blockers, + ) + accepted, duplicates, rejected = event_counts finally: - for session in reversed(opened): - try: - session.close() - except (OSError, RuntimeError, ValueError): - blockers.append( - f"{session.product_type.value}_WEBSOCKET_CLOSE_FAILED" - ) + self._close_sessions(opened, blockers) return WebSocketCollectionSummary( "COLLECTED" if not blockers and rejected == 0 else "DEGRADED", accepted, @@ -375,6 +338,88 @@ def collect_once(self, *, sync_run_id: str) -> WebSocketCollectionSummary: tuple(dict.fromkeys(blockers)), ) + def _open_sessions( + self, + *, + sync_run_id: str, + received_at: datetime, + blockers: list[str], + ) -> tuple[list[UserDataStreamSession], int, int, int]: + opened: list[UserDataStreamSession] = [] + accepted = duplicates = rejected = 0 + for session in self.sessions: + try: + session.open() + opened.append(session) + if self._append_subscription_event(session, sync_run_id, received_at): + accepted += 1 + else: + duplicates += 1 + except (OSError, RuntimeError, ValueError, WebSocketException) as error: + rejected += 1 + blockers.append( + f"{session.product_type.value}_WEBSOCKET_OPEN_FAILED:{type(error).__name__}" + ) + return opened, accepted, duplicates, rejected + + def _collect_events( + self, + *, + collector: AccountingWebSocketCollector, + opened: list[UserDataStreamSession], + sync_run_id: str, + deadline: float, + accepted: int, + duplicates: int, + rejected: int, + blockers: list[str], + ) -> tuple[int, int, int]: + while accepted + duplicates < self.event_limit and time.monotonic() < deadline: + made_progress = False + for session in opened: + if accepted + duplicates >= self.event_limit: + break + try: + event = session.receive_event( + timeout_seconds=self.receive_timeout_seconds + ) + except TimeoutError: + continue + except ( + OSError, + RuntimeError, + ValueError, + WebSocketException, + ) as error: + blockers.append( + f"{session.product_type.value}_WEBSOCKET_EVENT_FAILED:{type(error).__name__}" + ) + continue + result = collector.ingest_event( + event, + snapshot_id=sync_run_id, + sync_run_id=sync_run_id, + received_at=datetime.now(UTC), + ) + accepted += result.accepted_count + duplicates += result.duplicate_count + rejected += result.rejected_count + blockers.extend(result.blockers) + made_progress = True + if not made_progress and not opened: + break + return accepted, duplicates, rejected + + @staticmethod + def _close_sessions( + opened: list[UserDataStreamSession], blockers: list[str] + ) -> None: + for session in reversed(opened): + try: + session.close() + except (OSError, RuntimeError, ValueError, WebSocketException): + blockers.append(f"{session.product_type.value}_WEBSOCKET_CLOSE_FAILED") + def _append_subscription_event( self, session: UserDataStreamSession, diff --git a/src/ai4binance/application/runtime.py b/src/ai4binance/application/runtime.py index 86e043de..a4ef5971 100644 --- a/src/ai4binance/application/runtime.py +++ b/src/ai4binance/application/runtime.py @@ -12,6 +12,9 @@ from ai4binance.domain import Action +_PORTFOLIO_COST_BASIS_ASSET_LIMIT = 50 +_PAR_USDT_ASSETS = frozenset({"USDT", "FDUSD", "USDC"}) + class RuntimeState(StrEnum): STARTING = "STARTING" @@ -205,8 +208,15 @@ def run(self, now: datetime | None = None) -> DualMarketAdvisoryReport: futures_account = self._futures_account(symbol, created_at, futures_preflight) blockers = [*spot_preflight, *futures_preflight] cost_basis = self._cost_basis(symbol, spot_wallet, blockers) + portfolio_average_costs = self._portfolio_average_costs( + spot_wallet, + cost_basis, + blockers, + ) portfolio_analytics = self._portfolio_analytics( - spot_wallet, blockers, cost_basis + spot_wallet, + blockers, + portfolio_average_costs, ) spot_ready = spot_wallet is not None and not spot_preflight futures_ready = futures_account is not None and not futures_preflight @@ -318,22 +328,14 @@ def _portfolio_analytics( self, wallet: Any | None, blockers: list[str], - cost_basis: Any | None, + average_costs_usdt: Mapping[str, Decimal] | None, ) -> Any | None: if wallet is None or self.analytics_service is None: return None try: - costs = ( - {cost_basis.base_asset: cost_basis.average_cost_quote} - if cost_basis is not None - and cost_basis.average_cost_quote is not None - and not cost_basis.blockers - and cost_basis.quote_asset == "USDT" - else None - ) analytics = self.analytics_service.evaluate( wallet, - average_costs_usdt=costs, + average_costs_usdt=average_costs_usdt, ) except (OSError, RuntimeError, TypeError, ValueError): blockers.append("PORTFOLIO_ANALYTICS_FAILED") @@ -351,9 +353,11 @@ def _cost_basis( return None base_asset = symbol.removesuffix("USDT") or symbol balance = wallet.balance(base_asset) - wallet_quantity = ( - balance.free + balance.locked if balance is not None else Decimal("0") - ) + if balance is None: + return None + wallet_quantity = balance.free + balance.locked + if wallet_quantity <= Decimal("0"): + return None try: report = self.cost_basis_service.evaluate(symbol, wallet_quantity) except (OSError, RuntimeError, TypeError, ValueError): @@ -362,6 +366,60 @@ def _cost_basis( blockers.extend(report.blockers) return report + def _portfolio_average_costs( + self, + wallet: Any | None, + selected_cost_basis: Any | None, + blockers: list[str], + ) -> Mapping[str, Decimal] | None: + if wallet is None or self.cost_basis_service is None: + return None + costs: dict[str, Decimal] = {} + evaluated_assets: set[str] = set() + if selected_cost_basis is not None: + evaluated_assets.add(selected_cost_basis.base_asset) + if ( + selected_cost_basis.average_cost_quote is not None + and not selected_cost_basis.blockers + and selected_cost_basis.quote_asset == "USDT" + ): + costs[selected_cost_basis.base_asset] = ( + selected_cost_basis.average_cost_quote + ) + balances = tuple( + sorted( + ( + balance + for balance in wallet.balances + if balance.asset not in _PAR_USDT_ASSETS + and balance.free + balance.locked > Decimal("0") + ), + key=lambda balance: balance.asset, + ) + ) + if len(balances) > _PORTFOLIO_COST_BASIS_ASSET_LIMIT: + blockers.append("PORTFOLIO_COST_BASIS_SCOPE_EXCEEDED") + for balance in balances[:_PORTFOLIO_COST_BASIS_ASSET_LIMIT]: + if balance.asset in evaluated_assets: + continue + evaluated_assets.add(balance.asset) + try: + report = self.cost_basis_service.evaluate( + f"{balance.asset}USDT", + balance.free + balance.locked, + ) + except (OSError, RuntimeError, TypeError, ValueError): + blockers.append("COST_BASIS_RECONCILIATION_FAILED") + continue + blockers.extend(report.blockers) + if ( + report.average_cost_quote is not None + and not report.blockers + and report.quote_asset == "USDT" + ): + costs[report.base_asset] = report.average_cost_quote + return costs or None + @staticmethod def _unavailable_advisory(market: str, blockers: list[str]) -> MarketAdvisory: return MarketAdvisory( diff --git a/src/ai4binance/cli/futures_multitf.py b/src/ai4binance/cli/futures_multitf.py index ad205845..778f6661 100644 --- a/src/ai4binance/cli/futures_multitf.py +++ b/src/ai4binance/cli/futures_multitf.py @@ -27,6 +27,9 @@ estimate_measurable_trade_plan, ) from ai4binance.indicators import atr +from ai4binance.integrations.research_market_universe import ( + RESEARCH_MARKET_UNIVERSE_SOURCE, +) from ai4binance.ops import SingleInstanceLease from ai4binance.reporting import to_primitive from ai4binance.research.futures_multitf import ( @@ -57,6 +60,7 @@ _ADVERSE_OUTCOMES = frozenset({"INVALIDATED_FIRST", "ADVERSE"}) _TUNING_TRIGGER_STREAK = 3 _MAX_COMPLETED_TUNING_TRIGGERS = 200 +_UNIVERSE_MAX_AGE = timedelta(minutes=5) def _absolute(path: Path) -> Path: @@ -127,6 +131,20 @@ def _eligible_symbols(cache_root: Path) -> tuple[str, ...]: if path.is_symlink() or not path.is_file() or path.stat().st_size > 2_000_000: raise ValueError("FUTURES_MULTITF_UNIVERSE_UNAVAILABLE") payload = _load_mapping(path) + try: + observed_at = datetime.fromisoformat(str(payload["observed_at"])) + except (KeyError, TypeError, ValueError): + raise ValueError("FUTURES_MULTITF_UNIVERSE_INVALID") from None + if observed_at.utcoffset() is None: + raise ValueError("FUTURES_MULTITF_UNIVERSE_INVALID") + age = datetime.now(UTC) - observed_at + if ( + not timedelta(0) <= age <= _UNIVERSE_MAX_AGE + or payload.get("source") != RESEARCH_MARKET_UNIVERSE_SOURCE + or payload.get("execution_allowed") is not False + or payload.get("live_eligibility_status") != "LIVE_ORDER_BLOCKED" + ): + raise ValueError("FUTURES_MULTITF_UNIVERSE_INVALID") raw = payload.get("futures_symbols") if not isinstance(raw, list) or not raw or len(raw) > 5_000: raise ValueError("FUTURES_MULTITF_UNIVERSE_INVALID") @@ -147,6 +165,41 @@ def _eligible_symbols(cache_root: Path) -> tuple[str, ...]: return symbols +def _retain_current_universe_state( + state: Mapping[str, object], symbols: tuple[str, ...] +) -> dict[str, object]: + """Drop state projections and cursors that refer to an older universe.""" + allowed = frozenset(symbols) + previous = state.get("eligible_symbols") + same_universe = isinstance(previous, list) and tuple(previous) == symbols + retained = dict(state) if same_universe else {} + for field in ( + "attempted_windows", + "completed_windows", + "retry_after", + "unavailable_windows", + ): + raw = state.get(field) + retained[field] = ( + { + str(symbol): value + for symbol, value in raw.items() + if isinstance(symbol, str) and symbol in allowed + } + if isinstance(raw, Mapping) + else {} + ) + for field in ("active_symbol", "monitor_symbol"): + if retained.get(field) not in allowed: + retained.pop(field, None) + retained["eligible_symbols"] = list(symbols) + retained["universe_source"] = RESEARCH_MARKET_UNIVERSE_SOURCE + if not same_universe: + retained["completed_tuning_triggers"] = [] + retained.update(_SAFE_STATE) + return retained + + def _next_symbol(symbols: tuple[str, ...], cursor: object) -> str: """Return a stable round-robin Futures monitor target.""" @@ -581,13 +634,14 @@ def run_cycle(settings: Settings, *, observed_at: datetime) -> dict[str, object] start_day = end_day - timedelta(days=27) window = f"{start_day.isoformat()}_to_{end_day.isoformat()}" state = _load_mapping(state_path) + symbols = _eligible_symbols(cache_root) + state = _retain_current_universe_state(state, symbols) raw_attempted = state.get("attempted_windows", {}) attempted = ( {str(key): str(value) for key, value in raw_attempted.items()} if isinstance(raw_attempted, Mapping) else {} ) - symbols = _eligible_symbols(cache_root) raw_completed = state.get("completed_windows", {}) completed = dict(raw_completed) if isinstance(raw_completed, Mapping) else {} raw_retry = state.get("retry_after", {}) @@ -677,6 +731,8 @@ def retry_due(item: str) -> bool: "observed_at": now.isoformat(), "window": window, "eligible_symbol_count": len(symbols), + "eligible_symbols": list(symbols), + "universe_source": RESEARCH_MARKET_UNIVERSE_SOURCE, "attempted_symbol_count": sum( attempted.get(item) == window for item in symbols ), @@ -826,6 +882,8 @@ def retry_due(item: str) -> bool: "active_timeframe": active_timeframe, "phase": phase, "eligible_symbol_count": len(symbols), + "eligible_symbols": list(symbols), + "universe_source": RESEARCH_MARKET_UNIVERSE_SOURCE, "attempted_symbol_count": sum( attempted.get(item) == window for item in symbols ), diff --git a/src/ai4binance/cli/market_data.py b/src/ai4binance/cli/market_data.py index e1a6ccf4..9b47c2d3 100644 --- a/src/ai4binance/cli/market_data.py +++ b/src/ai4binance/cli/market_data.py @@ -30,6 +30,7 @@ MarketHistorySupervisor, MarketHistorySynchronizer, ) +from ai4binance.data.market_universe_retention import MarketUniverseRetention from ai4binance.domain.opportunity_observation import ( has_complete_measurable_opportunity, ) @@ -38,6 +39,10 @@ BinanceMarketUniverseProvider, ReadOnlyBinanceJsonTransport, ) +from ai4binance.integrations.research_market_universe import ( + ReadOnlyCoinGeckoJsonTransport, + ResearchMarketUniverseProvider, +) from ai4binance.ops.runtime import SingleInstanceLease _SAFE_STATE: dict[str, object] = { @@ -170,6 +175,11 @@ def build_continuous_market_history( spot=synchronizer.universe_provider.spot_transport, futures=synchronizer.universe_provider.futures_transport, initial_days=settings.market_history_initial_days, + enrichment_days=getattr( + settings, + "market_history_enrichment_days", + settings.market_history_initial_days, + ), pages_per_stream=settings.market_history_pages_per_stream, max_workers=settings.market_history_max_workers, minimum_candles=settings.minimum_closed_candles, @@ -188,6 +198,28 @@ def build_continuous_market_history( settings, root, include_futures=True ), on_symbol_screen=_build_opportunity_screen(settings, root), + retention=MarketUniverseRetention( + archive_root=_absolute(settings.dataset_directory), + source_cache_root=_absolute( + getattr( + settings, + "market_history_source_cache_directory", + settings.dataset_directory.parent / "market_sources", + ) + ), + opportunity_monitor_root=( + root / "runtime/artifacts/opportunity-radar/monitor" + ).resolve(), + futures_replay_root=( + root / "runtime/data/datasets/futures/multitf" + ).resolve(), + futures_artifact_roots=( + (root / "runtime/artifacts/validation/futures_multitf").resolve(), + ( + root / "runtime/artifacts/validation/futures_failure_tuning" + ).resolve(), + ), + ), ) @@ -394,7 +426,7 @@ def run_market_history_command( settings, synchronizer, root=Path.cwd(), - include_coin_m=True, + include_coin_m=getattr(settings, "market_history_coin_m_enabled", False), ) if command == "market-history-sync" and as_of is None: with SingleInstanceLease(synchronizer.state_path.with_suffix(".lock")): @@ -503,48 +535,64 @@ def build_market_history_synchronizer(settings: Settings) -> MarketHistorySynchr futures_budget = PublicRequestBudget( governor=WeightedRateLimitGovernor(bands=bands) ) - return MarketHistorySynchronizer( - universe_provider=BinanceMarketUniverseProvider( - spot_transport=MeteredPublicTransport( - ReadOnlyBinanceJsonTransport( - base_url="https://api.binance.com", - allowed_prefixes=("/api/v3/",), - timeout_seconds=settings.request_timeout_seconds, - max_attempts=1, - backoff_seconds=settings.request_backoff_seconds, - response_headers_observer=spot_budget.observe_headers, - ), - _absolute(settings.dataset_directory) / "spot" / "metadata", - budget=spot_budget, + binance_provider = BinanceMarketUniverseProvider( + spot_transport=MeteredPublicTransport( + ReadOnlyBinanceJsonTransport( + base_url="https://api.binance.com", + allowed_prefixes=("/api/v3/",), + timeout_seconds=settings.request_timeout_seconds, + max_attempts=1, + backoff_seconds=settings.request_backoff_seconds, + response_headers_observer=spot_budget.observe_headers, ), - futures_transport=MeteredPublicTransport( - ReadOnlyBinanceJsonTransport( - base_url="https://fapi.binance.com", - allowed_prefixes=("/fapi/v1/", "/futures/data/"), - timeout_seconds=settings.request_timeout_seconds, - max_attempts=1, - backoff_seconds=settings.request_backoff_seconds, - response_headers_observer=futures_budget.observe_headers, - ), - _absolute(settings.dataset_directory) / "usd_m_futures" / "metadata", - budget=futures_budget, + _absolute(settings.dataset_directory) / "spot" / "metadata", + budget=spot_budget, + ), + futures_transport=MeteredPublicTransport( + ReadOnlyBinanceJsonTransport( + base_url="https://fapi.binance.com", + allowed_prefixes=("/fapi/v1/", "/futures/data/"), + timeout_seconds=settings.request_timeout_seconds, + max_attempts=1, + backoff_seconds=settings.request_backoff_seconds, + response_headers_observer=futures_budget.observe_headers, ), - coin_m_transport=MeteredPublicTransport( - ReadOnlyBinanceJsonTransport( - base_url="https://dapi.binance.com", - allowed_prefixes=("/dapi/v1/", "/futures/data/"), - timeout_seconds=settings.request_timeout_seconds, - max_attempts=1, - response_headers_observer=futures_budget.observe_headers, - ), - _absolute(settings.dataset_directory) / "coin_m_futures" / "metadata", - budget=futures_budget, - ) - if settings.market_history_coin_m_enabled - else None, - quote_assets=settings.preferred_quote_assets, - futures_symbol_exclusions=settings.futures_symbol_exclusions, + _absolute(settings.dataset_directory) / "usd_m_futures" / "metadata", + budget=futures_budget, ), + coin_m_transport=MeteredPublicTransport( + ReadOnlyBinanceJsonTransport( + base_url="https://dapi.binance.com", + allowed_prefixes=("/dapi/v1/", "/futures/data/"), + timeout_seconds=settings.request_timeout_seconds, + max_attempts=1, + response_headers_observer=futures_budget.observe_headers, + ), + _absolute(settings.dataset_directory) / "coin_m_futures" / "metadata", + budget=futures_budget, + ) + if settings.market_history_coin_m_enabled + else None, + quote_assets=settings.preferred_quote_assets, + futures_symbol_exclusions=settings.futures_symbol_exclusions, + ) + universe_provider = ResearchMarketUniverseProvider( + binance=binance_provider, + market_cap_transport=ReadOnlyCoinGeckoJsonTransport( + timeout_seconds=settings.request_timeout_seconds, + max_attempts=min(settings.request_max_attempts, 3), + ), + wallet_balance_path=( + _absolute(settings.binance_accounting_directory) + / "spot" + / "balance_snapshots.jsonl" + ), + wallet_minimum_value_usdt=(settings.market_history_wallet_minimum_value_usdt), + wallet_maximum_age=timedelta(minutes=settings.accounting_freshness_minutes), + market_cap_asset_limit=settings.market_history_market_cap_asset_limit, + ) + return MarketHistorySynchronizer( + universe_provider=universe_provider, # type: ignore[arg-type] archive_root=_absolute(settings.dataset_directory), source_cache=BinanceVisionArchiveCache.with_network( _absolute(settings.market_history_source_cache_directory), diff --git a/src/ai4binance/cli/runtime.py b/src/ai4binance/cli/runtime.py index fc51f892..67e1779d 100644 --- a/src/ai4binance/cli/runtime.py +++ b/src/ai4binance/cli/runtime.py @@ -317,6 +317,10 @@ def _item(item: object) -> OpportunityReviewItem: def build_read_only_runtime(settings: Settings) -> ReadOnlyRuntimeCycle: """Build the resident wallet-first runtime without any write endpoint.""" public_acquisition = build_public_acquisition(settings) + runtime_symbol = _canonical_resident_runtime_symbol( + settings, + datetime.now(UTC), + ) context_loader = RuntimeResearchContextLoader( news_feed_path=settings.runtime_news_feed_path, social_feed_path=settings.runtime_social_feed_path, @@ -378,7 +382,7 @@ def build_read_only_runtime(settings: Settings) -> ReadOnlyRuntimeCycle: ), ) return ReadOnlyRuntimeCycle( - symbol=settings.symbol, + symbol=runtime_symbol, timeframes=settings.timeframes, spot_acquirer=RuntimeContextAcquirer(public_acquisition, context_loader), spot_wallet_service=spot_wallet_service, @@ -439,6 +443,27 @@ def build_read_only_runtime(settings: Settings) -> ReadOnlyRuntimeCycle: ) +def _canonical_resident_runtime_symbol( + settings: Settings, + observed_at: datetime, +) -> str: + """Select one deterministic resident symbol from the shared Spot universe.""" + + configured = tuple( + dict.fromkeys( + ( + settings.symbol, + *settings.fixed_symbols, + *settings.priority_watchlist, + ) + ) + ) + ranked = _virtual_market_ranked_symbols(settings, configured, observed_at) + if settings.symbol in ranked: + return settings.symbol + return ranked[0] if ranked else settings.symbol + + def run_runtime_command( command: str, settings: Settings, @@ -567,11 +592,15 @@ def run_virtual_market_daemon( from ai4binance.data.market_history_sync import ( read_cached_market_universe, ) + from ai4binance.integrations.research_market_universe import ( + RESEARCH_MARKET_UNIVERSE_SOURCE, + ) universe = read_cached_market_universe( settings.market_history_source_cache_directory / "universe-v3.json", clock(), + expected_source=RESEARCH_MARKET_UNIVERSE_SOURCE, ) configured = tuple( dict.fromkeys( @@ -986,69 +1015,140 @@ def _run_virtual_loss_tuning( ) -> dict[str, object]: """Run canonical Spot backtest/tuning after three same-day virtual losses.""" - safe_state = { + trigger, early_result = _virtual_loss_tuning_trigger(journal, observed_at) + if early_result is not None: + return early_result + if trigger is None: + raise RuntimeError("virtual loss tuning trigger resolution failed closed") + trigger_id = str(trigger.get("trigger_id", "")) + tuning_root = settings.validation_artifact_directory / "virtual_loss_tuning" + artifact_path = tuning_root / f"{trigger_id.rsplit(':', maxsplit=1)[-1]}.json" + existing_result = _existing_virtual_loss_tuning_result( + artifact_path, + trigger_id=trigger_id, + observed_at=observed_at, + ) + if existing_result is not None: + return existing_result + results, aggregate_blockers = _virtual_loss_tuning_subject_results( + settings, + trigger, + tuning_root, + ) + status = ( + "RESEARCH_TUNING_COMPLETED" + if results + and all(item.get("status") == "RESEARCH_TUNING_COMPLETED" for item in results) + else "RETRY_PENDING" + ) + payload = { + "schema_version": "VirtualLossTuningResult/v1", + "status": status, + "trigger_id": trigger_id, + "trigger": trigger, + "attempted_at": observed_at.astimezone(UTC).isoformat(), + "results": results, + "blockers": list(dict.fromkeys(aggregate_blockers)), + "parameter_application": "NOT_APPLIED", + **_virtual_loss_safe_state(), + } + write_json_object_verified( + artifact_path, + payload, + blocker="VIRTUAL_LOSS_TUNING_WRITE_FAILED", + subject_id=trigger_id, + indent=2, + durable=True, + ) + return {**payload, "artifact_path": str(artifact_path)} + + +def _virtual_loss_safe_state() -> dict[str, object]: + return { "execution_allowed": False, "promotion_status": "RESEARCH_ONLY", "live_eligibility_status": "LIVE_ORDER_BLOCKED", } + + +def _virtual_loss_tuning_trigger( + journal: VirtualWalletJournal, + observed_at: datetime, +) -> tuple[dict[str, object] | None, dict[str, object] | None]: try: trigger = journal.daily_loss_tuning_trigger(observed_at) except (OSError, RuntimeError, TypeError, ValueError): - return { + return None, { "status": "BLOCKED", "blockers": ["VIRTUAL_LOSS_TUNING_TRIGGER_UNAVAILABLE"], "parameter_application": "NOT_APPLIED", - **safe_state, + **_virtual_loss_safe_state(), } if trigger.get("status") != "TRIGGERED": - return {**trigger, "parameter_application": "NOT_APPLIED"} + return None, {**trigger, "parameter_application": "NOT_APPLIED"} trigger_id = str(trigger.get("trigger_id", "")) - if not re.fullmatch(r"virtual-loss-tuning:[0-9a-f]{24}", trigger_id): + if re.fullmatch(r"virtual-loss-tuning:[0-9a-f]{24}", trigger_id): + return trigger, None + return None, { + "status": "BLOCKED", + "blockers": ["VIRTUAL_LOSS_TUNING_TRIGGER_INVALID"], + "parameter_application": "NOT_APPLIED", + **_virtual_loss_safe_state(), + } + + +def _existing_virtual_loss_tuning_result( + artifact_path: Path, + *, + trigger_id: str, + observed_at: datetime, +) -> dict[str, object] | None: + existing = _load_json_mapping(artifact_path) + if not existing: + return None + safe_state = _virtual_loss_safe_state() + if ( + existing.get("trigger_id") != trigger_id + or existing.get("execution_allowed") is not False + or existing.get("promotion_status") != "RESEARCH_ONLY" + or existing.get("live_eligibility_status") != "LIVE_ORDER_BLOCKED" + ): return { "status": "BLOCKED", - "blockers": ["VIRTUAL_LOSS_TUNING_TRIGGER_INVALID"], + "blockers": ["VIRTUAL_LOSS_TUNING_ARTIFACT_INVALID"], "parameter_application": "NOT_APPLIED", **safe_state, } - tuning_root = settings.validation_artifact_directory / "virtual_loss_tuning" - artifact_path = tuning_root / f"{trigger_id.rsplit(':', maxsplit=1)[-1]}.json" - existing = _load_json_mapping(artifact_path) - if existing: - if ( - existing.get("trigger_id") != trigger_id - or existing.get("execution_allowed") is not False - or existing.get("promotion_status") != "RESEARCH_ONLY" - or existing.get("live_eligibility_status") != "LIVE_ORDER_BLOCKED" - ): - return { - "status": "BLOCKED", - "blockers": ["VIRTUAL_LOSS_TUNING_ARTIFACT_INVALID"], - "parameter_application": "NOT_APPLIED", - **safe_state, - } - if existing.get("status") == "RESEARCH_TUNING_COMPLETED": - return { - "status": "ALREADY_REVIEWED", - "trigger_id": trigger_id, - "artifact_path": str(artifact_path), - "parameter_application": "NOT_APPLIED", - **safe_state, - } - attempted_at = existing.get("attempted_at") - try: - previous_attempt = datetime.fromisoformat(str(attempted_at)).astimezone(UTC) - except (TypeError, ValueError): - previous_attempt = observed_at.astimezone(UTC) - timedelta(hours=1) - if observed_at.astimezone(UTC) - previous_attempt < timedelta(minutes=15): - return { - "status": "RETRY_PENDING", - "trigger_id": trigger_id, - "artifact_path": str(artifact_path), - "blockers": existing.get("blockers", []), - "parameter_application": "NOT_APPLIED", - **safe_state, - } + if existing.get("status") == "RESEARCH_TUNING_COMPLETED": + return { + "status": "ALREADY_REVIEWED", + "trigger_id": trigger_id, + "artifact_path": str(artifact_path), + "parameter_application": "NOT_APPLIED", + **safe_state, + } + attempted_at = existing.get("attempted_at") + try: + previous_attempt = datetime.fromisoformat(str(attempted_at)).astimezone(UTC) + except (TypeError, ValueError): + previous_attempt = observed_at.astimezone(UTC) - timedelta(hours=1) + if observed_at.astimezone(UTC) - previous_attempt >= timedelta(minutes=15): + return None + return { + "status": "RETRY_PENDING", + "trigger_id": trigger_id, + "artifact_path": str(artifact_path), + "blockers": existing.get("blockers", []), + "parameter_application": "NOT_APPLIED", + **safe_state, + } + +def _virtual_loss_tuning_subject_results( + settings: Settings, + trigger: Mapping[str, object], + tuning_root: Path, +) -> tuple[list[dict[str, object]], list[str]]: from ai4binance.application.validation_pipeline import VALIDATED_PLAYBOOKS from ai4binance.data import DatasetIntegrityError, ParquetOHLCVArchive from ai4binance.validation_pipeline_runtime import ( @@ -1065,7 +1165,7 @@ def _run_virtual_loss_tuning( ) runtime = HistoricalValidationRuntime( report_directory=settings.backtest_report_directory / "virtual_loss_tuning", - position_notional_to_equity_ratio=(SPOT_VALIDATION_NOTIONAL_TO_EQUITY_RATIO), + position_notional_to_equity_ratio=SPOT_VALIDATION_NOTIONAL_TO_EQUITY_RATIO, ) results: list[dict[str, object]] = [] aggregate_blockers: list[str] = [] @@ -1073,29 +1173,23 @@ def _run_virtual_loss_tuning( if not isinstance(raw_subject, Mapping): aggregate_blockers.append("VIRTUAL_LOSS_TUNING_SUBJECT_INVALID") continue - market = str(raw_subject.get("market", "")).upper() - symbol = str(raw_subject.get("symbol", "")).upper() - timeframe = str(raw_subject.get("timeframe", "")) - playbook = str(raw_subject.get("strategy_id", "")) subject = { - "market": market, - "symbol": symbol, - "timeframe": timeframe, - "strategy_id": playbook, + "market": str(raw_subject.get("market", "")).upper(), + "symbol": str(raw_subject.get("symbol", "")).upper(), + "timeframe": str(raw_subject.get("timeframe", "")), + "strategy_id": str(raw_subject.get("strategy_id", "")), } - if market != "SPOT" or playbook not in VALIDATED_PLAYBOOKS: - blocker = "VIRTUAL_LOSS_TUNING_SUBJECT_UNSUPPORTED" - aggregate_blockers.append(blocker) - results.append( - {"subject": subject, "status": "BLOCKED", "blockers": [blocker]} - ) - continue try: - candles = archive.read(symbol, timeframe) + if ( + subject["market"] != "SPOT" + or subject["strategy_id"] not in VALIDATED_PLAYBOOKS + ): + raise LookupError("unsupported virtual loss tuning subject") + candles = archive.read(str(subject["symbol"]), str(subject["timeframe"])) result = runtime.validate_one( - symbol, - timeframe, - playbook, + str(subject["symbol"]), + str(subject["timeframe"]), + str(subject["strategy_id"]), candles, artifact_directory=tuning_root / "evidence", ) @@ -1103,19 +1197,11 @@ def _run_virtual_loss_tuning( backtest = cast(Any, result.backtest) if tuning is None or backtest is None: raise ValueError("VIRTUAL_LOSS_TUNING_RESULT_INCOMPLETE") + except LookupError: + blocker = "VIRTUAL_LOSS_TUNING_SUBJECT_UNSUPPORTED" + aggregate_blockers.append(blocker) results.append( - { - "subject": subject, - "status": "RESEARCH_TUNING_COMPLETED", - "dataset_sha256": runtime.dataset_sha256(candles), - "backtest_metrics": to_primitive(backtest.metrics), - "tuning_report_id": tuning.report_id, - "candidate_count": tuning.search_space.candidate_count, - "selected_parameters": to_primitive(tuning.selected_parameters), - "tuning_blockers": list(tuning.blockers), - "validation_blockers": list(result.blockers), - "parameter_application": "NOT_APPLIED", - } + {"subject": subject, "status": "BLOCKED", "blockers": [blocker]} ) except ( DatasetIntegrityError, @@ -1129,34 +1215,24 @@ def _run_virtual_loss_tuning( results.append( {"subject": subject, "status": "BLOCKED", "blockers": [blocker]} ) + else: + results.append( + { + "subject": subject, + "status": "RESEARCH_TUNING_COMPLETED", + "dataset_sha256": runtime.dataset_sha256(candles), + "backtest_metrics": to_primitive(backtest.metrics), + "tuning_report_id": tuning.report_id, + "candidate_count": tuning.search_space.candidate_count, + "selected_parameters": to_primitive(tuning.selected_parameters), + "tuning_blockers": list(tuning.blockers), + "validation_blockers": list(result.blockers), + "parameter_application": "NOT_APPLIED", + } + ) if not subjects: aggregate_blockers.append("VIRTUAL_LOSS_TUNING_SUBJECTS_MISSING") - status = ( - "RESEARCH_TUNING_COMPLETED" - if results - and all(item.get("status") == "RESEARCH_TUNING_COMPLETED" for item in results) - else "RETRY_PENDING" - ) - payload = { - "schema_version": "VirtualLossTuningResult/v1", - "status": status, - "trigger_id": trigger_id, - "trigger": trigger, - "attempted_at": observed_at.astimezone(UTC).isoformat(), - "results": results, - "blockers": list(dict.fromkeys(aggregate_blockers)), - "parameter_application": "NOT_APPLIED", - **safe_state, - } - write_json_object_verified( - artifact_path, - payload, - blocker="VIRTUAL_LOSS_TUNING_WRITE_FAILED", - subject_id=trigger_id, - indent=2, - durable=True, - ) - return {**payload, "artifact_path": str(artifact_path)} + return results, aggregate_blockers def _virtual_wallet_report_root(state_path: Path) -> Path: @@ -1919,32 +1995,80 @@ def _validate_runtime_research_traceability( opportunity_payload = _load_json_mapping(opportunity_path) if opportunity_payload is None: - blockers.append("RUNTIME_OPPORTUNITY_REPORT_UNAVAILABLE") - report = { - "report_id": "runtime-research-trace-validation:unavailable", - "generated_at": generated_at.isoformat(), - "status": "BLOCKED", - "opportunity_report_path": str(opportunity_path), - "report_path": str(report_path), - "opportunity_count": 0, - "validated_count": 0, - "mismatch_count": 0, - "all_opportunities_traceable": False, - "mismatches": (), - "blockers": tuple(blockers), - "execution_allowed": False, - "live_eligibility_status": "LIVE_ORDER_BLOCKED", - } + report = _unavailable_runtime_trace_report( + generated_at=generated_at, + opportunity_path=opportunity_path, + report_path=report_path, + ) _persist_trace_report(report_path, report) return report, 2 + feed_paths, feed_indices = _runtime_trace_sources(settings) + opportunities = _mapping_sequence(opportunity_payload.get("opportunities")) + mismatch_entries, validated_count = _runtime_trace_mismatches( + opportunities, + feed_indices, + ) + all_traceable = not mismatch_entries + report_id = _runtime_trace_report_id(opportunity_payload) + report = { + "report_id": report_id, + "generated_at": generated_at.isoformat(), + "status": "PASS" if all_traceable else "BLOCKED", + "opportunity_report_path": str(opportunity_path), + "report_path": str(report_path), + "feed_paths": {name: str(path) for name, path in feed_paths.items()}, + "opportunity_count": len(opportunities), + "validated_count": validated_count, + "mismatch_count": len(mismatch_entries), + "all_opportunities_traceable": all_traceable, + "mismatches": tuple(mismatch_entries), + "blockers": tuple(blockers), + "execution_allowed": False, + "live_eligibility_status": "LIVE_ORDER_BLOCKED", + } + + persist_blocker = _persist_trace_report(report_path, report) + if persist_blocker is not None: + blockers.append(persist_blocker) + report["status"] = "BLOCKED" + report["blockers"] = tuple(dict.fromkeys(blockers)) + return report, 0 if report["status"] == "PASS" else 2 + + +def _unavailable_runtime_trace_report( + *, + generated_at: datetime, + opportunity_path: Path, + report_path: Path, +) -> dict[str, object]: + return { + "report_id": "runtime-research-trace-validation:unavailable", + "generated_at": generated_at.isoformat(), + "status": "BLOCKED", + "opportunity_report_path": str(opportunity_path), + "report_path": str(report_path), + "opportunity_count": 0, + "validated_count": 0, + "mismatch_count": 0, + "all_opportunities_traceable": False, + "mismatches": (), + "blockers": ("RUNTIME_OPPORTUNITY_REPORT_UNAVAILABLE",), + "execution_allowed": False, + "live_eligibility_status": "LIVE_ORDER_BLOCKED", + } + + +def _runtime_trace_sources( + settings: Settings, +) -> tuple[dict[str, Path], dict[str, dict[str, str]]]: feed_paths = { "news": _resolve_path(settings.runtime_news_feed_path), "social": _resolve_path(settings.runtime_social_feed_path), "content": _resolve_path(settings.runtime_content_feed_path), "technology": _resolve_path(settings.runtime_technology_feed_path), } - feed_indices = { + return feed_paths, { "news": _jsonl_source_index(feed_paths["news"], ("event_id",)), "social": _jsonl_source_index(feed_paths["social"], ("event_id",)), "content": _jsonl_source_index(feed_paths["content"], ("event_id",)), @@ -1954,54 +2078,20 @@ def _validate_runtime_research_traceability( ), } - opportunities = _mapping_sequence(opportunity_payload.get("opportunities")) - mismatch_entries: list[dict[str, object]] = [] - validated_count = 0 +def _runtime_trace_mismatches( + opportunities: tuple[Mapping[str, object], ...], + feed_indices: Mapping[str, Mapping[str, str]], +) -> tuple[list[dict[str, object]], int]: + mismatches: list[dict[str, object]] = [] + validated_count = 0 for index, opportunity in enumerate(opportunities): opportunity_id = _required_text(opportunity.get("opportunity_id")) if opportunity_id is None: opportunity_id = f"opportunity-{index + 1}" - errors: list[str] = [] - - primary_source_url = _required_text(opportunity.get("primary_source_url")) - trace_urls = _text_sequence(opportunity.get("trace_urls")) - trace_evidence = _mapping_sequence(opportunity.get("trace_evidence")) - if primary_source_url is None or not _is_allowed_url(primary_source_url): - errors.append("PRIMARY_SOURCE_URL_MISSING_OR_INVALID") - if not trace_urls: - errors.append("TRACE_URLS_MISSING") - if not trace_evidence: - errors.append("TRACE_EVIDENCE_MISSING") - if ( - primary_source_url is not None - and trace_urls - and primary_source_url not in trace_urls - ): - errors.append("PRIMARY_SOURCE_NOT_IN_TRACE_URLS") - - for trace_entry in trace_evidence: - feed = _required_text(trace_entry.get("feed")) - evidence_id = _required_text(trace_entry.get("evidence_id")) - source_url = _required_text(trace_entry.get("source_url")) - if feed is None or evidence_id is None or source_url is None: - errors.append("TRACE_EVIDENCE_SHAPE_INVALID") - continue - index_payload = feed_indices.get(feed) - if index_payload is None: - errors.append(f"TRACE_EVIDENCE_FEED_UNKNOWN:{feed}") - continue - expected_url = index_payload.get(evidence_id) - if expected_url is None: - errors.append(f"TRACE_EVIDENCE_ID_NOT_FOUND:{feed}:{evidence_id}") - continue - if expected_url != source_url: - errors.append(f"TRACE_EVIDENCE_URL_MISMATCH:{feed}:{evidence_id}") - if source_url not in trace_urls: - errors.append(f"TRACE_EVIDENCE_URL_NOT_LISTED:{feed}:{evidence_id}") - + errors = _runtime_trace_errors(opportunity, feed_indices) if errors: - mismatch_entries.append( + mismatches.append( { "opportunity_id": opportunity_id, "errors": tuple(dict.fromkeys(errors)), @@ -2009,37 +2099,61 @@ def _validate_runtime_research_traceability( ) else: validated_count += 1 - - all_traceable = not mismatch_entries + return mismatches, validated_count + + +def _runtime_trace_errors( + opportunity: Mapping[str, object], + feed_indices: Mapping[str, Mapping[str, str]], +) -> list[str]: + errors: list[str] = [] + primary_source_url = _required_text(opportunity.get("primary_source_url")) + trace_urls = _text_sequence(opportunity.get("trace_urls")) + trace_evidence = _mapping_sequence(opportunity.get("trace_evidence")) + if primary_source_url is None or not _is_allowed_url(primary_source_url): + errors.append("PRIMARY_SOURCE_URL_MISSING_OR_INVALID") + if not trace_urls: + errors.append("TRACE_URLS_MISSING") + if not trace_evidence: + errors.append("TRACE_EVIDENCE_MISSING") + if primary_source_url is not None and primary_source_url not in trace_urls: + errors.append("PRIMARY_SOURCE_NOT_IN_TRACE_URLS") + for trace_entry in trace_evidence: + errors.extend( + _runtime_trace_entry_errors(trace_entry, trace_urls, feed_indices) + ) + return errors + + +def _runtime_trace_entry_errors( + trace_entry: Mapping[str, object], + trace_urls: tuple[str, ...], + feed_indices: Mapping[str, Mapping[str, str]], +) -> list[str]: + feed = _required_text(trace_entry.get("feed")) + evidence_id = _required_text(trace_entry.get("evidence_id")) + source_url = _required_text(trace_entry.get("source_url")) + if feed is None or evidence_id is None or source_url is None: + return ["TRACE_EVIDENCE_SHAPE_INVALID"] + index_payload = feed_indices.get(feed) + if index_payload is None: + return [f"TRACE_EVIDENCE_FEED_UNKNOWN:{feed}"] + expected_url = index_payload.get(evidence_id) + if expected_url is None: + return [f"TRACE_EVIDENCE_ID_NOT_FOUND:{feed}:{evidence_id}"] + errors: list[str] = [] + if expected_url != source_url: + errors.append(f"TRACE_EVIDENCE_URL_MISMATCH:{feed}:{evidence_id}") + if source_url not in trace_urls: + errors.append(f"TRACE_EVIDENCE_URL_NOT_LISTED:{feed}:{evidence_id}") + return errors + + +def _runtime_trace_report_id(opportunity_payload: Mapping[str, object]) -> str: report_id_source = _required_text(opportunity_payload.get("snapshot_id")) - report_id = ( - f"runtime-research-trace-validation:{report_id_source}" - if report_id_source is not None - else "runtime-research-trace-validation:unknown" - ) - report = { - "report_id": report_id, - "generated_at": generated_at.isoformat(), - "status": "PASS" if all_traceable else "BLOCKED", - "opportunity_report_path": str(opportunity_path), - "report_path": str(report_path), - "feed_paths": {name: str(path) for name, path in feed_paths.items()}, - "opportunity_count": len(opportunities), - "validated_count": validated_count, - "mismatch_count": len(mismatch_entries), - "all_opportunities_traceable": all_traceable, - "mismatches": tuple(mismatch_entries), - "blockers": tuple(blockers), - "execution_allowed": False, - "live_eligibility_status": "LIVE_ORDER_BLOCKED", - } - - persist_blocker = _persist_trace_report(report_path, report) - if persist_blocker is not None: - blockers.append(persist_blocker) - report["status"] = "BLOCKED" - report["blockers"] = tuple(dict.fromkeys(blockers)) - return report, 0 if report["status"] == "PASS" else 2 + if report_id_source is None: + return "runtime-research-trace-validation:unknown" + return f"runtime-research-trace-validation:{report_id_source}" def _persist_trace_report(path: Path, payload: dict[str, object]) -> str | None: diff --git a/src/ai4binance/compatibility/opportunity_monitor.py b/src/ai4binance/compatibility/opportunity_monitor.py index 9d3bfc02..77d2e448 100644 --- a/src/ai4binance/compatibility/opportunity_monitor.py +++ b/src/ai4binance/compatibility/opportunity_monitor.py @@ -30,6 +30,9 @@ has_complete_measurable_opportunity, has_complete_measurable_trade_plan, ) +from ai4binance.integrations.research_market_universe import ( + RESEARCH_MARKET_UNIVERSE_SOURCE, +) from ai4binance.opportunity_intelligence import ( TIMEFRAME_DURATIONS, ChartPatternLifecycleState, @@ -72,7 +75,11 @@ def market_symbols(cache: Path, market: str, now: datetime) -> tuple[str, ...]: """Return only the currently verified Binance market universe for the UI.""" if market not in MARKET_TIMEFRAMES: raise ValueError("monitor market is invalid") - universe = read_cached_market_universe(cache / "universe-v3.json", now) + universe = read_cached_market_universe( + cache / "universe-v3.json", + now, + expected_source=RESEARCH_MARKET_UNIVERSE_SOURCE, + ) if universe is None: return () return universe.spot_symbols if market == "SPOT" else universe.futures_symbols diff --git a/src/ai4binance/config.py b/src/ai4binance/config.py index 08da6464..5ae3d287 100644 --- a/src/ai4binance/config.py +++ b/src/ai4binance/config.py @@ -51,14 +51,17 @@ class Settings(BaseSettings): market_history_state_path: Path = Path("runtime/state/market-history-latest.json") market_history_interval_seconds: float = 21_600.0 market_history_live_interval_seconds: float = 300.0 - # Keep enough native history to satisfy the governed 365-day Spot OOS - # observation floor, with a bounded buffer for publication lag and gaps. - market_history_initial_days: int = 400 + # Baseline collection is narrower than the governed 365-day OOS promotion + # floor. Insufficient OOS evidence remains research-only and blocked. + market_history_initial_days: int = 90 + market_history_enrichment_days: int = 30 market_history_pages_per_stream: int = 32 market_history_max_workers: int = 8 market_history_opportunity_workers: int = 2 market_history_local_candles: bool = True - market_history_coin_m_enabled: bool = True + market_history_coin_m_enabled: bool = False + market_history_wallet_minimum_value_usdt: Decimal = Decimal("1") + market_history_market_cap_asset_limit: int = 20 market_depth_enabled: bool = True public_rate_limit_soft: float = 0.70 public_rate_limit_warning: float = 0.80 @@ -423,6 +426,31 @@ def validate_market_history_initial_days(cls, value: int) -> int: raise ValueError("market history initial days must be between 1 and 3650") return value + @field_validator("market_history_enrichment_days") + @classmethod + def validate_market_history_enrichment_days(cls, value: int) -> int: + if not 1 <= value <= 3650: + raise ValueError( + "market history enrichment days must be between 1 and 3650" + ) + return value + + @field_validator("market_history_market_cap_asset_limit") + @classmethod + def validate_market_history_market_cap_asset_limit(cls, value: int) -> int: + if not 1 <= value <= 50: + raise ValueError("market history market-cap limit must be between 1 and 50") + return value + + @field_validator("market_history_wallet_minimum_value_usdt") + @classmethod + def validate_market_history_wallet_minimum_value_usdt( + cls, value: Decimal + ) -> Decimal: + if not value.is_finite() or value <= Decimal("0"): + raise ValueError("market history wallet minimum value must be positive") + return value + @field_validator("market_history_pages_per_stream") @classmethod def validate_market_history_pages(cls, value: int) -> int: diff --git a/src/ai4binance/data/market_depth.py b/src/ai4binance/data/market_depth.py index 0de54fbc..9d63e246 100644 --- a/src/ai4binance/data/market_depth.py +++ b/src/ai4binance/data/market_depth.py @@ -176,6 +176,44 @@ def append(self, records: list[DepthRecord]) -> None: ), ) + def retain_streams(self, allowed: Mapping[str, tuple[str, ...]]) -> dict[str, int]: + """Remove journal rows that are outside the current research universe.""" + allowed_pairs: set[tuple[str, str]] = set() + for market, symbols in allowed.items(): + if market not in _URLS: + raise ValueError("invalid depth retention market") + for symbol in symbols: + normalized = symbol.strip().upper() + if normalized != symbol or not normalized.isalnum(): + raise ValueError("invalid depth retention symbol") + allowed_pairs.add((market, normalized)) + + with self.lock, self.connection: + observed = { + (str(row[0]), str(row[1])) + for row in self.connection.execute( + "SELECT market,symbol FROM depth_heads UNION " + "SELECT DISTINCT market,symbol FROM depth_events" + ).fetchall() + } + removed_streams = observed - allowed_pairs + removed_events = 0 + removed_heads = 0 + for market, symbol in sorted(removed_streams): + removed_events += self.connection.execute( + "DELETE FROM depth_events WHERE market=? AND symbol=?", + (market, symbol), + ).rowcount + removed_heads += self.connection.execute( + "DELETE FROM depth_heads WHERE market=? AND symbol=?", + (market, symbol), + ).rowcount + return { + "removed_stream_count": len(removed_streams), + "removed_event_count": removed_events, + "removed_head_count": removed_heads, + } + def close(self) -> None: with self.lock: self.connection.close() @@ -284,6 +322,7 @@ def start(self, symbols: Mapping[str, tuple[str, ...]]) -> None: self.stop_event = Event() self.symbols = dict(symbols) self.journal = DepthJournal(self.root / "depth.sqlite3") + retention = self.journal.retain_streams(self.symbols) for market, names in self.symbols.items(): if market not in self.transports: raise ValueError("depth market transport is missing") @@ -311,6 +350,7 @@ def start(self, symbols: Mapping[str, tuple[str, ...]]) -> None: "groups": self.groups, "coverage": "PUBLIC_L2_SNAPSHOT_LIMITED_PLUS_DIFFS", "offline_depth_backfill": "UNAVAILABLE", + "retention": retention, "blockers": ([] if self.threads else ["NO_ELIGIBLE_DEPTH_TARGETS"]), **_SAFE_STATE, }, diff --git a/src/ai4binance/data/market_history_continuous.py b/src/ai4binance/data/market_history_continuous.py index 98c4f9a1..671d1eab 100644 --- a/src/ai4binance/data/market_history_continuous.py +++ b/src/ai4binance/data/market_history_continuous.py @@ -29,6 +29,7 @@ MarketHistorySynchronizer, _csv_rows, ) +from ai4binance.data.market_universe_retention import MarketUniverseRetention from ai4binance.data.timeframes import timeframe_duration from ai4binance.exchange.client import JsonTransport from ai4binance.exchange.rate_limit import ( @@ -37,10 +38,10 @@ public_request_weight, ) from ai4binance.infrastructure.persistence.safe_json import ( - DestinationVerificationError, - write_json_object_verified, + DestinationVerificationError as LegacyDestinationVerificationError, ) from ai4binance.schemas import OHLCVCandle +from ai4binance.storage import DestinationVerificationError, write_json_object_verified _DEFAULT_KLINE_INTERVAL = timedelta(minutes=5) _TIMEFRAME_REFRESH_SECONDS: dict[str, int] = { @@ -263,7 +264,9 @@ def _save(path: Path, payload: Mapping[str, object]) -> None: def _recoverable_error_code(error: Exception) -> str: """Expose only repository-owned verification codes, never exception detail.""" - if isinstance(error, DestinationVerificationError): + if isinstance( + error, (DestinationVerificationError, LegacyDestinationVerificationError) + ): code = str(error).strip() if re.fullmatch(r"[A-Z][A-Z0-9_]{2,127}", code): return code @@ -474,6 +477,7 @@ class ContinuousMarketHistory: spot: JsonTransport futures: JsonTransport initial_days: int = 201 + enrichment_days: int = 30 pages_per_stream: int = 2 archive_downloaded_bytes: int = 0 archive_request_count: int = 0 @@ -488,9 +492,14 @@ class ContinuousMarketHistory: on_symbol_ready: SymbolReadyHandler | None = field(default=None, repr=False) on_symbol_screen: SymbolReadyHandler | None = field(default=None, repr=False) clock: Callable[[], datetime] = field(default=lambda: datetime.now(UTC), repr=False) + retention: MarketUniverseRetention | None = field(default=None, repr=False) def __post_init__(self) -> None: - if not 1 <= self.initial_days <= 3650 or not 1 <= self.pages_per_stream <= 32: + if ( + not 1 <= self.initial_days <= 3650 + or not 1 <= self.enrichment_days <= 3650 + or not 1 <= self.pages_per_stream <= 32 + ): raise ValueError("continuous collection bounds are invalid") if not 1 <= self.max_workers <= 8: raise ValueError("collection workers must be between 1 and 8") @@ -518,6 +527,9 @@ def sync_cycle(self, *, observed_at: datetime) -> dict[str, object]: now = observed_at.astimezone(UTC) before = self._network_totals() universe = self.history._eligible_universe(now, force_refresh=True) + retention_result = None + if self.retention is not None and not universe.blockers: + retention_result = self.retention.prune(universe) self.coin_m_contracts = { symbol: (pair, contract) for symbol, pair, contract in universe.coin_m_contracts @@ -1046,6 +1058,7 @@ def publish_progress( publish_progress(None, force=True) if not universe.blockers: with ThreadPoolExecutor(max_workers=self.max_workers) as executor: + enrichment_stream_keys: set[tuple[str, str, str]] = set() def collect_stream( market: str, @@ -1054,14 +1067,29 @@ def collect_stream( kind: str, timeframe: str | None, ) -> dict[str, object]: - result = self._collect_stream( + if ( market, symbol, - transport, - now, - kind=kind, - timeframe=timeframe, - ) + str(timeframe), + ) in enrichment_stream_keys: + result = self._collect_stream( + market, + symbol, + transport, + now, + kind=kind, + timeframe=timeframe, + history_days=self.enrichment_days, + ) + else: + result = self._collect_stream( + market, + symbol, + transport, + now, + kind=kind, + timeframe=timeframe, + ) record_coverage( result, market=market, @@ -1172,6 +1200,7 @@ def enqueue_enrichment(identity: tuple[str, str]) -> None: for tf in _ENRICHMENT_TIMEFRAMES: if (*identity, "klines", tf) not in queued_keys: background.append((*identity, transport, "klines", tf)) + enrichment_stream_keys.add((*identity, tf)) label = market_labels[identity[0]] with coverage_lock: coverage_expected[label][tf] += 1 @@ -1410,10 +1439,15 @@ def complete_active_request() -> None: ], "collection_plan": self._collection_plan(), "initial_history_days": self.initial_days, + "enrichment_history_days": self.enrichment_days, "spot_universe_count": len(universe.spot_symbols), "futures_universe_count": len(universe.futures_symbols), "coin_m_universe_count": len(universe.coin_m_symbols), "excluded_asset_count": len(universe.excluded_assets), + "selected_assets": list(universe.selected_assets), + "wallet_assets": list(universe.wallet_assets), + "market_cap_assets": list(universe.market_cap_assets), + "universe_retention": retention_result, "archive_root": str(self.history.archive_root), "completed_symbols": completed_symbols, "total_symbols": total_symbols, @@ -1564,7 +1598,7 @@ def _interleaved_stream_work( self, market_work: tuple[tuple[str, str, JsonTransport], ...], ) -> tuple[tuple[str, str, JsonTransport, str, str | None], ...]: - """Finish watched symbols first, then schedule fair universe backfill.""" + """Finish watched symbols first, then keep the screen universe fresh.""" streams: list[tuple[str, str, JsonTransport, str, str | None]] = [] priority_set = frozenset(self.priority_symbols) @@ -1578,9 +1612,25 @@ def _interleaved_stream_work( for market, symbol, transport in priority_work: for kind, timeframe in self._collection_kinds(market): streams.append((market, symbol, transport, kind, timeframe)) - # Keep each background symbol's direct streams contiguous. Workers may - # fetch later symbols concurrently, while the first fully current - # symbol can immediately enter the bounded opportunity-analysis stage. + if self.on_symbol_screen is not None: + # A complete top-50 cycle can exceed the 15m freshness budget. Run + # each screening timeframe across the full background universe + # before moving to the next timeframe, so early symbols do not + # become stale while later symbols are still being collected. + screen_kinds = tuple( + ("klines", timeframe) for timeframe in _SCREEN_TIMEFRAMES + ) + for kind, timeframe in screen_kinds: + for market, symbol, transport in background_work: + if (kind, timeframe) in self._collection_kinds(market): + streams.append((market, symbol, transport, kind, timeframe)) + for market, symbol, transport in background_work: + for kind, timeframe in self._collection_kinds(market): + if (kind, timeframe) not in screen_kinds: + streams.append((market, symbol, transport, kind, timeframe)) + return tuple(streams) + # Non-staged full-history runs retain symbol-contiguous streams so the + # first complete symbol can become available during long backfills. for market, symbol, transport in background_work: for kind, timeframe in self._collection_kinds(market): streams.append((market, symbol, transport, kind, timeframe)) @@ -1637,6 +1687,8 @@ def _write_collection_progress( "schema_version": "2.0", "status": "COLLECTING", "collection_plan": self._collection_plan(), + "initial_history_days": self.initial_days, + "enrichment_history_days": self.enrichment_days, "cycle_started_at": cycle_started_at.isoformat(), "observed_at": datetime.now(UTC).isoformat(), "active_market": active_market, @@ -1743,6 +1795,7 @@ def _collect_stream( *, kind: str, timeframe: str | None, + history_days: int | None = None, ) -> dict[str, object]: """Collect exactly one market-data stream with fail-closed evidence.""" @@ -1761,7 +1814,13 @@ def _collect_stream( } else: result = self._candles( - market, symbol, kind, transport, now, timeframe="5m" + market, + symbol, + kind, + transport, + now, + timeframe="5m", + history_days=history_days, ) elif kind in {"funding", "open_interest"}: result = self._details(symbol, kind, now, market=market) @@ -1769,7 +1828,13 @@ def _collect_stream( if timeframe is None: raise ValueError("candle timeframe is required") result = self._candles( - market, symbol, kind, transport, now, timeframe=timeframe + market, + symbol, + kind, + transport, + now, + timeframe=timeframe, + history_days=history_days, ) except (OSError, ValueError, ArithmeticError, ExchangeError) as error: result = { @@ -2018,14 +2083,15 @@ def _record_archive_source(self, source: object) -> None: self.archive_downloaded_bytes += int(getattr(source, "downloaded_bytes", 0)) def _candle_initial_start( - self, now: datetime, interval: timedelta + self, now: datetime, interval: timedelta, *, history_days: int | None = None ) -> datetime | None: """Keep the configured history horizon and deterministic candle floor.""" interval_seconds = int(interval.total_seconds()) anchor = now.replace(second=0, microsecond=0) configured_anchor = anchor.replace(hour=0, minute=0) - start = configured_anchor - timedelta(days=self.initial_days) + requested_days = self.initial_days if history_days is None else history_days + start = configured_anchor - timedelta(days=requested_days) if self.minimum_candles is not None: # Never narrow a requested historical horizon to the minimum # analysis window. Conversely, daily data retains the existing @@ -2035,6 +2101,135 @@ def _candle_initial_start( aligned_epoch = int(start.timestamp()) // interval_seconds * interval_seconds return datetime.fromtimestamp(aligned_epoch, tz=UTC) + @staticmethod + def _flush_vision_pending( + *, + archive: ParquetOHLCVArchive, + dataset_symbol: str, + timeframe: str, + pending_candles: list[OHLCVCandle], + pending_source_hashes: list[str], + generated_at: datetime, + state: dict[str, object], + progress_path: Path, + next_at: datetime, + replace_conflicts_from_sources: tuple[str, ...], + ) -> None: + if not pending_candles: + return + source_digest = ( + pending_source_hashes[0] + if len(pending_source_hashes) == 1 + else sha256("".join(pending_source_hashes).encode()).hexdigest() + ) + archive.update( + dataset_symbol, + timeframe, + tuple(pending_candles), + source=f"BINANCE_VISION_DIRECT_{timeframe.upper()}_SHA256:{source_digest}", + generated_at=generated_at, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) + state["next_at"] = next_at.isoformat() + _save(progress_path, state) + pending_candles.clear() + pending_source_hashes.clear() + + @staticmethod + def _vision_first_available(state: Mapping[str, object]) -> datetime | None: + raw_first_available = state.get("first_available_at") + known_first_available = ( + datetime.fromisoformat(str(raw_first_available)) + if raw_first_available + else None + ) + if ( + known_first_available is not None + and known_first_available.utcoffset() is None + ): + raise ValueError("collection first available boundary is invalid") + return known_first_available + + def _handle_missing_vision_archive( + self, + *, + archive: ParquetOHLCVArchive, + dataset_symbol: str, + timeframe: str, + pending_candles: list[OHLCVCandle], + pending_source_hashes: list[str], + generated_at: datetime, + state: dict[str, object], + progress_path: Path, + cursor: datetime, + archive_end: datetime, + replace_conflicts_from_sources: tuple[str, ...], + ) -> tuple[datetime, dict[str, object] | None, bool]: + known_first_available = self._vision_first_available(state) + self._flush_vision_pending( + archive=archive, + dataset_symbol=dataset_symbol, + timeframe=timeframe, + pending_candles=pending_candles, + pending_source_hashes=pending_source_hashes, + generated_at=generated_at, + state=state, + progress_path=progress_path, + next_at=cursor, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) + if known_first_available is None or archive_end <= known_first_available: + state["next_at"] = archive_end.isoformat() + _save(progress_path, state) + return archive_end, None, True + return ( + cursor, + { + "status": "UNAVAILABLE", + "next_at": cursor.isoformat(), + "reason": "BINANCE_VISION_ARCHIVE_UNAVAILABLE", + }, + False, + ) + + def _validate_vision_chunk_start( + self, + *, + archive: ParquetOHLCVArchive, + dataset_symbol: str, + timeframe: str, + pending_candles: list[OHLCVCandle], + pending_source_hashes: list[str], + generated_at: datetime, + state: dict[str, object], + progress_path: Path, + cursor: datetime, + first_timestamp: datetime, + replace_conflicts_from_sources: tuple[str, ...], + ) -> tuple[datetime, dict[str, object] | None]: + if first_timestamp <= cursor or not state.get("first_available_at"): + return cursor, None + known_first_available = self._vision_first_available(state) + if first_timestamp == known_first_available: + return first_timestamp, None + self._flush_vision_pending( + archive=archive, + dataset_symbol=dataset_symbol, + timeframe=timeframe, + pending_candles=pending_candles, + pending_source_hashes=pending_source_hashes, + generated_at=generated_at, + state=state, + progress_path=progress_path, + next_at=cursor, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) + return cursor, { + "status": "UNAVAILABLE", + "next_at": cursor.isoformat(), + "reason": "KLINE_GAP", + } + def _vision_history( self, *, @@ -2059,29 +2254,6 @@ def _vision_history( pending_candles: list[OHLCVCandle] = [] pending_source_hashes: list[str] = [] - def flush_pending() -> None: - if not pending_candles: - return - source_digest = ( - pending_source_hashes[0] - if len(pending_source_hashes) == 1 - else sha256("".join(pending_source_hashes).encode()).hexdigest() - ) - archive.update( - dataset_symbol, - timeframe, - tuple(pending_candles), - source=( - f"BINANCE_VISION_DIRECT_{timeframe.upper()}_SHA256:{source_digest}" - ), - generated_at=now, - replace_conflicts_from_sources=replace_conflicts_from_sources, - ) - state["next_at"] = cursor.isoformat() - _save(progress_path, state) - pending_candles.clear() - pending_source_hashes.clear() - while cursor < closed_history_end: month_start = cursor.replace( day=1, hour=0, minute=0, second=0, microsecond=0 @@ -2121,34 +2293,24 @@ def flush_pending() -> None: # also applies when a retained dataset already proved a later # first-available boundary and the requested history horizon # is subsequently extended backwards. - raw_first_available = state.get("first_available_at") - known_first_available = ( - datetime.fromisoformat(str(raw_first_available)) - if raw_first_available - else None + cursor, unavailable, should_continue = ( + self._handle_missing_vision_archive( + archive=archive, + dataset_symbol=dataset_symbol, + timeframe=timeframe, + pending_candles=pending_candles, + pending_source_hashes=pending_source_hashes, + generated_at=now, + state=state, + progress_path=progress_path, + cursor=cursor, + archive_end=archive_end, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) ) - if ( - known_first_available is not None - and known_first_available.utcoffset() is None - ): - raise ValueError( - "collection first available boundary is invalid" - ) from None - if ( - known_first_available is None - or archive_end <= known_first_available - ): - flush_pending() - cursor = archive_end - state["next_at"] = cursor.isoformat() - _save(progress_path, state) + if should_continue: continue - flush_pending() - return cursor, { - "status": "UNAVAILABLE", - "next_at": cursor.isoformat(), - "reason": "BINANCE_VISION_ARCHIVE_UNAVAILABLE", - } + return cursor, unavailable except ValueError as error: if ( timeframe in _DASHBOARD_REFRESH_TIMEFRAMES @@ -2156,35 +2318,68 @@ def flush_pending() -> None: ): break raise - if candles[0].timestamp > cursor and state.get("first_available_at"): - known_first_available = datetime.fromisoformat( - str(state["first_available_at"]) - ) - if known_first_available.utcoffset() is None: - raise ValueError("collection first available boundary is invalid") - if candles[0].timestamp == known_first_available: - cursor = known_first_available - else: - flush_pending() - return cursor, { - "status": "UNAVAILABLE", - "next_at": cursor.isoformat(), - "reason": "KLINE_GAP", - } + cursor, gap = self._validate_vision_chunk_start( + archive=archive, + dataset_symbol=dataset_symbol, + timeframe=timeframe, + pending_candles=pending_candles, + pending_source_hashes=pending_source_hashes, + generated_at=now, + state=state, + progress_path=progress_path, + cursor=cursor, + first_timestamp=candles[0].timestamp, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) + if gap is not None: + return cursor, gap state.setdefault("first_available_at", candles[0].timestamp.isoformat()) pending_candles.extend(candles) pending_source_hashes.append(str(source.sha256)) cursor = candles[-1].timestamp + interval if cursor < archive_end: - flush_pending() + self._flush_vision_pending( + archive=archive, + dataset_symbol=dataset_symbol, + timeframe=timeframe, + pending_candles=pending_candles, + pending_source_hashes=pending_source_hashes, + generated_at=now, + state=state, + progress_path=progress_path, + next_at=cursor, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) return cursor, { "status": "UNAVAILABLE", "next_at": cursor.isoformat(), "reason": "KLINE_GAP", } if len(pending_source_hashes) >= _VISION_ARCHIVE_BATCH_SIZE: - flush_pending() - flush_pending() + self._flush_vision_pending( + archive=archive, + dataset_symbol=dataset_symbol, + timeframe=timeframe, + pending_candles=pending_candles, + pending_source_hashes=pending_source_hashes, + generated_at=now, + state=state, + progress_path=progress_path, + next_at=cursor, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) + self._flush_vision_pending( + archive=archive, + dataset_symbol=dataset_symbol, + timeframe=timeframe, + pending_candles=pending_candles, + pending_source_hashes=pending_source_hashes, + generated_at=now, + state=state, + progress_path=progress_path, + next_at=cursor, + replace_conflicts_from_sources=replace_conflicts_from_sources, + ) return cursor, None def _candles( @@ -2196,6 +2391,7 @@ def _candles( now: datetime, *, timeframe: str = "5m", + history_days: int | None = None, ) -> dict[str, object]: if timeframe not in MARKET_HISTORY_TIMEFRAMES: raise ValueError("candle timeframe is invalid") @@ -2222,7 +2418,9 @@ def _candles( progress_directory, now, extend_history=True, - initial_start=self._candle_initial_start(now, interval), + initial_start=self._candle_initial_start( + now, interval, history_days=history_days + ), maximum_cursor=end, cursor_tolerance=interval, ) diff --git a/src/ai4binance/data/market_history_sync.py b/src/ai4binance/data/market_history_sync.py index ef50ea15..0e13ce6d 100644 --- a/src/ai4binance/data/market_history_sync.py +++ b/src/ai4binance/data/market_history_sync.py @@ -58,6 +58,7 @@ def read_cached_market_universe( observed_at: datetime, *, max_age: timedelta = _UNIVERSE_CACHE_MAX_AGE, + expected_source: str | None = None, ) -> BinanceEligibleMarketSnapshot | None: """Read a bounded, safety-constrained cached market universe. @@ -85,6 +86,9 @@ def read_cached_market_universe( if ( payload.get("execution_allowed") is not False or payload.get("live_eligibility_status") != "LIVE_ORDER_BLOCKED" + or ( + expected_source is not None and payload.get("source") != expected_source + ) ): return None return BinanceEligibleMarketSnapshot( @@ -97,6 +101,10 @@ def read_cached_market_universe( (str(asset), tuple(reasons)) for asset, reasons in payload["excluded_assets"] ), + selected_assets=tuple(payload.get("selected_assets", ())), + wallet_assets=tuple(payload.get("wallet_assets", ())), + market_cap_assets=tuple(payload.get("market_cap_assets", ())), + source=str(payload.get("source", "BINANCE_PUBLIC_EXCHANGE_INFO")), ) except (AttributeError, KeyError, OSError, TypeError, ValueError): return None @@ -401,27 +409,40 @@ def _eligible_universe( cadence. It prevents a current cache entry from masking a delisting until the next full historical backfill cycle. """ + cache_filename = getattr(self.universe_provider, "cache_filename", None) + cache_source = getattr(self.universe_provider, "cache_source", None) cache_path = self.source_cache.root / ( - "universe-v3.json" + str(cache_filename) + if cache_filename + else "universe-v3.json" if getattr(self.universe_provider, "coin_m_transport", None) else "universe-v2.json" ) if not force_refresh: - cached = read_cached_market_universe(cache_path, observed_at) + cached = read_cached_market_universe( + cache_path, observed_at, expected_source=cache_source + ) if cached is not None: return cached + priority_snapshot = getattr( + self.universe_provider, "priority_eligible_market_snapshot", None + ) top_volume_snapshot = getattr( self.universe_provider, "top_volume_eligible_market_snapshot", None ) - snapshot = ( - top_volume_snapshot( + if callable(priority_snapshot): + snapshot = priority_snapshot() + elif callable(top_volume_snapshot): + snapshot = top_volume_snapshot( max_symbols_per_market=_COLLECTION_MAX_SYMBOLS_PER_MARKET ) - if callable(top_volume_snapshot) - else self.universe_provider.eligible_market_snapshot() - ) + else: + snapshot = self.universe_provider.eligible_market_snapshot() + snapshot = cast(BinanceEligibleMarketSnapshot, snapshot) if snapshot.blockers: - cached = read_cached_market_universe(cache_path, observed_at) + cached = read_cached_market_universe( + cache_path, observed_at, expected_source=cache_source + ) if cached is not None: return cached return snapshot @@ -436,6 +457,10 @@ def _eligible_universe( "futures_symbols": snapshot.futures_symbols, "coin_m_contracts": snapshot.coin_m_contracts, "excluded_assets": snapshot.excluded_assets, + "selected_assets": snapshot.selected_assets, + "wallet_assets": snapshot.wallet_assets, + "market_cap_assets": snapshot.market_cap_assets, + "source": snapshot.source, "execution_allowed": False, "live_eligibility_status": "LIVE_ORDER_BLOCKED", }, diff --git a/src/ai4binance/data/market_universe_retention.py b/src/ai4binance/data/market_universe_retention.py new file mode 100644 index 00000000..343c45a5 --- /dev/null +++ b/src/ai4binance/data/market_universe_retention.py @@ -0,0 +1,157 @@ +"""Bounded retention for the canonical research market universe.""" + +from __future__ import annotations + +import re +import shutil +from dataclasses import dataclass +from pathlib import Path + +from ai4binance.integrations.binance import BinanceEligibleMarketSnapshot + +_SYMBOL = re.compile(r"^[A-Z0-9]{4,24}$") + + +@dataclass(frozen=True, slots=True) +class MarketUniverseRetention: + """Remove only symbol-scoped generated data outside the selected universe.""" + + archive_root: Path + source_cache_root: Path + opportunity_monitor_root: Path + futures_replay_root: Path | None = None + futures_artifact_roots: tuple[Path, ...] = () + + def prune(self, universe: BinanceEligibleMarketSnapshot) -> dict[str, object]: + if universe.blockers: + return { + "status": "BLOCKED", + "removed_directory_count": 0, + "removed_bytes": 0, + "blockers": ["UNIVERSE_RETENTION_SELECTION_UNAVAILABLE"], + } + spot = frozenset(universe.spot_symbols) + futures = frozenset(universe.futures_symbols) + candidates: list[Path] = [] + for name, allowed in ( + ("spot", spot), + ("usd_m_futures", futures), + ("usd_m_futures_mark", futures), + ("usd_m_futures_index", futures), + ): + candidates.extend( + self._out_of_scope_children(self.archive_root / name, allowed) + ) + for name in ( + "coin_m_futures", + "coin_m_futures_mark", + "coin_m_futures_index", + ): + root = self.archive_root / name + if root.is_dir() and not root.is_symlink(): + candidates.append(root) + candidates.extend( + self._out_of_scope_children(self.opportunity_monitor_root / "SPOT", spot) + ) + candidates.extend( + self._out_of_scope_children( + self.opportunity_monitor_root / "USD_M_FUTURES", futures + ) + ) + candidates.extend(self._source_candidates("spot", spot)) + candidates.extend(self._source_candidates("futures/um", futures)) + if self.futures_replay_root is not None: + candidates.extend( + self._out_of_scope_flat_files(self.futures_replay_root, futures) + ) + for root in self.futures_artifact_roots: + candidates.extend(self._out_of_scope_children(root, futures)) + coin_m = self.source_cache_root / "data" / "futures" / "cm" + if coin_m.is_dir() and not coin_m.is_symlink(): + candidates.append(coin_m) + + unique = tuple(sorted(set(candidates), key=lambda path: str(path).casefold())) + removed_bytes = sum(self._directory_bytes(path) for path in unique) + for path in unique: + self._remove_verified_directory(path) + return { + "status": "PRUNED", + "removed_directory_count": len(unique), + "removed_bytes": removed_bytes, + "blockers": [], + } + + def _source_candidates(self, segment: str, allowed: frozenset[str]) -> list[Path]: + root = self.source_cache_root / "data" / Path(segment) + if not root.is_dir() or root.is_symlink(): + return [] + candidates: list[Path] = [] + for cadence in ("daily", "monthly"): + cadence_root = root / cadence + if not cadence_root.is_dir() or cadence_root.is_symlink(): + continue + for kind_root in cadence_root.iterdir(): + if not kind_root.is_dir() or kind_root.is_symlink(): + continue + candidates.extend(self._out_of_scope_children(kind_root, allowed)) + return candidates + + @staticmethod + def _out_of_scope_children(root: Path, allowed: frozenset[str]) -> list[Path]: + if not root.is_dir() or root.is_symlink(): + return [] + return [ + child + for child in root.iterdir() + if child.is_dir() + and not child.is_symlink() + and _SYMBOL.fullmatch(child.name) + and child.name not in allowed + ] + + @staticmethod + def _directory_bytes(path: Path) -> int: + if path.is_file(): + return path.stat().st_size + total = 0 + for child in path.rglob("*"): + if child.is_file() and not child.is_symlink(): + total += child.stat().st_size + return total + + def _remove_verified_directory(self, path: Path) -> None: + resolved = path.resolve() + allowed_roots = tuple( + root.resolve() + for root in ( + self.archive_root, + self.source_cache_root, + self.opportunity_monitor_root, + *((self.futures_replay_root,) if self.futures_replay_root else ()), + *self.futures_artifact_roots, + ) + ) + if path.is_symlink() or not any( + root in resolved.parents for root in allowed_roots + ): + raise ValueError("universe retention target is outside the allowed roots") + if resolved.is_dir(): + shutil.rmtree(resolved) + else: + resolved.unlink() + + @staticmethod + def _out_of_scope_flat_files(root: Path, allowed: frozenset[str]) -> list[Path]: + if not root.is_dir() or root.is_symlink(): + return [] + candidates: list[Path] = [] + for child in root.iterdir(): + if not child.is_file() or child.is_symlink(): + continue + symbol = child.name.split("-", maxsplit=1)[0] + if _SYMBOL.fullmatch(symbol) and symbol not in allowed: + candidates.append(child) + return candidates + + +__all__ = ("MarketUniverseRetention",) diff --git a/src/ai4binance/domain/universe.py b/src/ai4binance/domain/universe.py index c67c2a80..e2280928 100644 --- a/src/ai4binance/domain/universe.py +++ b/src/ai4binance/domain/universe.py @@ -22,6 +22,14 @@ "PYUSD", "USD1", "USDE", + "USDS", + "USDG", + "USDD", + "RLUSD", + "FRAX", + "LUSD", + "GHO", + "EURC", "U", } ) diff --git a/src/ai4binance/integrations/__init__.py b/src/ai4binance/integrations/__init__.py index ece2acc5..5e3e4ddf 100644 --- a/src/ai4binance/integrations/__init__.py +++ b/src/ai4binance/integrations/__init__.py @@ -1 +1,8 @@ """External integration adapters for read-only AI4Binance boundaries.""" + +from ai4binance.integrations.research_market_universe import ( + ReadOnlyCoinGeckoJsonTransport, + ResearchMarketUniverseProvider, +) + +__all__ = ("ReadOnlyCoinGeckoJsonTransport", "ResearchMarketUniverseProvider") diff --git a/src/ai4binance/integrations/binance/market_universe_provider.py b/src/ai4binance/integrations/binance/market_universe_provider.py index f7e3e13a..377e5a87 100644 --- a/src/ai4binance/integrations/binance/market_universe_provider.py +++ b/src/ai4binance/integrations/binance/market_universe_provider.py @@ -72,25 +72,17 @@ class BinanceEligibleMarketSnapshot: execution_allowed: bool = False live_eligibility_status: str = "LIVE_ORDER_BLOCKED" coin_m_contracts: tuple[tuple[str, str, str], ...] = () + selected_assets: tuple[str, ...] = () + wallet_assets: tuple[str, ...] = () + market_cap_assets: tuple[str, ...] = () @property def coin_m_symbols(self) -> tuple[str, ...]: return tuple(item[0] for item in self.coin_m_contracts) def __post_init__(self) -> None: - if self.coin_m_symbols != tuple(sorted(set(self.coin_m_symbols))): - raise ValueError("COIN-M identities must be sorted and unique") - for symbol, pair, contract in self.coin_m_contracts: - if ( - not re.fullmatch(r"[A-Z0-9]{2,24}_(?:PERP|[0-9]{6})", symbol) - or not pair.isalnum() - or not symbol.startswith(pair + "_") - or contract not in {"PERPETUAL", "CURRENT_QUARTER", "NEXT_QUARTER"} - ): - raise ValueError("COIN-M contract identity is invalid") - for symbol in (*self.spot_symbols, *self.futures_symbols): - if not _PUBLIC_SYMBOL_PATTERN.fullmatch(symbol) or symbol != symbol.upper(): - raise ValueError("eligible market symbol identity is invalid") + _validate_coin_m_contracts(self.coin_m_contracts) + _validate_market_symbols(self.spot_symbols, self.futures_symbols) if self.spot_symbols != tuple( sorted(set(self.spot_symbols)) ) or self.futures_symbols != tuple(sorted(set(self.futures_symbols))): @@ -104,6 +96,11 @@ def __post_init__(self) -> None: raise ValueError("eligible market exclusions are invalid") if any(not blocker.strip() for blocker in self.blockers): raise ValueError("eligible market blockers cannot contain blanks") + _validate_market_assets( + self.selected_assets, + self.wallet_assets, + self.market_cap_assets, + ) if ( self.execution_allowed or self.live_eligibility_status != "LIVE_ORDER_BLOCKED" @@ -111,6 +108,38 @@ def __post_init__(self) -> None: raise ValueError("eligible market snapshot cannot grant execution") +def _validate_coin_m_contracts( + contracts: tuple[tuple[str, str, str], ...], +) -> None: + symbols = tuple(item[0] for item in contracts) + if symbols != tuple(sorted(set(symbols))): + raise ValueError("COIN-M identities must be sorted and unique") + for symbol, pair, contract in contracts: + if ( + not re.fullmatch(r"[A-Z0-9]{2,24}_(?:PERP|[0-9]{6})", symbol) + or not pair.isalnum() + or not symbol.startswith(pair + "_") + or contract not in {"PERPETUAL", "CURRENT_QUARTER", "NEXT_QUARTER"} + ): + raise ValueError("COIN-M contract identity is invalid") + + +def _validate_market_symbols( + spot_symbols: tuple[str, ...], futures_symbols: tuple[str, ...] +) -> None: + for symbol in (*spot_symbols, *futures_symbols): + if not _PUBLIC_SYMBOL_PATTERN.fullmatch(symbol) or symbol != symbol.upper(): + raise ValueError("eligible market symbol identity is invalid") + + +def _validate_market_assets(*groups: tuple[str, ...]) -> None: + for assets in groups: + if assets != tuple(dict.fromkeys(assets)) or any( + not asset.isalnum() or asset != asset.upper() for asset in assets + ): + raise ValueError("eligible market asset identities are invalid") + + @dataclass(frozen=True, slots=True) class ReadOnlyBinanceJsonTransport: """HTTPS GET-only transport for public Spot or Futures market endpoints.""" diff --git a/src/ai4binance/integrations/research_market_universe.py b/src/ai4binance/integrations/research_market_universe.py new file mode 100644 index 00000000..989b7972 --- /dev/null +++ b/src/ai4binance/integrations/research_market_universe.py @@ -0,0 +1,439 @@ +"""Research-only wallet and public market-cap universe selection.""" + +from __future__ import annotations + +import json +import math +from collections.abc import Callable, Mapping, Sequence +from dataclasses import dataclass, field +from datetime import UTC, datetime, timedelta +from decimal import Decimal, InvalidOperation +from pathlib import Path +from typing import Protocol +from urllib.error import HTTPError, URLError +from urllib.parse import urlencode +from urllib.request import Request, urlopen + +from ai4binance.core.errors import ( + ExchangeHttpError, + ExchangePayloadError, + ExchangeRateLimitError, + ExchangeTransportError, +) +from ai4binance.domain.universe import classify_asset_eligibility +from ai4binance.integrations.binance.market_universe_provider import ( + BinanceEligibleMarketSnapshot, + BinanceMarketUniverseProvider, +) +from ai4binance.storage import read_bounded_jsonl_tail + +_ZERO = Decimal("0") +RESEARCH_MARKET_UNIVERSE_SOURCE = "BINANCE_WALLET_AND_COINGECKO_MARKET_CAP" + + +class _RetryableMarketCapError(Exception): + def __init__(self, final_error: Exception) -> None: + super().__init__(str(final_error)) + self.final_error = final_error + + +def _read_market_cap_response( + request: Request, *, timeout_seconds: float, maximum_bytes: int +) -> list[object]: + try: + with urlopen( # noqa: S310 # nosec B310 + request, timeout=timeout_seconds + ) as response: + payload = response.read(maximum_bytes + 1) + except HTTPError as error: + if error.code in {418, 429}: + raise ExchangeRateLimitError(error.code, None, request.full_url) from None + failure = ExchangeHttpError( + f"public market-cap HTTP {error.code} at /api/v3/coins/markets" + ) + if error.code < 500: + raise failure from None + raise _RetryableMarketCapError(failure) from None + except (TimeoutError, URLError, OSError): + raise _RetryableMarketCapError( + ExchangeTransportError("public market-cap request failed") + ) from None + if len(payload) > maximum_bytes: + raise ExchangePayloadError("public market-cap payload is too large") + try: + decoded = json.loads(payload.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError): + raise ExchangePayloadError("public market-cap JSON is invalid") from None + if not isinstance(decoded, list): + raise ExchangePayloadError("public market-cap payload is invalid") + return decoded + + +class PublicMarketCapTransport(Protocol): + """Credential-free JSON transport used only for public market ranking.""" + + def get_json( + self, + path: str, + params: Mapping[str, str | int | bool] | None = None, + ) -> object: ... + + +@dataclass(frozen=True, slots=True) +class ReadOnlyCoinGeckoJsonTransport: + """Bounded HTTPS GET transport for the public CoinGecko market endpoint.""" + + base_url: str = "https://api.coingecko.com" + timeout_seconds: float = 10.0 + max_attempts: int = 2 + max_response_bytes: int = 2_000_000 + + def __post_init__(self) -> None: + normalized = self.base_url.strip().rstrip("/") + if normalized != "https://api.coingecko.com": + raise ValueError("CoinGecko public base URL is outside the allowlist") + if self.timeout_seconds <= 0 or not 1 <= self.max_attempts <= 3: + raise ValueError("CoinGecko transport bounds are invalid") + if not 1024 <= self.max_response_bytes <= 8_000_000: + raise ValueError("CoinGecko response bound is invalid") + object.__setattr__(self, "base_url", normalized) + + def get_json( + self, + path: str, + params: Mapping[str, str | int | bool] | None = None, + ) -> object: + if path != "/api/v3/coins/markets" or "?" in path or "#" in path: + raise ValueError("CoinGecko path is outside the configured boundary") + query = urlencode(sorted((params or {}).items())) + request = Request( # noqa: S310 # nosec B310 + f"{self.base_url}{path}?{query}", + headers={"Accept": "application/json", "User-Agent": "AI4Binance/0.1"}, + method="GET", + ) + for attempt in range(1, self.max_attempts + 1): + try: + return _read_market_cap_response( + request, + timeout_seconds=self.timeout_seconds, + maximum_bytes=self.max_response_bytes, + ) + except _RetryableMarketCapError as error: + if attempt == self.max_attempts: + raise error.final_error from None + raise ExchangeTransportError("public market-cap request failed") + + +@dataclass(frozen=True, slots=True) +class WalletAssetSnapshot: + """Secret-safe latest Spot wallet asset selection.""" + + assets: tuple[str, ...] + observed_at: datetime + sync_run_id: str + + +def _wallet_record( + raw: bytes, +) -> tuple[datetime, str, Mapping[str, object]] | None: + value = json.loads(raw.decode("utf-8")) + payload = value.get("payload") if isinstance(value, Mapping) else None + envelope = payload.get("envelope") if isinstance(payload, Mapping) else None + if not isinstance(payload, Mapping) or not isinstance(envelope, Mapping): + return None + sync_run_id = envelope.get("sync_run_id") + event_time = envelope.get("event_time") or envelope.get("received_at") + if not isinstance(sync_run_id, str) or not isinstance(event_time, str): + return None + parsed = datetime.fromisoformat(event_time) + if parsed.tzinfo is None or parsed.utcoffset() is None: + return None + return parsed.astimezone(UTC), sync_run_id, payload + + +def _eligible_wallet_asset( + payload: Mapping[str, object], *, minimum_value_usdt: Decimal +) -> str | None: + if payload.get("execution_allowed") is not False: + return None + asset = payload.get("asset") + try: + market_value = Decimal(str(payload.get("market_value_usdt"))) + except (InvalidOperation, TypeError, ValueError): + return None + if ( + isinstance(asset, str) + and asset.isalnum() + and asset == asset.upper() + and market_value.is_finite() + and market_value > minimum_value_usdt + ): + return asset + return None + + +def _validate_wallet_selection_request( + path: Path, + *, + observed_at: datetime, + minimum_value_usdt: Decimal, + maximum_age: timedelta, +) -> None: + if observed_at.tzinfo is None or observed_at.utcoffset() is None: + raise ValueError("wallet universe observation time must be timezone-aware") + if ( + not minimum_value_usdt.is_finite() + or minimum_value_usdt <= _ZERO + or maximum_age <= timedelta(0) + ): + raise ValueError("wallet universe bounds must be positive") + if path.is_symlink() or not path.is_file(): + raise OSError("wallet balance snapshot is unavailable") + + +def read_wallet_assets_above_value( + path: Path, + *, + observed_at: datetime, + minimum_value_usdt: Decimal, + maximum_age: timedelta, +) -> WalletAssetSnapshot: + """Read one complete latest Spot snapshot without retaining stale assets.""" + + _validate_wallet_selection_request( + path, + observed_at=observed_at, + minimum_value_usdt=minimum_value_usdt, + maximum_age=maximum_age, + ) + records: list[tuple[datetime, str, Mapping[str, object]]] = [] + for raw in read_bounded_jsonl_tail(path, max_lines=2_000, max_bytes=16_000_000): + record = _wallet_record(raw) + if record is not None: + records.append(record) + if not records: + raise ValueError("wallet balance snapshot contains no valid records") + latest_time, latest_run, _ = max(records, key=lambda item: item[0]) + age = observed_at.astimezone(UTC) - latest_time + if age < timedelta(minutes=-1) or age > maximum_age: + raise ValueError("wallet balance snapshot is stale") + assets: list[str] = [] + for _, sync_run_id, payload in records: + if sync_run_id != latest_run: + continue + asset = _eligible_wallet_asset(payload, minimum_value_usdt=minimum_value_usdt) + if asset is not None: + assets.append(asset) + return WalletAssetSnapshot( + assets=tuple(sorted(set(assets))), + observed_at=latest_time, + sync_run_id=latest_run, + ) + + +@dataclass(frozen=True, slots=True) +class ResearchMarketUniverseProvider: + """Select wallet holdings plus filtered Binance-listed market-cap leaders.""" + + binance: BinanceMarketUniverseProvider + market_cap_transport: PublicMarketCapTransport + wallet_balance_path: Path + wallet_minimum_value_usdt: Decimal = Decimal("1") + wallet_maximum_age: timedelta = timedelta(minutes=30) + market_cap_asset_limit: int = 20 + market_cap_page_size: int = 100 + clock: Callable[[], datetime] = field( + default=lambda: datetime.now(UTC), repr=False, compare=False + ) + cache_filename: str = "universe-v3.json" + cache_source: str = RESEARCH_MARKET_UNIVERSE_SOURCE + + def __post_init__(self) -> None: + if not 1 <= self.market_cap_asset_limit <= 50: + raise ValueError("market-cap asset limit must be between 1 and 50") + if not self.market_cap_asset_limit <= self.market_cap_page_size <= 250: + raise ValueError("market-cap page size is invalid") + + @property + def spot_transport(self) -> object: + return self.binance.spot_transport + + @property + def futures_transport(self) -> object: + return self.binance.futures_transport + + @property + def coin_m_transport(self) -> None: + return None + + def eligible_market_snapshot(self) -> BinanceEligibleMarketSnapshot: + return self.binance.eligible_market_snapshot() + + def priority_eligible_market_snapshot(self) -> BinanceEligibleMarketSnapshot: + """Return the strict union required by the research collection policy.""" + + now = self._now() + metadata = self.binance.eligible_market_snapshot() + if metadata.blockers: + return metadata + try: + wallet = read_wallet_assets_above_value( + self.wallet_balance_path, + observed_at=now, + minimum_value_usdt=self.wallet_minimum_value_usdt, + maximum_age=self.wallet_maximum_age, + ) + except (OSError, ValueError, TypeError, json.JSONDecodeError): + return self._blocked(metadata, "WALLET_UNIVERSE_UNAVAILABLE_OR_STALE") + try: + rows = self.market_cap_transport.get_json( + "/api/v3/coins/markets", + { + "vs_currency": "usd", + "order": "market_cap_desc", + "per_page": self.market_cap_page_size, + "page": 1, + "sparkline": "false", + "include_rehypothecated": "false", + }, + ) + ranked_assets = self._ranked_assets(rows, metadata) + except ( + ExchangeHttpError, + ExchangePayloadError, + ExchangeRateLimitError, + ExchangeTransportError, + TypeError, + ValueError, + ): + return self._blocked(metadata, "PUBLIC_MARKET_CAP_UNIVERSE_UNAVAILABLE") + if len(ranked_assets) < self.market_cap_asset_limit: + return self._blocked(metadata, "PUBLIC_MARKET_CAP_UNIVERSE_INCOMPLETE") + + spot_by_asset = self._symbols_by_asset(metadata.spot_symbols) + futures_by_asset = self._symbols_by_asset(metadata.futures_symbols) + selected_assets: list[str] = [] + excluded = list(metadata.excluded_assets) + for asset in (*wallet.assets, *ranked_assets): + classification = classify_asset_eligibility( + asset, + spot_symbols=spot_by_asset.get(asset, ()), + futures_symbols=futures_by_asset.get(asset, ()), + ) + if classification.eligible: + selected_assets.append(asset) + else: + excluded.append((asset, classification.exclusion_reasons)) + selected = tuple(sorted(set(selected_assets))) + spot = tuple( + sorted( + symbol + for asset in selected + if (symbol := self._preferred_symbol(spot_by_asset.get(asset, ()))) + ) + ) + futures = tuple( + sorted( + symbol + for asset in selected + if (symbol := self._preferred_symbol(futures_by_asset.get(asset, ()))) + ) + ) + if not spot and not futures: + return self._blocked(metadata, "RESEARCH_MARKET_UNIVERSE_EMPTY") + return BinanceEligibleMarketSnapshot( + spot_symbols=spot, + futures_symbols=futures, + excluded_assets=tuple(sorted(set(excluded))), + source=RESEARCH_MARKET_UNIVERSE_SOURCE, + selected_assets=selected, + wallet_assets=tuple( + asset for asset in wallet.assets if asset in frozenset(selected) + ), + market_cap_assets=tuple( + asset for asset in ranked_assets if asset in frozenset(selected) + ), + ) + + def _ranked_assets( + self, rows: object, metadata: BinanceEligibleMarketSnapshot + ) -> tuple[str, ...]: + if not isinstance(rows, Sequence) or isinstance(rows, (str, bytes)): + raise ValueError("market-cap rows must be an array") + spot_by_asset = self._symbols_by_asset(metadata.spot_symbols) + futures_by_asset = self._symbols_by_asset(metadata.futures_symbols) + accepted: list[tuple[int, str]] = [] + seen: set[str] = set() + for raw in rows: + if not isinstance(raw, Mapping): + raise ValueError("market-cap row must be an object") + asset = str(raw.get("symbol", "")).strip().upper() + name = str(raw.get("name", "")).strip() + rank = raw.get("market_cap_rank") + cap = raw.get("market_cap") + if ( + not asset.isalnum() + or asset in seen + or not isinstance(rank, int) + or rank < 1 + or not isinstance(cap, (int, float)) + or isinstance(cap, bool) + or not math.isfinite(float(cap)) + or float(cap) <= 0 + ): + continue + classification = classify_asset_eligibility( + asset, + spot_symbols=spot_by_asset.get(asset, ()), + futures_symbols=futures_by_asset.get(asset, ()), + metadata_name=name, + ) + if not classification.eligible: + continue + accepted.append((rank, asset)) + seen.add(asset) + accepted.sort() + return tuple(asset for _, asset in accepted[: self.market_cap_asset_limit]) + + def _symbols_by_asset(self, symbols: tuple[str, ...]) -> dict[str, tuple[str, ...]]: + grouped: dict[str, list[str]] = {} + quotes = tuple(sorted(self.binance.quote_assets, key=len, reverse=True)) + for symbol in symbols: + quote = next((item for item in quotes if symbol.endswith(item)), None) + if quote is None or len(symbol) == len(quote): + continue + grouped.setdefault(symbol[: -len(quote)], []).append(symbol) + return {asset: tuple(sorted(items)) for asset, items in grouped.items()} + + def _preferred_symbol(self, symbols: tuple[str, ...]) -> str | None: + for quote in self.binance.quote_assets: + match = next((symbol for symbol in symbols if symbol.endswith(quote)), None) + if match is not None: + return match + return symbols[0] if symbols else None + + def _blocked( + self, metadata: BinanceEligibleMarketSnapshot, blocker: str + ) -> BinanceEligibleMarketSnapshot: + return BinanceEligibleMarketSnapshot( + spot_symbols=(), + futures_symbols=(), + excluded_assets=metadata.excluded_assets, + blockers=(blocker,), + source=RESEARCH_MARKET_UNIVERSE_SOURCE, + ) + + def _now(self) -> datetime: + value = self.clock() + if value.tzinfo is None: + raise ValueError("research universe clock must be timezone-aware") + return value.astimezone(UTC) + + +__all__ = ( + "RESEARCH_MARKET_UNIVERSE_SOURCE", + "ReadOnlyCoinGeckoJsonTransport", + "ResearchMarketUniverseProvider", + "WalletAssetSnapshot", + "read_wallet_assets_above_value", +) diff --git a/src/ai4binance/local_dashboard/server.py.in b/src/ai4binance/local_dashboard/server.py.in index 01a3515e..536b33ee 100644 --- a/src/ai4binance/local_dashboard/server.py.in +++ b/src/ai4binance/local_dashboard/server.py.in @@ -811,8 +811,7 @@ def market_depth_projection(data, database_path: Path): status = ( "NO_ELIGIBLE_DEPTH_TARGETS" if declared_status == "NO_ELIGIBLE_DEPTH_TARGETS" - else - "READY" + else "READY" if fully_observed else ( "COLLECTING" @@ -879,6 +878,19 @@ def futures_research_projection(data): } +def futures_research_readiness_status(projection: object) -> str: + """Classify normal bounded research progress without hiding real blockers.""" + status = projection.get("status") if isinstance(projection, dict) else None + blockers = projection.get("blockers") if isinstance(projection, dict) else None + if blockers: + return "DEGRADED" + if status == "CURRENT": + return "READY" + if status in {"RUNNING", "PROCESSED", "WAITING_FOR_RETRY"}: + return "COLLECTING" + return "DEGRADED" + + def virtual_simulation_projection(data): """Expose bounded all-symbol virtual-simulation coverage to the dashboard.""" projection = ( @@ -1042,7 +1054,7 @@ def dashboard_health_findings(result): history = result.get("market_history") or {} history_source = (result.get("sources") or {}).get("market_history") or {} - if history_source.get("status") != "READY": + if history_source.get("status") not in {"READY", "COLLECTING"}: completed = history.get("completed_streams") total = history.get("total_streams") add( @@ -1060,7 +1072,7 @@ def dashboard_health_findings(result): depth = result.get("market_depth") or {} depth_source = (result.get("sources") or {}).get("market_depth") or {} - if depth_source.get("status") != "READY": + if depth_source.get("status") not in {"READY", "COLLECTING"}: add( "MARKET_DEPTH_PARTIAL_COVERAGE", "MARKET_DATA", @@ -1078,7 +1090,7 @@ def dashboard_health_findings(result): research = result.get("futures_research") or {} research_source = (result.get("sources") or {}).get("futures_research") or {} - if research_source.get("status") != "READY": + if research_source.get("status") not in {"READY", "COLLECTING"}: add( "FUTURES_RESEARCH_NOT_READY", "RESEARCH", @@ -1138,6 +1150,33 @@ def dashboard_health_findings(result): return findings[:20] +def operational_readiness_projection( + findings: list[dict[str, object]], + services: list[dict[str, object]], +) -> dict[str, object]: + """Separate actionable readiness failures from advisory P3 observations.""" + required_services = [ + item + for item in services + if isinstance(item, dict) and item.get("required") is True + ] + ready_services = [ + item for item in required_services if item.get("status") == "READY" + ] + high_priority_findings = [ + item + for item in findings + if isinstance(item, dict) and item.get("severity") in {"P1", "P2"} + ] + return { + "status": "DEGRADED" if high_priority_findings else "READY", + "finding_count": len(findings), + "high_priority_finding_count": len(high_priority_findings), + "required_service_ready_count": len(ready_services), + "required_service_count": len(required_services), + } + + def warm_dashboard_dependencies() -> None: """Warm bounded read-only adapters before accepting loopback requests.""" # These imports are deliberately deferred from module import so the file can @@ -1275,12 +1314,9 @@ def snapshot(config): meta["status"] = result[name]["status"] elif name == "futures_research": result[name] = futures_research_projection(data) - if meta["status"] == "CURRENT" and result[name].get("status") == "CURRENT": - meta["freshness_status"] = "CURRENT" - meta["status"] = "READY" - elif meta["status"] == "CURRENT" and result[name].get("status") != "READY": + if meta["status"] == "CURRENT": meta["freshness_status"] = "CURRENT" - meta["status"] = "DEGRADED" + meta["status"] = futures_research_readiness_status(result[name]) elif name == "learning": result[name] = { "lesson_count": len(data.get("lessons", [])), @@ -1354,21 +1390,9 @@ def snapshot(config): source["producer_status"] = producer.get("status") source["producer_blockers"] = producer.get("blockers") or [] result["health_findings"] = dashboard_health_findings(result) - required_services = [ - item for item in result["services"] if item.get("required") is True - ] - ready_services = [ - item for item in required_services if item.get("status") == "READY" - ] - result["operational_readiness"] = { - "status": "DEGRADED" if result["health_findings"] else "READY", - "finding_count": len(result["health_findings"]), - "high_priority_finding_count": sum( - item.get("severity") in {"P1", "P2"} for item in result["health_findings"] - ), - "required_service_ready_count": len(ready_services), - "required_service_count": len(required_services), - } + result["operational_readiness"] = operational_readiness_projection( + result["health_findings"], result["services"] + ) return result diff --git a/src/ai4binance/skills/discovery_pipeline.py b/src/ai4binance/skills/discovery_pipeline.py index 13a730ae..14d6ccf0 100644 --- a/src/ai4binance/skills/discovery_pipeline.py +++ b/src/ai4binance/skills/discovery_pipeline.py @@ -208,6 +208,11 @@ def _append_stage_review( review: ContinuousDiscoveryStageReview, ) -> None: stage_reviews.append(review) + if review.stage_id is ContinuousDiscoveryStageId.FILTER: + # A deterministic filter rejection is a successful safety outcome. Keep + # its reasons on the candidate review and admission record without + # degrading the discovery service or the whole-system status. + return blockers.extend( blocker for blocker in review.blockers diff --git a/tests/test_accounting_collectors.py b/tests/test_accounting_collectors.py index 9b013f2e..f0afbc58 100644 --- a/tests/test_accounting_collectors.py +++ b/tests/test_accounting_collectors.py @@ -7,6 +7,8 @@ from pathlib import Path from typing import Any, cast +from websockets.exceptions import WebSocketException + import ai4binance.accounting.ui_reports as accounting_ui_reports from ai4binance.accounting import ( AccountingFileReconciler, @@ -33,6 +35,10 @@ MS = 1_784_367_000_000 +class ExpectedWebSocketFailure(WebSocketException): + """Deterministic provider transport failure for collector boundary tests.""" + + @dataclass(frozen=True, slots=True) class FakeRestSource: bad_spot_order_price: bool = False @@ -504,6 +510,36 @@ def test_user_stream_service_records_subscription_and_websocket_events( assert futures_ws.closed is True +def test_user_stream_collector_contains_provider_websocket_failures( + tmp_path: Path, +) -> None: + class ProviderFailureSession: + product_type = ProductType.SPOT + + def open(self) -> None: + raise ExpectedWebSocketFailure("provider rejected connection") + + def receive_event(self, *, timeout_seconds: float) -> dict[str, object]: + raise AssertionError(f"unexpected receive timeout {timeout_seconds}") + + def close(self) -> None: + raise AssertionError("unopened session must not be closed") + + service = AccountingUserStreamCollectorService( + ledger=BinanceAccountLedger(tmp_path, account_id="acct-a"), + sessions=(ProviderFailureSession(),), + event_limit=1, + collect_seconds=1.0, + receive_timeout_seconds=0.1, + ) + + result = service.collect_once(sync_run_id="ws-provider-failure") + + assert result.status == "DEGRADED" + assert result.rejected_count == 1 + assert result.blockers == ("SPOT_WEBSOCKET_OPEN_FAILED:ExpectedWebSocketFailure",) + + def test_file_reconciler_and_ui_report_use_rest_and_websocket_evidence( tmp_path: Path, ) -> None: diff --git a/tests/test_config_reporting.py b/tests/test_config_reporting.py index 1beb66e3..797f37ef 100644 --- a/tests/test_config_reporting.py +++ b/tests/test_config_reporting.py @@ -142,9 +142,13 @@ def test_settings_reject_minimum_history_above_request_limit() -> None: Settings(candle_limit=100, minimum_closed_candles=200) -def test_settings_defaults_to_governed_spot_oos_history_horizon() -> None: +def test_settings_defaults_to_bounded_research_history_horizons() -> None: settings = Settings() - assert settings.market_history_initial_days == 400 + assert settings.market_history_initial_days == 90 + assert settings.market_history_enrichment_days == 30 + assert settings.market_history_coin_m_enabled is False + assert settings.market_history_market_cap_asset_limit == 20 + assert settings.market_history_wallet_minimum_value_usdt == Decimal("1") assert settings.max_data_workers == 4 assert settings.market_history_max_workers == 8 assert settings.market_history_opportunity_workers == 2 diff --git a/tests/test_continuous_skill_discovery.py b/tests/test_continuous_skill_discovery.py index e0c41f4d..117bd1b0 100644 --- a/tests/test_continuous_skill_discovery.py +++ b/tests/test_continuous_skill_discovery.py @@ -358,6 +358,7 @@ def test_rejected_filter_keeps_library_record_and_marker_readiness( ) assert report.drafts_created == 0 + assert report.blockers == () assert report.admission_records[0].status == "READY_FOR_LIBRARY_PR" assert report.admission_records[0].approval_marker_present is True admission_payload = json.loads( diff --git a/tests/test_data_acquisition.py b/tests/test_data_acquisition.py index a1d84431..8157d9a7 100644 --- a/tests/test_data_acquisition.py +++ b/tests/test_data_acquisition.py @@ -113,6 +113,7 @@ def test_local_snapshot_priority_reuses_canonical_liquidity_order( tmp_path: Path, ) -> None: from ai4binance.cli.runtime import ( + _canonical_resident_runtime_symbol, _virtual_market_priority_symbols, _virtual_market_ranked_symbols, ) @@ -142,12 +143,17 @@ def test_local_snapshot_priority_reuses_canonical_liquidity_order( ) assert _virtual_market_priority_symbols(settings, (), ("BTCUSDT",), NOW) == () assert _virtual_market_ranked_symbols(settings, (), NOW) == ("SOLUSDT",) + assert _canonical_resident_runtime_symbol(settings, NOW) == "SOLUSDT" assert ( _virtual_market_priority_symbols( settings, (), ("SOLUSDT",), NOW + timedelta(hours=1) ) == () ) + assert ( + _canonical_resident_runtime_symbol(settings, NOW + timedelta(hours=1)) + == settings.symbol + ) def test_local_archive_supplies_candles_to_shared_snapshot(tmp_path: Path) -> None: diff --git a/tests/test_futures_multitf.py b/tests/test_futures_multitf.py index 6fcb7934..d0aec1ac 100644 --- a/tests/test_futures_multitf.py +++ b/tests/test_futures_multitf.py @@ -139,6 +139,7 @@ def test_current_cycle_keeps_the_local_futures_opportunity_monitor_running( state_path.write_text( json.dumps( { + "eligible_symbols": ["BTCUSDT", "ETHUSDT"], "completed_windows": { "BTCUSDT": current_window, "ETHUSDT": current_window, @@ -439,36 +440,89 @@ def test_runner_rejects_incomplete_timeframe_set(tmp_path: Path) -> None: def test_eligible_symbols_validates_and_normalizes_cached_universe( tmp_path: Path, ) -> None: - from ai4binance.cli.futures_multitf import _eligible_symbols + from ai4binance.cli.futures_multitf import ( + _eligible_symbols, + _retain_current_universe_state, + ) + from ai4binance.integrations.research_market_universe import ( + RESEARCH_MARKET_UNIVERSE_SOURCE, + ) + + base = { + "observed_at": datetime.now(UTC).isoformat(), + "source": RESEARCH_MARKET_UNIVERSE_SOURCE, + "execution_allowed": False, + "live_eligibility_status": "LIVE_ORDER_BLOCKED", + } with pytest.raises(ValueError, match="UNIVERSE_UNAVAILABLE"): _eligible_symbols(tmp_path) path = tmp_path / "universe-v3.json" - path.write_text(json.dumps({"futures_symbols": "BTCUSDT"}), encoding="utf-8") + path.write_text( + json.dumps({**base, "futures_symbols": "BTCUSDT"}), encoding="utf-8" + ) with pytest.raises(ValueError, match="UNIVERSE_INVALID"): _eligible_symbols(tmp_path) - path.write_text(json.dumps({"futures_symbols": ["../bad", 7]}), encoding="utf-8") + path.write_text( + json.dumps( + { + **base, + "observed_at": datetime.now().replace(microsecond=0).isoformat(), + "futures_symbols": ["BTCUSDT"], + } + ), + encoding="utf-8", + ) + with pytest.raises(ValueError, match="UNIVERSE_INVALID"): + _eligible_symbols(tmp_path) + + path.write_text( + json.dumps({**base, "futures_symbols": ["../bad", 7]}), + encoding="utf-8", + ) with pytest.raises(ValueError, match="UNIVERSE_EMPTY"): _eligible_symbols(tmp_path) path.write_text( json.dumps( { + **base, "futures_symbols": [ " ethusdt ", "BTCUSDT", "BTCUSDT", "not-valid!", 7, - ] + ], } ), encoding="utf-8", ) assert _eligible_symbols(tmp_path) == ("BTCUSDT", "ETHUSDT") + retained = _retain_current_universe_state( + { + "eligible_symbols": ["OLDUSDT"], + "attempted_windows": {"OLDUSDT": "old", "BTCUSDT": "current"}, + "completed_windows": {"OLDUSDT": "old"}, + "retry_after": {"OLDUSDT": "old"}, + "unavailable_windows": {"OLDUSDT": "old"}, + "active_symbol": "OLDUSDT", + "monitor_symbol": "OLDUSDT", + "completed_tuning_triggers": ["old-trigger"], + }, + ("BTCUSDT", "ETHUSDT"), + ) + assert retained["attempted_windows"] == {"BTCUSDT": "current"} + assert retained["completed_windows"] == {} + assert retained["retry_after"] == {} + assert retained["unavailable_windows"] == {} + assert retained["completed_tuning_triggers"] == [] + assert "active_symbol" not in retained + assert "monitor_symbol" not in retained + def test_learning_projection_handles_unavailable_and_bounded_collections( monkeypatch: pytest.MonkeyPatch, diff --git a/tests/test_local_dashboard_source.py b/tests/test_local_dashboard_source.py index 9cfa6c7c..4c17f597 100644 --- a/tests/test_local_dashboard_source.py +++ b/tests/test_local_dashboard_source.py @@ -71,12 +71,59 @@ def test_dashboard_exposes_useful_cycle_and_producer_health() -> None: assert "DISABLED" in views -def test_dashboard_recognizes_completed_futures_research() -> None: - server = (SOURCE / "server.py.in").read_text(encoding="utf-8") +def test_dashboard_keeps_collection_progress_and_p3_advice_out_of_readiness() -> None: + module = runpy.run_path(str(SOURCE / "server.py.in")) + findings = module["dashboard_health_findings"]( + { + "market_history": {"completed_streams": 12, "total_streams": 301}, + "market_depth": {}, + "futures_research": {}, + "learning": {"lesson_count": 0, "experiment_count": 0}, + "services": [], + "sources": { + "market_history": {"status": "COLLECTING"}, + "market_depth": {"status": "COLLECTING"}, + "futures_research": {"status": "COLLECTING"}, + "auto_audit": {"status": "CURRENT"}, + }, + } + ) + projection = module["operational_readiness_projection"](findings, []) + + assert [item["finding_id"] for item in findings] == ["LEARNING_SUMMARY_EMPTY"] + assert projection == { + "status": "READY", + "finding_count": 1, + "high_priority_finding_count": 0, + "required_service_ready_count": 0, + "required_service_count": 0, + } + + +def test_dashboard_degrades_for_high_priority_health_findings() -> None: + module = runpy.run_path(str(SOURCE / "server.py.in")) + projection = module["operational_readiness_projection"]( + [{"severity": "P2"}], + [{"required": True, "status": "DEGRADED"}], + ) + + assert projection["status"] == "DEGRADED" + assert projection["high_priority_finding_count"] == 1 + assert projection["required_service_ready_count"] == 0 + - assert 'result[name].get("status") == "CURRENT"' in server - assert 'meta["status"] = "READY"' in server - assert '"unavailable_symbol_count"' in server +def test_dashboard_classifies_futures_research_progress_and_failures() -> None: + module = runpy.run_path(str(SOURCE / "server.py.in")) + classify = module["futures_research_readiness_status"] + + assert classify({"status": "CURRENT", "blockers": []}) == "READY" + assert classify({"status": "RUNNING", "blockers": []}) == "COLLECTING" + assert classify({"status": "PROCESSED", "blockers": []}) == "COLLECTING" + assert classify({"status": "WAITING_FOR_RETRY", "blockers": []}) == "COLLECTING" + assert ( + classify({"status": "CURRENT_WITH_SOURCE_GAPS", "blockers": []}) == "DEGRADED" + ) + assert classify({"status": "CURRENT", "blockers": ["SOURCE_GAP"]}) == "DEGRADED" def test_virtual_market_separates_trade_records_from_potential_opportunities() -> None: diff --git a/tests/test_market_depth.py b/tests/test_market_depth.py index 04463f50..53a0fcd7 100644 --- a/tests/test_market_depth.py +++ b/tests/test_market_depth.py @@ -128,6 +128,49 @@ def test_checkpoint_compacts_superseded_depth_events(tmp_path: Path) -> None: journal.close() +def test_depth_journal_removes_streams_outside_current_universe( + tmp_path: Path, +) -> None: + journal = DepthJournal(tmp_path / "depth.sqlite3") + journal.append( + [ + ("spot", "BTCUSDT", "snapshot", snapshot(), NOW.timestamp()), + ("spot", "ETHUSDT", "snapshot", snapshot(), NOW.timestamp()), + ( + "usd_m_futures", + "ETHUSDT", + "snapshot", + snapshot(), + NOW.timestamp(), + ), + ] + ) + + result = journal.retain_streams({"spot": ("BTCUSDT",)}) + + assert result == { + "removed_stream_count": 2, + "removed_event_count": 2, + "removed_head_count": 2, + } + assert journal.connection.execute( + "SELECT market,symbol FROM depth_heads" + ).fetchall() == [("spot", "BTCUSDT")] + assert journal.connection.execute( + "SELECT market,symbol FROM depth_events" + ).fetchall() == [("spot", "BTCUSDT")] + journal.close() + + +def test_depth_journal_rejects_invalid_retention_identity(tmp_path: Path) -> None: + journal = DepthJournal(tmp_path / "depth.sqlite3") + with pytest.raises(ValueError, match="market"): + journal.retain_streams({"invalid": ("BTCUSDT",)}) + with pytest.raises(ValueError, match="symbol"): + journal.retain_streams({"spot": ("btc/usdt",)}) + journal.close() + + @pytest.mark.parametrize("market", ["spot", "usd_m_futures", "coin_m_futures"]) def test_gap_and_unsynchronized_are_never_readable(tmp_path: Path, market: str) -> None: journal = DepthJournal(tmp_path / "depth.db") diff --git a/tests/test_market_history_continuous.py b/tests/test_market_history_continuous.py index 74aa1a4e..c1dff4c9 100644 --- a/tests/test_market_history_continuous.py +++ b/tests/test_market_history_continuous.py @@ -839,6 +839,37 @@ def test_background_stream_plan_finishes_each_symbol_before_the_next( ] +def test_staged_background_stream_plan_refreshes_each_timeframe_across_universe( + tmp_path: Path, +) -> None: + instance = collector(tmp_path, Transport()) + instance.on_symbol_screen = lambda *_args: {} + transport = Transport() + + streams = instance._interleaved_stream_work( + ( + ("spot", "AUSDT", transport), + ("spot", "BUSDT", transport), + ("usd_m_futures", "CUSDT", transport), + ) + ) + screening = [ + (market, symbol, timeframe) + for market, symbol, _, kind, timeframe in streams + if kind == "klines" and timeframe in _SCREEN_TIMEFRAMES + ] + + assert screening == [ + (market, symbol, timeframe) + for timeframe in _SCREEN_TIMEFRAMES + for market, symbol in ( + ("spot", "AUSDT"), + ("spot", "BUSDT"), + ("usd_m_futures", "CUSDT"), + ) + ] + + def test_symbol_with_incomplete_stream_does_not_start_opportunity_analysis( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: @@ -2054,9 +2085,9 @@ def test_direct_timeframe_bootstrap_uses_its_own_closed_candle_window( minimum_candles=200, ) - assert instance._candle_initial_start(NOW, timedelta(minutes=5)) == ( - datetime(2026, 6, 4, 0, 0, tzinfo=UTC) - ) + assert instance._candle_initial_start( + NOW, timedelta(minutes=5), history_days=30 + ) == (NOW.replace(hour=0, minute=0) - timedelta(days=30)) assert instance._candle_initial_start(NOW, timedelta(days=1)) == ( NOW.replace(hour=0, minute=0) - timedelta(days=201) ) diff --git a/tests/test_opportunity_monitor.py b/tests/test_opportunity_monitor.py index 64f8a43d..1ac67123 100644 --- a/tests/test_opportunity_monitor.py +++ b/tests/test_opportunity_monitor.py @@ -86,7 +86,9 @@ def test_monitor_helper_boundaries_and_research_estimates( "Universe", (), {"spot_symbols": ("BTCUSDT",), "futures_symbols": ("ETHUSDT",)} )() monkeypatch.setattr( - monitor_module, "read_cached_market_universe", lambda *_args: universe + monitor_module, + "read_cached_market_universe", + lambda *_args, **_kwargs: universe, ) assert monitor_module.market_symbols(tmp_path, "USD_M_FUTURES", NOW) == ("ETHUSDT",) diff --git a/tests/test_research_market_universe.py b/tests/test_research_market_universe.py new file mode 100644 index 00000000..4f93c5ec --- /dev/null +++ b/tests/test_research_market_universe.py @@ -0,0 +1,289 @@ +from __future__ import annotations + +import json +from datetime import UTC, datetime, timedelta +from decimal import Decimal +from pathlib import Path +from types import SimpleNamespace + +import pytest + +from ai4binance.data.market_universe_retention import MarketUniverseRetention +from ai4binance.integrations.binance import BinanceEligibleMarketSnapshot +from ai4binance.integrations.research_market_universe import ( + ResearchMarketUniverseProvider, + read_wallet_assets_above_value, +) + +NOW = datetime(2026, 9, 26, 1, 0, tzinfo=UTC) +MARKET_CAP_ASSETS = ( + "BTC", + "ETH", + "BNB", + "SOL", + "XRP", + "DOGE", + "ADA", + "TRX", + "AVAX", + "LINK", + "DOT", + "BCH", + "LTC", + "NEAR", + "AAVE", + "UNI", + "ICP", + "FIL", + "XLM", + "HBAR", +) + + +class _MarketCapTransport: + def __init__(self, rows: list[dict[str, object]]) -> None: + self.rows = rows + self.calls: list[tuple[str, object]] = [] + + def get_json(self, path: str, params: object = None) -> object: + self.calls.append((path, params)) + return self.rows + + +def _wallet_record( + *, asset: str, value: str, run_id: str, observed_at: datetime +) -> dict[str, object]: + return { + "timestamp": observed_at.isoformat(), + "payload": { + "asset": asset, + "market_value_usdt": value, + "execution_allowed": False, + "envelope": { + "sync_run_id": run_id, + "event_time": observed_at.isoformat(), + }, + }, + } + + +def _write_wallet(path: Path, records: list[dict[str, object]]) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text( + "".join(json.dumps(record) + "\n" for record in records), + encoding="utf-8", + ) + + +def _metadata() -> BinanceEligibleMarketSnapshot: + assets = (*MARKET_CAP_ASSETS, "ATOM", "USDT", "WBTC") + return BinanceEligibleMarketSnapshot( + spot_symbols=tuple(sorted(f"{asset}USDT" for asset in assets)), + futures_symbols=tuple( + sorted(f"{asset}USDT" for asset in assets if asset not in {"USDT", "WBTC"}) + ), + ) + + +def test_wallet_reader_uses_only_latest_complete_run_and_strict_value_floor( + tmp_path: Path, +) -> None: + path = tmp_path / "balance_snapshots.jsonl" + _write_wallet( + path, + [ + _wallet_record( + asset="OLD", + value="100", + run_id="old", + observed_at=NOW - timedelta(minutes=5), + ), + _wallet_record(asset="ATOM", value="1.01", run_id="new", observed_at=NOW), + _wallet_record(asset="DUST", value="1", run_id="new", observed_at=NOW), + ], + ) + + snapshot = read_wallet_assets_above_value( + path, + observed_at=NOW, + minimum_value_usdt=Decimal("1"), + maximum_age=timedelta(minutes=30), + ) + + assert snapshot.assets == ("ATOM",) + assert snapshot.sync_run_id == "new" + + +def test_research_universe_unions_wallet_and_filtered_market_cap( + tmp_path: Path, +) -> None: + wallet_path = tmp_path / "balance_snapshots.jsonl" + _write_wallet( + wallet_path, + [_wallet_record(asset="ATOM", value="2", run_id="new", observed_at=NOW)], + ) + rows = [ + {"symbol": "usdt", "name": "Tether", "market_cap_rank": 1, "market_cap": 1}, + { + "symbol": "wbtc", + "name": "Wrapped Bitcoin", + "market_cap_rank": 2, + "market_cap": 1, + }, + { + "symbol": "usds", + "name": "USDS", + "market_cap_rank": 3, + "market_cap": 1, + }, + *[ + { + "symbol": asset.lower(), + "name": asset, + "market_cap_rank": rank, + "market_cap": 1_000_000 - rank, + } + for rank, asset in enumerate(MARKET_CAP_ASSETS, start=4) + ], + ] + metadata = _metadata() + binance = SimpleNamespace( + quote_assets=("USDT", "USDC"), + spot_transport=object(), + futures_transport=object(), + eligible_market_snapshot=lambda: metadata, + ) + transport = _MarketCapTransport(rows) + provider = ResearchMarketUniverseProvider( + binance=binance, # type: ignore[arg-type] + market_cap_transport=transport, + wallet_balance_path=wallet_path, + clock=lambda: NOW, + ) + + result = provider.priority_eligible_market_snapshot() + + assert result.blockers == () + assert result.wallet_assets == ("ATOM",) + assert result.market_cap_assets == MARKET_CAP_ASSETS + assert len(result.market_cap_assets) == 20 + assert "USDT" not in result.selected_assets + assert "WBTC" not in result.selected_assets + assert "USDS" not in result.selected_assets + assert "ATOMUSDT" in result.spot_symbols + assert result.execution_allowed is False + assert result.live_eligibility_status == "LIVE_ORDER_BLOCKED" + assert transport.calls[0][0] == "/api/v3/coins/markets" + + +def test_research_universe_fails_closed_for_stale_wallet(tmp_path: Path) -> None: + wallet_path = tmp_path / "balance_snapshots.jsonl" + _write_wallet( + wallet_path, + [ + _wallet_record( + asset="ATOM", + value="2", + run_id="old", + observed_at=NOW - timedelta(hours=1), + ) + ], + ) + metadata = _metadata() + provider = ResearchMarketUniverseProvider( + binance=SimpleNamespace( + quote_assets=("USDT",), + spot_transport=object(), + futures_transport=object(), + eligible_market_snapshot=lambda: metadata, + ), # type: ignore[arg-type] + market_cap_transport=_MarketCapTransport([]), + wallet_balance_path=wallet_path, + clock=lambda: NOW, + ) + + result = provider.priority_eligible_market_snapshot() + + assert result.spot_symbols == () + assert result.blockers == ("WALLET_UNIVERSE_UNAVAILABLE_OR_STALE",) + + +def test_retention_removes_only_out_of_scope_symbol_directories( + tmp_path: Path, +) -> None: + archive = tmp_path / "market" + sources = tmp_path / "market_sources" + monitor = tmp_path / "monitor" + replay = tmp_path / "replay" + validation = tmp_path / "validation" + keep = archive / "spot" / "BTCUSDT" + drop = archive / "spot" / "OLDUSDT" + metadata = archive / "spot" / "metadata" + coin_m_archive = archive / "coin_m_futures" / "metadata" + source_keep = sources / "data/spot/daily/klines/BTCUSDT" + source_drop = sources / "data/spot/daily/klines/OLDUSDT" + coin_m = sources / "data/futures/cm/daily/klines/BTCUSD_PERP" + monitor_keep = monitor / "SPOT" / "BTCUSDT" + monitor_drop = monitor / "SPOT" / "OLDUSDT" + replay.mkdir() + replay_keep = replay / "BTCUSDT-5m-window.json" + replay_drop = replay / "OLDUSDT-5m-window.json" + replay_keep.write_text("{}", encoding="utf-8") + replay_drop.write_text("{}", encoding="utf-8") + validation_keep = validation / "BTCUSDT" + validation_drop = validation / "OLDUSDT" + for directory in ( + keep, + drop, + metadata, + coin_m_archive, + source_keep, + source_drop, + coin_m, + monitor_keep, + monitor_drop, + validation_keep, + validation_drop, + ): + directory.mkdir(parents=True) + (directory / "value.json").write_text("{}", encoding="utf-8") + retention = MarketUniverseRetention( + archive, + sources, + monitor, + futures_replay_root=replay, + futures_artifact_roots=(validation,), + ) + universe = BinanceEligibleMarketSnapshot( + spot_symbols=("BTCUSDT",), futures_symbols=("BTCUSDT",) + ) + + result = retention.prune(universe) + + assert result["status"] == "PRUNED" + assert keep.is_dir() + assert metadata.is_dir() + assert not coin_m_archive.exists() + assert source_keep.is_dir() + assert monitor_keep.is_dir() + assert replay_keep.is_file() + assert validation_keep.is_dir() + assert not drop.exists() + assert not source_drop.exists() + assert not coin_m.exists() + assert not monitor_drop.exists() + assert not replay_drop.exists() + assert not validation_drop.exists() + + +@pytest.mark.parametrize("invalid_value", [Decimal("0"), Decimal("NaN")]) +def test_wallet_reader_rejects_invalid_floor( + tmp_path: Path, invalid_value: Decimal +) -> None: + with pytest.raises(ValueError, match="bounds must be positive"): + read_wallet_assets_above_value( + tmp_path / "missing.jsonl", + observed_at=NOW, + minimum_value_usdt=invalid_value, + maximum_age=timedelta(minutes=30), + ) diff --git a/tests/test_runtime_cycle.py b/tests/test_runtime_cycle.py index 924002c3..e69e17fa 100644 --- a/tests/test_runtime_cycle.py +++ b/tests/test_runtime_cycle.py @@ -354,6 +354,47 @@ def trades( return rows +class MultiAssetSpotWalletReader(SpotWalletReader): + def account(self) -> object: + payload = super().account() + assert isinstance(payload, dict) + payload["balances"] = [ + {"asset": "HOT", "free": "10", "locked": "0"}, + {"asset": "XRP", "free": "4", "locked": "0"}, + {"asset": "USDT", "free": "100", "locked": "0"}, + ] + return payload + + +class MultiAssetSpotPrices: + def ticker_price(self, symbol: str) -> Decimal: + return {"HOTUSDT": Decimal("1"), "XRPUSDT": Decimal("2")}[symbol] + + +class MultiAssetSpotHistory: + def trades( + self, symbol: str, *, from_id: int | None = None, limit: int = 1_000 + ) -> object: + assert from_id == 0 + assert limit == 1_000 + quantity, quote = { + "HOTUSDT": ("10", "5"), + "XRPUSDT": ("4", "4"), + }[symbol] + return [ + { + "id": 1, + "price": str(Decimal(quote) / Decimal(quantity)), + "qty": quantity, + "quoteQty": quote, + "commission": "0", + "commissionAsset": "USDT", + "isBuyer": True, + "time": 1, + } + ] + + def test_runtime_attaches_read_only_cost_basis_to_portfolio_analytics() -> None: runtime = ReadOnlyRuntimeCycle( symbol="HOTUSDT", @@ -408,6 +449,29 @@ def test_runtime_does_not_use_blocked_cost_basis_for_portfolio_pnl() -> None: assert report.live_eligibility_status == "LIVE_ORDER_BLOCKED" +def test_runtime_reconciles_cost_basis_for_every_valued_wallet_asset() -> None: + runtime = ReadOnlyRuntimeCycle( + symbol="HOTUSDT", + timeframes=("1h",), + spot_acquirer=RecordingAcquisition(), + spot_wallet_service=WalletSnapshotService(MultiAssetSpotWalletReader()), + futures_account_service=FuturesAccountSnapshotService(FuturesReader()), + futures_advisor=DerivativesCollectorStub(), + research_service=research_service(), + investment_manager=RuntimeInvestmentManager(InvestmentManagementAssistant()), + analytics_service=PortfolioAnalyticsService(MultiAssetSpotPrices()), + cost_basis_service=CostBasisService(MultiAssetSpotHistory()), + ) + + report = runtime.run(NOW) + + assert report.portfolio_analytics is not None + valued = {item.asset: item for item in report.portfolio_analytics.valued_assets} + assert valued["HOT"].average_cost_usdt == Decimal("0.5") + assert valued["XRP"].average_cost_usdt == Decimal("1") + assert "PORTFOLIO_COST_BASIS_UNAVAILABLE" not in report.blockers + + def test_runtime_values_spot_wallet_even_when_market_cycle_is_blocked() -> None: runtime = ReadOnlyRuntimeCycle( symbol="HOTUSDT", diff --git a/tests/test_service_manifest.py b/tests/test_service_manifest.py index 04708992..d095822b 100644 --- a/tests/test_service_manifest.py +++ b/tests/test_service_manifest.py @@ -54,7 +54,7 @@ def test_service_manifest_is_shared_safe_and_complete() -> None: assert scheduled["futures-multitf"].useful_state_file == ( "futures-multitf-latest.json" ) - assert scheduled["futures-multitf"].enabled is False + assert scheduled["futures-multitf"].enabled is True for spec in specs: assert spec.startup_trigger == "AtLogOn" assert spec.allow_start_if_on_batteries is True @@ -315,6 +315,22 @@ def test_ykb_daily_report_task_refreshes_stale_report_on_logon() -> None: assert "live_eligibility_status" in script_text assert "LIVE_ORDER_BLOCKED" in script_text assert "Invoke-PythonUtf8Command -ArgumentList @(" in script_text + assert "report_status" in script_text + assert "report_blockers" in script_text + assert "$reportExitCode -eq 2" in script_text + assert '$humanReport.Status -ne "RUNNING_WITH_BLOCKERS"' in script_text + + +def test_primary_local_reasoning_publishes_resident_pid_and_lock() -> None: + script_text = (Path("scripts") / "primary_local_reasoning_agent.ps1").read_text( + encoding="utf-8" + ) + + assert "runtime\\state\\primary-local-reasoning.lock" in script_text + assert "pid = $PID" in script_text + assert "Write-PrimaryLock" in script_text + assert "Remove-PrimaryLock" in script_text + assert 'live_eligibility_status = "LIVE_ORDER_BLOCKED"' in script_text def test_ykb_daily_report_script_uses_utf8_redirect_helper() -> None: